File
Blob: firmware/vendor/sctp-proto/src/association/mod.rs
| 1 | use crate::association::state::{AckMode, AckState, AssociationState}; |
| 2 | use crate::association::stats::AssociationStats; |
| 3 | use crate::chunk::Chunk; |
| 4 | use crate::chunk::ErrorCauseUnrecognizedChunkType; |
| 5 | use crate::chunk::USER_INITIATED_ABORT; |
| 6 | use crate::chunk::chunk_abort::ChunkAbort; |
| 7 | use crate::chunk::chunk_cookie_ack::ChunkCookieAck; |
| 8 | use crate::chunk::chunk_cookie_echo::ChunkCookieEcho; |
| 9 | use crate::chunk::chunk_error::ChunkError; |
| 10 | use crate::chunk::chunk_forward_tsn::{ChunkForwardTsn, ChunkForwardTsnStream}; |
| 11 | use crate::chunk::chunk_heartbeat::ChunkHeartbeat; |
| 12 | use crate::chunk::chunk_heartbeat_ack::ChunkHeartbeatAck; |
| 13 | use crate::chunk::chunk_i_forward_tsn::ChunkIForwardTsn; |
| 14 | use crate::chunk::chunk_init::{ChunkInit, ChunkInitAck}; |
| 15 | use crate::chunk::chunk_payload_data::{ChunkPayloadData, PayloadProtocolIdentifier}; |
| 16 | use crate::chunk::chunk_reconfig::ChunkReconfig; |
| 17 | use crate::chunk::chunk_selective_ack::ChunkSelectiveAck; |
| 18 | use crate::chunk::chunk_shutdown::ChunkShutdown; |
| 19 | use crate::chunk::chunk_shutdown_ack::ChunkShutdownAck; |
| 20 | use crate::chunk::chunk_shutdown_complete::ChunkShutdownComplete; |
| 21 | use crate::chunk::chunk_type::CT_FORWARD_TSN; |
| 22 | use crate::config::COMMON_HEADER_SIZE; |
| 23 | use crate::config::DATA_CHUNK_HEADER_SIZE; |
| 24 | use crate::config::DEFAULT_SCTP_PORT; |
| 25 | use crate::config::{ServerConfig, TransportConfig}; |
| 26 | use crate::error::{Error, Result}; |
| 27 | use crate::packet::{CommonHeader, Packet}; |
| 28 | use crate::param::Param; |
| 29 | use crate::param::param_heartbeat_info::ParamHeartbeatInfo; |
| 30 | use crate::param::param_outgoing_reset_request::ParamOutgoingResetRequest; |
| 31 | use crate::param::param_reconfig_response::{ParamReconfigResponse, ReconfigResult}; |
| 32 | use crate::param::param_state_cookie::ParamStateCookie; |
| 33 | use crate::param::param_supported_extensions::ParamSupportedExtensions; |
| 34 | use crate::queue::payload_queue::PayloadQueue; |
| 35 | use crate::queue::pending_queue::PendingQueue; |
| 36 | use crate::queue::reassembly_queue::ReassemblyQueue; |
| 37 | use crate::shared::{AssociationEventInner, AssociationId, EndpointEvent, EndpointEventInner}; |
| 38 | use crate::util::{sna16lt, sna32gt, sna32gte, sna32lt, sna32lte}; |
| 39 | use crate::{AssociationEvent, Payload, Side, Transmit}; |
| 40 | use stream::{ReliabilityType, Stream, StreamEvent, StreamId, StreamResetError, StreamState}; |
| 41 | use timer::{ACK_INTERVAL, RtoManager, Timer, TimerTable}; |
| 42 | |
| 43 | use crate::association::stream::RecvSendState; |
| 44 | use alloc::boxed::Box; |
| 45 | use alloc::collections::VecDeque; |
| 46 | use alloc::string::String; |
| 47 | use alloc::sync::Arc; |
| 48 | use alloc::vec; |
| 49 | use alloc::vec::Vec; |
| 50 | use bytes::Bytes; |
| 51 | use core::net::{IpAddr, SocketAddr}; |
| 52 | use core::num::NonZeroU32; |
| 53 | use core::str::FromStr; |
| 54 | use core::time::Duration; |
| 55 | use log::{debug, error, trace, warn}; |
| 56 | use rand::random; |
| 57 | use rustc_hash::{FxHashMap, FxHashSet}; |
| 58 | use std::collections::HashMap; |
| 59 | use std::time::Instant; |
| 60 | use thiserror::Error; |
| 61 | |
| 62 | pub(crate) mod state; |
| 63 | pub(crate) mod stats; |
| 64 | pub(crate) mod stream; |
| 65 | mod timer; |
| 66 | |
| 67 | #[cfg(test)] |
| 68 | mod association_test; |
| 69 | #[cfg(test)] |
| 70 | mod receive_limits_test; |
| 71 | |
| 72 | /// Reasons why an association might be lost |
| 73 | #[non_exhaustive] |
| 74 | #[derive(Debug, Error, Clone, PartialEq)] |
| 75 | pub enum AssociationError { |
| 76 | /// Handshake failed |
| 77 | #[error("{0}")] |
| 78 | HandshakeFailed(#[from] Error), |
| 79 | /// The peer violated the SCTP specification as understood by this implementation |
| 80 | #[error("transport error")] |
| 81 | TransportError, |
| 82 | /// The peer's SCTP stack aborted the association |
| 83 | #[error("aborted by peer")] |
| 84 | AssociationClosed, |
| 85 | /// The peer closed the association |
| 86 | #[error("closed by peer")] |
| 87 | ApplicationClosed, |
| 88 | /// The peer is unable to continue processing this association, usually due to having restarted |
| 89 | #[error("reset by peer")] |
| 90 | Reset, |
| 91 | /// Communication with the peer has lapsed for longer than the negotiated idle timeout |
| 92 | /// |
| 93 | /// If neither side is sending keep-alives, an association will time out after a long enough idle |
| 94 | /// period even if the peer is still reachable |
| 95 | #[error("timed out")] |
| 96 | TimedOut, |
| 97 | /// The local application closed the association |
| 98 | #[error("closed")] |
| 99 | LocallyClosed, |
| 100 | } |
| 101 | |
| 102 | /// Events of interest to the application |
| 103 | #[non_exhaustive] |
| 104 | #[derive(Debug)] |
| 105 | pub enum Event { |
| 106 | /// The association was successfully established |
| 107 | Connected, |
| 108 | /// The association handshake failed |
| 109 | /// |
| 110 | /// Emitted if the handshake (INIT/COOKIE exchange) fails. |
| 111 | HandshakeFailed { |
| 112 | /// Reason that the handshake failed |
| 113 | reason: AssociationError, |
| 114 | }, |
| 115 | /// The association was lost |
| 116 | /// |
| 117 | /// Emitted if the peer closes the association or an error is encountered |
| 118 | AssociationLost { |
| 119 | /// Reason that the association was closed |
| 120 | reason: AssociationError, |
| 121 | }, |
| 122 | /// Stream events |
| 123 | Stream(StreamEvent), |
| 124 | /// One or more application datagrams have been received |
| 125 | DatagramReceived, |
| 126 | } |
| 127 | |
| 128 | /// Multiset of stream ids owing a terminal reset event. |
| 129 | /// Each `Finished` for an id arms one result, a re-created id is a distinct |
| 130 | /// generation and owes its own. |
| 131 | #[derive(Debug, Default)] |
| 132 | struct PendingResetCompletions(FxHashMap<StreamId, usize>); |
| 133 | |
| 134 | impl PendingResetCompletions { |
| 135 | /// Arm one completion for this id. |
| 136 | fn insert(&mut self, id: StreamId) { |
| 137 | *self.0.entry(id).or_default() += 1; |
| 138 | } |
| 139 | |
| 140 | fn contains(&self, id: &StreamId) -> bool { |
| 141 | self.0.contains_key(id) |
| 142 | } |
| 143 | |
| 144 | /// Consume one armed completion for this id. |
| 145 | fn take_one(&mut self, id: StreamId) -> bool { |
| 146 | if let Some(count) = self.0.get_mut(&id) { |
| 147 | *count -= 1; |
| 148 | if *count == 0 { |
| 149 | self.0.remove(&id); |
| 150 | } |
| 151 | true |
| 152 | } else { |
| 153 | false |
| 154 | } |
| 155 | } |
| 156 | |
| 157 | fn clear(&mut self) { |
| 158 | self.0.clear(); |
| 159 | } |
| 160 | } |
| 161 | |
| 162 | #[derive(Debug, Copy, Clone)] |
| 163 | enum DeferredForwardTsnKind { |
| 164 | Ordered { |
| 165 | last_ssn: u16, |
| 166 | new_cumulative_tsn: u32, |
| 167 | }, |
| 168 | Unordered { |
| 169 | new_cumulative_tsn: u32, |
| 170 | }, |
| 171 | } |
| 172 | |
| 173 | #[derive(Debug, Copy, Clone)] |
| 174 | struct DeferredForwardTsn { |
| 175 | /// The reset boundary ending the generation this update belongs to. |
| 176 | /// None identifies the active tail generation after all accepted resets. |
| 177 | generation_boundary: Option<u32>, |
| 178 | kind: DeferredForwardTsnKind, |
| 179 | } |
| 180 | |
| 181 | ///Association represents an SCTP association |
| 182 | //13.2. Parameters Necessary per Association (i.e., the TCB) |
| 183 | //Peer : Tag value to be sent in every packet and is received |
| 184 | //Verification: in the INIT or INIT ACK chunk. |
| 185 | //Tag : |
| 186 | // |
| 187 | //My : Tag expected in every inbound packet and sent in the |
| 188 | //Verification: INIT or INIT ACK chunk. |
| 189 | // |
| 190 | //Tag : |
| 191 | //State : A state variable indicating what state the association |
| 192 | // : is in, i.e., COOKIE-WAIT, COOKIE-ECHOED, ESTABLISHED, |
| 193 | // : SHUTDOWN-PENDING, SHUTDOWN-SENT, SHUTDOWN-RECEIVED, |
| 194 | // : SHUTDOWN-ACK-SENT. |
| 195 | // |
| 196 | // No Closed state is illustrated since if a |
| 197 | // association is Closed its TCB SHOULD be removed. |
| 198 | #[derive(Debug)] |
| 199 | pub struct Association { |
| 200 | side: Side, |
| 201 | state: AssociationState, |
| 202 | handshake_completed: bool, |
| 203 | max_send_message_size: u32, |
| 204 | max_receive_message_size: u32, |
| 205 | receive_limits: Option<crate::ReceiveLimits>, |
| 206 | inflight_queue_length: usize, |
| 207 | will_send_shutdown: bool, |
| 208 | bytes_received: usize, |
| 209 | bytes_sent: usize, |
| 210 | |
| 211 | peer_verification_tag: u32, |
| 212 | my_verification_tag: u32, |
| 213 | my_next_tsn: u32, |
| 214 | peer_last_tsn: u32, |
| 215 | // for RTT measurement |
| 216 | min_tsn2measure_rtt: u32, |
| 217 | will_send_forward_tsn: bool, |
| 218 | will_retransmit_fast: bool, |
| 219 | will_retransmit_reconfig: bool, |
| 220 | |
| 221 | will_send_shutdown_ack: bool, |
| 222 | will_send_shutdown_complete: bool, |
| 223 | |
| 224 | // Reconfig |
| 225 | my_next_rsn: u32, |
| 226 | reconfigs: FxHashMap<u32, ChunkReconfig>, |
| 227 | /// Stream ids semantically covered by an outgoing reset. Reset-all stays |
| 228 | /// compact on the wire, so its completion ids cannot be recovered from the |
| 229 | /// serialized parameter itself. |
| 230 | reconfig_reset_streams: FxHashMap<u32, Vec<StreamId>>, |
| 231 | /// The one request whose Re-configuration Timer is running. |
| 232 | active_reconfig: Option<u32>, |
| 233 | /// Application-initiated resets waiting for the active request to finish. |
| 234 | pending_reset_streams: VecDeque<StreamId>, |
| 235 | reconfig_requests: FxHashMap<u32, ParamOutgoingResetRequest>, |
| 236 | /// DATA received above an InProgress reset boundary. A copy is retained |
| 237 | /// after cumulative TSN processing removes it from `payload_queue`, so it |
| 238 | /// can remain withheld until the old stream generation is drained. |
| 239 | deferred_reset_data: FxHashMap<u32, ChunkPayloadData>, |
| 240 | max_completed_reconfig_rsn: Option<u32>, |
| 241 | /// Re-configuration Request Sequence Number most recently accepted from |
| 242 | /// the peer, initialized to the peer's initial TSN minus one. |
| 243 | peer_last_reconfig_rsn: u32, |
| 244 | /// Whether `peer_last_reconfig_rsn` has been initialized from an INIT or a |
| 245 | /// first request in test/manual construction. |
| 246 | peer_reconfig_rsn_initialized: bool, |
| 247 | /// Stream ids for which `StreamEvent::Finished` is queued but whose terminal |
| 248 | /// reset event is still awaited. |
| 249 | pending_reset_completions: PendingResetCompletions, |
| 250 | /// Stream ids whose latest completed reset was unsuccessful. |
| 251 | failed_reset_streams: FxHashSet<StreamId>, |
| 252 | /// Reset boundaries for generations that cannot finish until an older |
| 253 | /// application-readable generation with the same stream id is drained. |
| 254 | retiring_streams: FxHashMap<StreamId, VecDeque<u32>>, |
| 255 | /// Number of queued Finished events whose generation has actually drained. |
| 256 | ready_stream_finishes: PendingResetCompletions, |
| 257 | /// Finished events already delivered to the application and therefore |
| 258 | /// eligible to be followed by one terminal reset event. |
| 259 | delivered_stream_finishes: PendingResetCompletions, |
| 260 | /// Forward-TSN updates that belong to a successor hidden behind a |
| 261 | /// still-readable stream generation. |
| 262 | deferred_forward_tsns: FxHashMap<StreamId, VecDeque<DeferredForwardTsn>>, |
| 263 | |
| 264 | // Non-RFC internal data |
| 265 | remote_addr: SocketAddr, |
| 266 | local_ip: Option<IpAddr>, |
| 267 | source_port: u16, |
| 268 | destination_port: u16, |
| 269 | my_max_num_inbound_streams: u16, |
| 270 | my_max_num_outbound_streams: u16, |
| 271 | my_cookie: Option<ParamStateCookie>, |
| 272 | |
| 273 | payload_queue: PayloadQueue, |
| 274 | inflight_queue: PayloadQueue, |
| 275 | pending_queue: PendingQueue, |
| 276 | control_queue: VecDeque<Packet>, |
| 277 | stream_queue: VecDeque<u16>, |
| 278 | |
| 279 | pub(crate) mtu: u32, |
| 280 | // max DATA chunk payload size |
| 281 | max_payload_size: u32, |
| 282 | cumulative_tsn_ack_point: u32, |
| 283 | advanced_peer_tsn_ack_point: u32, |
| 284 | use_forward_tsn: bool, |
| 285 | |
| 286 | pub(crate) rto_mgr: RtoManager, |
| 287 | timers: TimerTable, |
| 288 | |
| 289 | // Congestion control parameters |
| 290 | max_receive_buffer_size: u32, |
| 291 | // my congestion window size |
| 292 | pub(crate) cwnd: u32, |
| 293 | // calculated peer's receiver windows size |
| 294 | rwnd: u32, |
| 295 | // slow start threshold |
| 296 | pub(crate) ssthresh: u32, |
| 297 | partial_bytes_acked: u32, |
| 298 | pub(crate) in_fast_recovery: bool, |
| 299 | fast_recover_exit_point: u32, |
| 300 | |
| 301 | // Chunks stored for retransmission |
| 302 | stored_init: Option<ChunkInit>, |
| 303 | stored_cookie_echo: Option<ChunkCookieEcho>, |
| 304 | pub(crate) streams: FxHashMap<StreamId, StreamState>, |
| 305 | |
| 306 | events: VecDeque<Event>, |
| 307 | endpoint_events: VecDeque<EndpointEventInner>, |
| 308 | error: Option<AssociationError>, |
| 309 | |
| 310 | // per inbound packet context |
| 311 | delayed_ack_triggered: bool, |
| 312 | immediate_ack_triggered: bool, |
| 313 | |
| 314 | pub(crate) stats: AssociationStats, |
| 315 | ack_state: AckState, |
| 316 | |
| 317 | // for testing |
| 318 | pub(crate) ack_mode: AckMode, |
| 319 | } |
| 320 | |
| 321 | impl Default for Association { |
| 322 | fn default() -> Self { |
| 323 | Association { |
| 324 | side: Side::default(), |
| 325 | state: AssociationState::default(), |
| 326 | handshake_completed: false, |
| 327 | max_send_message_size: 0, |
| 328 | max_receive_message_size: 0, |
| 329 | receive_limits: None, |
| 330 | inflight_queue_length: 0, |
| 331 | will_send_shutdown: false, |
| 332 | bytes_received: 0, |
| 333 | bytes_sent: 0, |
| 334 | |
| 335 | peer_verification_tag: 0, |
| 336 | my_verification_tag: 0, |
| 337 | my_next_tsn: 0, |
| 338 | peer_last_tsn: 0, |
| 339 | // for RTT measurement |
| 340 | min_tsn2measure_rtt: 0, |
| 341 | will_send_forward_tsn: false, |
| 342 | will_retransmit_fast: false, |
| 343 | will_retransmit_reconfig: false, |
| 344 | |
| 345 | will_send_shutdown_ack: false, |
| 346 | will_send_shutdown_complete: false, |
| 347 | |
| 348 | // Reconfig |
| 349 | my_next_rsn: 0, |
| 350 | reconfigs: FxHashMap::default(), |
| 351 | reconfig_reset_streams: FxHashMap::default(), |
| 352 | active_reconfig: None, |
| 353 | pending_reset_streams: VecDeque::default(), |
| 354 | reconfig_requests: FxHashMap::default(), |
| 355 | deferred_reset_data: FxHashMap::default(), |
| 356 | max_completed_reconfig_rsn: None, |
| 357 | peer_last_reconfig_rsn: 0, |
| 358 | peer_reconfig_rsn_initialized: false, |
| 359 | pending_reset_completions: PendingResetCompletions::default(), |
| 360 | failed_reset_streams: FxHashSet::default(), |
| 361 | retiring_streams: FxHashMap::default(), |
| 362 | ready_stream_finishes: PendingResetCompletions::default(), |
| 363 | delivered_stream_finishes: PendingResetCompletions::default(), |
| 364 | deferred_forward_tsns: FxHashMap::default(), |
| 365 | |
| 366 | // Non-RFC internal data |
| 367 | remote_addr: SocketAddr::from_str("0.0.0.0:0").unwrap(), |
| 368 | local_ip: None, |
| 369 | source_port: 0, |
| 370 | destination_port: 0, |
| 371 | my_max_num_inbound_streams: 0, |
| 372 | my_max_num_outbound_streams: 0, |
| 373 | my_cookie: None, |
| 374 | |
| 375 | payload_queue: PayloadQueue::default(), |
| 376 | inflight_queue: PayloadQueue::default(), |
| 377 | pending_queue: PendingQueue::default(), |
| 378 | control_queue: VecDeque::default(), |
| 379 | stream_queue: VecDeque::default(), |
| 380 | |
| 381 | mtu: 0, |
| 382 | // max DATA chunk payload size |
| 383 | max_payload_size: 0, |
| 384 | cumulative_tsn_ack_point: 0, |
| 385 | advanced_peer_tsn_ack_point: 0, |
| 386 | use_forward_tsn: false, |
| 387 | |
| 388 | rto_mgr: RtoManager::default(), |
| 389 | timers: TimerTable::default(), |
| 390 | |
| 391 | // Congestion control parameters |
| 392 | max_receive_buffer_size: 0, |
| 393 | // my congestion window size |
| 394 | cwnd: 0, |
| 395 | // calculated peer's receiver windows size |
| 396 | rwnd: 0, |
| 397 | // slow start threshold |
| 398 | ssthresh: 0, |
| 399 | partial_bytes_acked: 0, |
| 400 | in_fast_recovery: false, |
| 401 | fast_recover_exit_point: 0, |
| 402 | |
| 403 | // Chunks stored for retransmission |
| 404 | stored_init: None, |
| 405 | stored_cookie_echo: None, |
| 406 | streams: FxHashMap::default(), |
| 407 | |
| 408 | events: VecDeque::default(), |
| 409 | endpoint_events: VecDeque::default(), |
| 410 | error: None, |
| 411 | |
| 412 | // per inbound packet context |
| 413 | delayed_ack_triggered: false, |
| 414 | immediate_ack_triggered: false, |
| 415 | |
| 416 | stats: AssociationStats::default(), |
| 417 | ack_state: AckState::default(), |
| 418 | |
| 419 | // for testing |
| 420 | ack_mode: AckMode::default(), |
| 421 | } |
| 422 | } |
| 423 | } |
| 424 | |
| 425 | impl Association { |
| 426 | fn new_common( |
| 427 | config: Arc<TransportConfig>, |
| 428 | max_payload_size: u32, |
| 429 | remote_addr: SocketAddr, |
| 430 | local_ip: Option<IpAddr>, |
| 431 | side: Side, |
| 432 | verification_tag: u32, |
| 433 | initial_tsn: u32, |
| 434 | ) -> Self { |
| 435 | // It's a bit strange, but we're going backwards from the calculation in |
| 436 | // config.rs to get max_payload_size from INITIAL_MTU. |
| 437 | let mtu = max_payload_size + COMMON_HEADER_SIZE + DATA_CHUNK_HEADER_SIZE; |
| 438 | |
| 439 | // RFC 4960 Sec 7.2.1 |
| 440 | // The initial cwnd before DATA transmission or after a sufficiently |
| 441 | // long idle period MUST be set to min(4*MTU, max (2*MTU, 4380bytes)). |
| 442 | let cwnd = (4 * mtu).min((2 * mtu).max(4380)); |
| 443 | |
| 444 | Association { |
| 445 | side, |
| 446 | handshake_completed: false, |
| 447 | max_receive_buffer_size: config.max_receive_buffer_size(), |
| 448 | max_send_message_size: config.max_send_message_size(), |
| 449 | max_receive_message_size: config.max_receive_message_size(), |
| 450 | receive_limits: config.receive_limits(), |
| 451 | my_max_num_outbound_streams: config.max_num_outbound_streams(), |
| 452 | my_max_num_inbound_streams: config.max_num_inbound_streams(), |
| 453 | max_payload_size, |
| 454 | |
| 455 | rto_mgr: RtoManager::new( |
| 456 | config.rto_initial_ms(), |
| 457 | config.rto_min_ms(), |
| 458 | config.rto_max_ms(), |
| 459 | ), |
| 460 | timers: TimerTable::new( |
| 461 | config.max_init_retransmits(), |
| 462 | config.max_data_retransmits(), |
| 463 | config.rto_max_ms(), |
| 464 | ), |
| 465 | |
| 466 | mtu, |
| 467 | cwnd, |
| 468 | remote_addr, |
| 469 | local_ip, |
| 470 | |
| 471 | my_verification_tag: verification_tag, |
| 472 | my_next_tsn: initial_tsn, |
| 473 | my_next_rsn: initial_tsn, |
| 474 | min_tsn2measure_rtt: initial_tsn, |
| 475 | cumulative_tsn_ack_point: initial_tsn.wrapping_sub(1), |
| 476 | advanced_peer_tsn_ack_point: initial_tsn.wrapping_sub(1), |
| 477 | error: None, |
| 478 | |
| 479 | ..Default::default() |
| 480 | } |
| 481 | } |
| 482 | |
| 483 | pub(crate) fn new( |
| 484 | server_config: Option<Arc<ServerConfig>>, |
| 485 | config: Arc<TransportConfig>, |
| 486 | max_payload_size: u32, |
| 487 | local_aid: AssociationId, |
| 488 | remote_addr: SocketAddr, |
| 489 | local_ip: Option<IpAddr>, |
| 490 | now: Instant, |
| 491 | ) -> Self { |
| 492 | let side = if server_config.is_some() { |
| 493 | Side::Server |
| 494 | } else { |
| 495 | Side::Client |
| 496 | }; |
| 497 | |
| 498 | let tsn = random::<NonZeroU32>().get(); |
| 499 | |
| 500 | let mut this = Self::new_common( |
| 501 | config, |
| 502 | max_payload_size, |
| 503 | remote_addr, |
| 504 | local_ip, |
| 505 | side, |
| 506 | local_aid, |
| 507 | tsn, |
| 508 | ); |
| 509 | |
| 510 | this.source_port = DEFAULT_SCTP_PORT; |
| 511 | this.destination_port = DEFAULT_SCTP_PORT; |
| 512 | |
| 513 | if side.is_client() { |
| 514 | let mut init = ChunkInit { |
| 515 | initial_tsn: this.my_next_tsn, |
| 516 | num_outbound_streams: this.my_max_num_outbound_streams, |
| 517 | num_inbound_streams: this.my_max_num_inbound_streams, |
| 518 | initiate_tag: this.my_verification_tag, |
| 519 | advertised_receiver_window_credit: this.max_receive_buffer_size, |
| 520 | ..Default::default() |
| 521 | }; |
| 522 | init.set_supported_extensions(); |
| 523 | |
| 524 | this.set_state(AssociationState::CookieWait); |
| 525 | this.stored_init = Some(init); |
| 526 | let _ = this.send_init(); |
| 527 | this.timers |
| 528 | .start(Timer::T1Init, now, this.rto_mgr.get_rto()); |
| 529 | } |
| 530 | |
| 531 | this |
| 532 | } |
| 533 | |
| 534 | /// Creates a new association using out-of-band exchanged SNAP tokens (INIT chunks). |
| 535 | /// |
| 536 | /// This allows skipping the SCTP 4-way handshake (RFC 4960 Section 5.1) |
| 537 | /// by exchanging tokens out-of-band (e.g., via a signaling channel |
| 538 | /// using SDP `a=sctp-init`). The association immediately transitions to |
| 539 | /// the ESTABLISHED state. |
| 540 | /// |
| 541 | /// **Note:** When using SNAP, **both** peers must call |
| 542 | /// [`Endpoint::connect`](crate::Endpoint::connect). There is no |
| 543 | /// server-side SNAP via [`Endpoint::handle`](crate::Endpoint::handle). |
| 544 | /// |
| 545 | /// See [draft-hancke-tsvwg-snap-01](https://datatracker.ietf.org/doc/draft-hancke-tsvwg-snap/). |
| 546 | /// |
| 547 | /// # Arguments |
| 548 | /// * `config` - Transport configuration. |
| 549 | /// * `max_payload_size` - Maximum payload size. |
| 550 | /// * `remote_addr` - Remote socket address. |
| 551 | /// * `local_ip` - Optional local IP address. |
| 552 | /// * `local_init` - Parsed local token (INIT chunk). |
| 553 | /// * `remote_init` - Parsed remote token (INIT chunk). |
| 554 | /// |
| 555 | /// # Returns |
| 556 | /// A new association in the ESTABLISHED state, or an error if the |
| 557 | /// tokens are invalid. |
| 558 | pub(crate) fn new_with_out_of_band_init( |
| 559 | config: Arc<TransportConfig>, |
| 560 | max_payload_size: u32, |
| 561 | remote_addr: SocketAddr, |
| 562 | local_ip: Option<IpAddr>, |
| 563 | local_init: ChunkInit, |
| 564 | remote_init: ChunkInit, |
| 565 | ) -> Result<Self> { |
| 566 | // Derive side deterministically: the peer with the lower initiate_tag |
| 567 | // acts as server so that log lines are distinguishable. This is purely |
| 568 | // cosmetic โ equal tags are impossible because the caller in |
| 569 | // `connect_with_snap` rejects that case with `AidCollision`. |
| 570 | let side = if local_init.initiate_tag <= remote_init.initiate_tag { |
| 571 | Side::Server |
| 572 | } else { |
| 573 | Side::Client |
| 574 | }; |
| 575 | |
| 576 | // Use the TSN from our local INIT chunk |
| 577 | let tsn = local_init.initial_tsn; |
| 578 | |
| 579 | let mut this = Self::new_common( |
| 580 | config.clone(), |
| 581 | max_payload_size, |
| 582 | remote_addr, |
| 583 | local_ip, |
| 584 | side, |
| 585 | local_init.initiate_tag, |
| 586 | tsn, |
| 587 | ); |
| 588 | |
| 589 | // Negotiate stream counts: use the smaller of our config limit and |
| 590 | // what the remote offers (RFC 4960 ยง5.1.1 cross-negotiation). |
| 591 | this.my_max_num_inbound_streams = core::cmp::min( |
| 592 | config.max_num_inbound_streams(), |
| 593 | remote_init.num_outbound_streams, |
| 594 | ); |
| 595 | this.my_max_num_outbound_streams = core::cmp::min( |
| 596 | config.max_num_outbound_streams(), |
| 597 | remote_init.num_inbound_streams, |
| 598 | ); |
| 599 | |
| 600 | this.peer_verification_tag = remote_init.initiate_tag; |
| 601 | |
| 602 | this.source_port = DEFAULT_SCTP_PORT; |
| 603 | this.destination_port = DEFAULT_SCTP_PORT; |
| 604 | this.handshake_completed = true; |
| 605 | |
| 606 | this.apply_remote_init_params( |
| 607 | remote_init.initial_tsn, |
| 608 | remote_init.advertised_receiver_window_credit, |
| 609 | &remote_init.params, |
| 610 | "out-of-band init", |
| 611 | ); |
| 612 | |
| 613 | // Set state to ESTABLISHED - out-of-band init skips the handshake |
| 614 | this.set_state(AssociationState::Established); |
| 615 | this.events.push_back(Event::Connected); |
| 616 | |
| 617 | debug!( |
| 618 | "[{}] out-of-band init association established: my_tag={:#x} peer_tag={:#x} tsn={}", |
| 619 | this.side, this.my_verification_tag, this.peer_verification_tag, this.my_next_tsn |
| 620 | ); |
| 621 | |
| 622 | Ok(this) |
| 623 | } |
| 624 | |
| 625 | /// Returns application-facing event |
| 626 | /// |
| 627 | /// Associations should be polled for events after: |
| 628 | /// - a call was made to `handle_event` |
| 629 | /// - a call was made to `handle_timeout` |
| 630 | #[must_use] |
| 631 | pub fn poll(&mut self) -> Option<Event> { |
| 632 | // A reset boundary can become cumulative after its DATA was made |
| 633 | // readable. Keep events for that stream generation behind Finished, |
| 634 | // while allowing unrelated association and stream events to progress. |
| 635 | let mut blocked_streams = FxHashSet::default(); |
| 636 | let mut selected = None; |
| 637 | for (index, event) in self.events.iter().enumerate() { |
| 638 | if let Event::Stream(StreamEvent::Finished { id }) = event { |
| 639 | if self.state != AssociationState::Closed |
| 640 | && !self.ready_stream_finishes.contains(id) |
| 641 | { |
| 642 | blocked_streams.insert(*id); |
| 643 | continue; |
| 644 | } |
| 645 | } |
| 646 | |
| 647 | if let Event::Stream(stream_event) = event { |
| 648 | let stream_id = match stream_event { |
| 649 | StreamEvent::Opened { id } |
| 650 | | StreamEvent::Readable { id } |
| 651 | | StreamEvent::Writable { id } |
| 652 | | StreamEvent::Finished { id } |
| 653 | | StreamEvent::ResetComplete { id } |
| 654 | | StreamEvent::ResetFailed { id, .. } |
| 655 | | StreamEvent::Stopped { id, .. } |
| 656 | | StreamEvent::BufferedAmountLow { id } |
| 657 | | StreamEvent::BufferedAmountHigh { id } => Some(*id), |
| 658 | StreamEvent::Available => None, |
| 659 | }; |
| 660 | if let Some(stream_id) = stream_id { |
| 661 | if blocked_streams.contains(&stream_id) { |
| 662 | let terminal_for_delivered_generation = |
| 663 | matches!( |
| 664 | stream_event, |
| 665 | StreamEvent::ResetComplete { .. } | StreamEvent::ResetFailed { .. } |
| 666 | ) && self.delivered_stream_finishes.contains(&stream_id); |
| 667 | let readable_current_generation = matches!( |
| 668 | stream_event, |
| 669 | StreamEvent::Opened { .. } | StreamEvent::Readable { .. } |
| 670 | ); |
| 671 | if !terminal_for_delivered_generation && !readable_current_generation { |
| 672 | continue; |
| 673 | } |
| 674 | } |
| 675 | } |
| 676 | } |
| 677 | |
| 678 | selected = Some(index); |
| 679 | break; |
| 680 | } |
| 681 | |
| 682 | if let Some(index) = selected { |
| 683 | let x = self |
| 684 | .events |
| 685 | .remove(index) |
| 686 | .expect("selected event must exist"); |
| 687 | match &x { |
| 688 | Event::Stream(StreamEvent::Finished { id }) => { |
| 689 | self.ready_stream_finishes.take_one(*id); |
| 690 | if self.state != AssociationState::Closed { |
| 691 | self.delivered_stream_finishes.insert(*id); |
| 692 | } |
| 693 | } |
| 694 | Event::Stream( |
| 695 | StreamEvent::ResetComplete { id } | StreamEvent::ResetFailed { id, .. }, |
| 696 | ) => { |
| 697 | self.delivered_stream_finishes.take_one(*id); |
| 698 | } |
| 699 | _ => {} |
| 700 | } |
| 701 | return Some(x); |
| 702 | } |
| 703 | |
| 704 | /*TODO: if let Some(event) = self.streams.poll() { |
| 705 | return Some(Event::Stream(event)); |
| 706 | }*/ |
| 707 | |
| 708 | if let Some(err) = self.error.take() { |
| 709 | return Some(Event::HandshakeFailed { reason: err }); |
| 710 | } |
| 711 | |
| 712 | None |
| 713 | } |
| 714 | |
| 715 | /// Return endpoint-facing event |
| 716 | #[must_use] |
| 717 | pub fn poll_endpoint_event(&mut self) -> Option<EndpointEvent> { |
| 718 | self.endpoint_events.pop_front().map(EndpointEvent) |
| 719 | } |
| 720 | |
| 721 | /// Returns the next time at which `handle_timeout` should be called |
| 722 | /// |
| 723 | /// The value returned may change after: |
| 724 | /// - the application performed some I/O on the association |
| 725 | /// - a call was made to `handle_transmit` |
| 726 | /// - a call to `poll_transmit` returned `Some` |
| 727 | /// - a call was made to `handle_timeout` |
| 728 | #[must_use] |
| 729 | pub fn poll_timeout(&self) -> Option<Instant> { |
| 730 | self.timers.next_timeout() |
| 731 | } |
| 732 | |
| 733 | /// Returns packets to transmit |
| 734 | /// |
| 735 | /// Associations should be polled for transmit after: |
| 736 | /// - the application performed some I/O on the Association |
| 737 | /// - a call was made to `handle_event` |
| 738 | /// - a call was made to `handle_timeout` |
| 739 | #[must_use] |
| 740 | pub fn poll_transmit(&mut self, now: Instant) -> Option<Transmit> { |
| 741 | let (contents, _) = self.gather_outbound(now); |
| 742 | if contents.is_empty() { |
| 743 | None |
| 744 | } else { |
| 745 | trace!( |
| 746 | "[{}] sending {} bytes (total {} datagrams)", |
| 747 | self.side, |
| 748 | contents.iter().fold(0, |l, c| l + c.len()), |
| 749 | contents.len() |
| 750 | ); |
| 751 | Some(Transmit { |
| 752 | now, |
| 753 | remote: self.remote_addr, |
| 754 | payload: Payload::RawEncode(contents), |
| 755 | ecn: None, |
| 756 | local_ip: self.local_ip, |
| 757 | }) |
| 758 | } |
| 759 | } |
| 760 | |
| 761 | /// Process timer expirations |
| 762 | /// |
| 763 | /// Executes protocol logic, potentially preparing signals (including application `Event`s, |
| 764 | /// `EndpointEvent`s and outgoing datagrams) that should be extracted through the relevant |
| 765 | /// methods. |
| 766 | /// |
| 767 | /// It is most efficient to call this immediately after the system clock reaches the latest |
| 768 | /// `Instant` that was output by `poll_timeout`; however spurious extra calls will simply |
| 769 | /// no-op and therefore are safe. |
| 770 | pub fn handle_timeout(&mut self, now: Instant) { |
| 771 | for &timer in &Timer::VALUES { |
| 772 | let (expired, failure, n_rtos) = self.timers.is_expired(timer, now); |
| 773 | if !expired { |
| 774 | continue; |
| 775 | } |
| 776 | self.timers.set(timer, None); |
| 777 | |
| 778 | if timer == Timer::Ack { |
| 779 | self.on_ack_timeout(); |
| 780 | } else if failure { |
| 781 | self.on_retransmission_failure(timer); |
| 782 | } else { |
| 783 | self.on_retransmission_timeout(timer, n_rtos); |
| 784 | self.timers.start(timer, now, self.rto_mgr.get_rto()); |
| 785 | } |
| 786 | } |
| 787 | } |
| 788 | |
| 789 | /// Process `AssociationEvent`s generated by the associated `Endpoint` |
| 790 | /// |
| 791 | /// Will execute protocol logic upon receipt of an association event, in turn preparing signals |
| 792 | /// (including application `Event`s, `EndpointEvent`s and outgoing datagrams) that should be |
| 793 | /// extracted through the relevant methods. |
| 794 | pub fn handle_event(&mut self, event: AssociationEvent) { |
| 795 | match event.0 { |
| 796 | AssociationEventInner::Datagram(transmit) => { |
| 797 | // If this packet could initiate a migration and we're a client or a server that |
| 798 | // forbids migration, drop the datagram. This could be relaxed to heuristically |
| 799 | // permit NAT-rebinding-like migration. |
| 800 | /*TODO:if remote != self.remote && self.server_config.as_ref().map_or(true, |x| !x.migration) |
| 801 | { |
| 802 | trace!("discarding packet from unrecognized peer {}", remote); |
| 803 | return; |
| 804 | }*/ |
| 805 | |
| 806 | if let Payload::PartialDecode(partial_decode) = transmit.payload { |
| 807 | trace!( |
| 808 | "[{}] receiving {} bytes", |
| 809 | self.side, |
| 810 | COMMON_HEADER_SIZE as usize + partial_decode.remaining.len() |
| 811 | ); |
| 812 | |
| 813 | let pkt = match partial_decode.finish() { |
| 814 | Ok(p) => p, |
| 815 | Err(err) => { |
| 816 | warn!("[{}] unable to parse SCTP packet {}", self.side, err); |
| 817 | return; |
| 818 | } |
| 819 | }; |
| 820 | |
| 821 | if let Err(err) = self.handle_inbound(pkt, transmit.now) { |
| 822 | error!("handle_inbound got err: {}", err); |
| 823 | let _ = self.close(); |
| 824 | } |
| 825 | } else { |
| 826 | trace!("discarding invalid partial_decode"); |
| 827 | } |
| 828 | } //TODO: |
| 829 | } |
| 830 | } |
| 831 | |
| 832 | /// Returns Association statistics |
| 833 | pub fn stats(&self) -> AssociationStats { |
| 834 | self.stats |
| 835 | } |
| 836 | |
| 837 | /// Whether the Association is in the process of being established |
| 838 | /// |
| 839 | /// If this returns `false`, the Association may be either established or closed, signaled by the |
| 840 | /// emission of a `Connected` or `AssociationLost` message respectively. |
| 841 | pub fn is_handshaking(&self) -> bool { |
| 842 | !self.handshake_completed |
| 843 | } |
| 844 | |
| 845 | /// Whether the Association is closed |
| 846 | /// |
| 847 | /// Closed Associations cannot transport any further data. An association becomes closed when |
| 848 | /// either peer application intentionally closes it, or when either transport layer detects an |
| 849 | /// error such as a time-out or certificate validation failure. |
| 850 | /// |
| 851 | /// A `AssociationLost` event is emitted with details when the association becomes closed. |
| 852 | pub fn is_closed(&self) -> bool { |
| 853 | self.state == AssociationState::Closed |
| 854 | } |
| 855 | |
| 856 | /// Whether the Association has started SCTP shutdown, but is not closed yet |
| 857 | /// |
| 858 | /// Closing Associations may still need polling, timer-driven retransmission, and packet output |
| 859 | /// before they become fully closed. |
| 860 | pub fn is_closing(&self) -> bool { |
| 861 | self.state.is_closing() |
| 862 | } |
| 863 | |
| 864 | /// Whether there is no longer any need to keep the association around |
| 865 | /// |
| 866 | /// Closed associations become drained after a brief timeout to absorb any remaining in-flight |
| 867 | /// packets from the peer. All drained associations have been closed. |
| 868 | pub fn is_drained(&self) -> bool { |
| 869 | self.state.is_drained() |
| 870 | } |
| 871 | |
| 872 | /// Look up whether we're the client or server of this Association |
| 873 | pub fn side(&self) -> Side { |
| 874 | self.side |
| 875 | } |
| 876 | |
| 877 | /// The latest socket address for this Association's peer |
| 878 | pub fn remote_addr(&self) -> SocketAddr { |
| 879 | self.remote_addr |
| 880 | } |
| 881 | |
| 882 | /// Current best estimate of this Association's latency (round-trip-time) |
| 883 | pub fn rtt(&self) -> Duration { |
| 884 | Duration::from_millis(self.rto_mgr.get_rto()) |
| 885 | } |
| 886 | |
| 887 | /// The local IP address which was used when the peer established |
| 888 | /// the association |
| 889 | /// |
| 890 | /// This can be different from the address the endpoint is bound to, in case |
| 891 | /// the endpoint is bound to a wildcard address like `0.0.0.0` or `::`. |
| 892 | /// |
| 893 | /// This will return `None` for clients. |
| 894 | /// |
| 895 | /// Retrieving the local IP address is currently supported on the following |
| 896 | /// platforms: |
| 897 | /// - Linux |
| 898 | /// |
| 899 | /// On all non-supported platforms the local IP address will not be available, |
| 900 | /// and the method will return `None`. |
| 901 | pub fn local_ip(&self) -> Option<IpAddr> { |
| 902 | self.local_ip |
| 903 | } |
| 904 | |
| 905 | /// Shutdown initiates the shutdown sequence. The method blocks until the |
| 906 | /// shutdown sequence is completed and the association is closed, or until the |
| 907 | /// passed context is done, in which case the context's error is returned. |
| 908 | pub fn shutdown(&mut self) -> Result<()> { |
| 909 | debug!("[{}] closing association..", self.side); |
| 910 | |
| 911 | let state = self.state(); |
| 912 | if state != AssociationState::Established { |
| 913 | return Err(Error::ErrShutdownNonEstablished); |
| 914 | } |
| 915 | |
| 916 | // Attempt a graceful shutdown. |
| 917 | self.set_state(AssociationState::ShutdownPending); |
| 918 | |
| 919 | if self.inflight_queue_length == 0 { |
| 920 | // No more outstanding, send shutdown. |
| 921 | self.will_send_shutdown = true; |
| 922 | self.awake_write_loop(); |
| 923 | self.set_state(AssociationState::ShutdownSent); |
| 924 | } |
| 925 | |
| 926 | self.endpoint_events.push_back(EndpointEventInner::Drained); |
| 927 | |
| 928 | Ok(()) |
| 929 | } |
| 930 | |
| 931 | /// Close ends the SCTP Association and cleans up any state |
| 932 | pub fn close(&mut self) -> Result<()> { |
| 933 | if self.state() != AssociationState::Closed { |
| 934 | self.set_state(AssociationState::Closed); |
| 935 | |
| 936 | debug!("[{}] closing association..", self.side); |
| 937 | |
| 938 | self.close_all_timers(); |
| 939 | |
| 940 | for si in self.streams.keys().cloned().collect::<Vec<u16>>() { |
| 941 | self.unregister_stream(si, false); |
| 942 | } |
| 943 | |
| 944 | // AssociationLost stops any pending ResetComplete. |
| 945 | self.pending_reset_completions.clear(); |
| 946 | self.failed_reset_streams.clear(); |
| 947 | self.retiring_streams.clear(); |
| 948 | self.ready_stream_finishes.clear(); |
| 949 | self.delivered_stream_finishes.clear(); |
| 950 | self.deferred_forward_tsns.clear(); |
| 951 | self.pending_reset_streams.clear(); |
| 952 | self.reconfigs.clear(); |
| 953 | self.reconfig_reset_streams.clear(); |
| 954 | self.reconfig_requests.clear(); |
| 955 | self.deferred_reset_data.clear(); |
| 956 | self.control_queue.clear(); |
| 957 | self.active_reconfig = None; |
| 958 | self.will_retransmit_reconfig = false; |
| 959 | |
| 960 | self.events.push_back(Event::AssociationLost { |
| 961 | reason: AssociationError::AssociationClosed, |
| 962 | }); |
| 963 | |
| 964 | debug!("[{}] association closed", self.side); |
| 965 | debug!( |
| 966 | "[{}] stats nDATAs (in) : {}", |
| 967 | self.side, |
| 968 | self.stats.get_num_datas() |
| 969 | ); |
| 970 | debug!( |
| 971 | "[{}] stats nSACKs (in) : {}", |
| 972 | self.side, |
| 973 | self.stats.get_num_sacks() |
| 974 | ); |
| 975 | debug!( |
| 976 | "[{}] stats nT3Timeouts : {}", |
| 977 | self.side, |
| 978 | self.stats.get_num_t3timeouts() |
| 979 | ); |
| 980 | debug!( |
| 981 | "[{}] stats nAckTimeouts: {}", |
| 982 | self.side, |
| 983 | self.stats.get_num_ack_timeouts() |
| 984 | ); |
| 985 | debug!( |
| 986 | "[{}] stats nFastRetrans: {}", |
| 987 | self.side, |
| 988 | self.stats.get_num_fast_retrans() |
| 989 | ); |
| 990 | } |
| 991 | |
| 992 | Ok(()) |
| 993 | } |
| 994 | |
| 995 | /// open_stream opens a stream |
| 996 | pub fn open_stream( |
| 997 | &mut self, |
| 998 | stream_identifier: StreamId, |
| 999 | default_payload_type: PayloadProtocolIdentifier, |
| 1000 | ) -> Result<Stream<'_>> { |
| 1001 | if self.streams.contains_key(&stream_identifier) { |
| 1002 | return Err(Error::ErrStreamAlreadyExist); |
| 1003 | } |
| 1004 | |
| 1005 | if self.stream_reset_blocked(stream_identifier) { |
| 1006 | return Err(Error::ErrStreamResetPending); |
| 1007 | } |
| 1008 | |
| 1009 | if let Some(s) = self.create_stream(stream_identifier, false, default_payload_type) { |
| 1010 | Ok(s) |
| 1011 | } else { |
| 1012 | Err(Error::ErrStreamCreateFailed) |
| 1013 | } |
| 1014 | } |
| 1015 | |
| 1016 | /// accept_stream accepts a stream |
| 1017 | pub fn accept_stream(&mut self) -> Option<Stream<'_>> { |
| 1018 | self.stream_queue |
| 1019 | .pop_front() |
| 1020 | .map(move |stream_identifier| Stream { |
| 1021 | stream_identifier, |
| 1022 | association: self, |
| 1023 | }) |
| 1024 | } |
| 1025 | |
| 1026 | /// stream returns a stream |
| 1027 | pub fn stream(&mut self, stream_identifier: StreamId) -> Result<Stream<'_>> { |
| 1028 | if !self.streams.contains_key(&stream_identifier) { |
| 1029 | Err(Error::ErrStreamNotExisted) |
| 1030 | } else { |
| 1031 | Ok(Stream { |
| 1032 | stream_identifier, |
| 1033 | association: self, |
| 1034 | }) |
| 1035 | } |
| 1036 | } |
| 1037 | |
| 1038 | /// stream_ids returns a list of all active stream identifiers |
| 1039 | pub fn stream_ids(&self) -> Vec<StreamId> { |
| 1040 | self.streams.keys().cloned().collect() |
| 1041 | } |
| 1042 | |
| 1043 | /// bytes_sent returns the number of bytes sent |
| 1044 | pub(crate) fn bytes_sent(&self) -> usize { |
| 1045 | self.bytes_sent |
| 1046 | } |
| 1047 | |
| 1048 | /// bytes_received returns the number of bytes received |
| 1049 | pub(crate) fn bytes_received(&self) -> usize { |
| 1050 | self.bytes_received |
| 1051 | } |
| 1052 | |
| 1053 | /// max_send_message_size returns the maximum message size you can send. |
| 1054 | pub(crate) fn max_send_message_size(&self) -> u32 { |
| 1055 | self.max_send_message_size |
| 1056 | } |
| 1057 | |
| 1058 | /// set_max_send_message_size sets the maximum message size you can send. |
| 1059 | pub fn set_max_send_message_size(&mut self, value: u32) { |
| 1060 | self.max_send_message_size = value; |
| 1061 | } |
| 1062 | |
| 1063 | /// max_receive_message_size returns the maximum message size accepted. |
| 1064 | pub(crate) fn max_receive_message_size(&self) -> u32 { |
| 1065 | self.max_receive_message_size |
| 1066 | } |
| 1067 | |
| 1068 | /// set_max_receive_message_size sets the maximum message size accepted. |
| 1069 | pub(crate) fn set_max_receive_message_size(&mut self, value: u32) { |
| 1070 | self.max_receive_message_size = value; |
| 1071 | } |
| 1072 | |
| 1073 | /// max_message_size returns the maximum message size you can send. |
| 1074 | #[deprecated(note = "Use max_send_message_size instead")] |
| 1075 | pub(crate) fn max_message_size(&self) -> u32 { |
| 1076 | self.max_send_message_size() |
| 1077 | } |
| 1078 | |
| 1079 | /// set_max_message_size sets the maximum message size you can send. |
| 1080 | #[deprecated(note = "Use set_max_send_message_size instead")] |
| 1081 | pub(crate) fn set_max_message_size(&mut self, value: u32) { |
| 1082 | self.set_max_send_message_size(value) |
| 1083 | } |
| 1084 | |
| 1085 | /// Push one [`StreamEvent::ResetComplete`] for each candidate id whose reset |
| 1086 | /// handshake has completed and `Finished` has already been fired for it. |
| 1087 | fn emit_reset_complete(&mut self, ids: impl IntoIterator<Item = StreamId>) { |
| 1088 | for id in ids { |
| 1089 | let notify = self.pending_reset_completions.take_one(id); |
| 1090 | // Requests are serialized, so a success supersedes any older failure |
| 1091 | // for the same stream id. A newer pending generation still blocks reuse. |
| 1092 | self.failed_reset_streams.remove(&id); |
| 1093 | let reset_blocked = self.stream_reset_blocked(id); |
| 1094 | if !reset_blocked { |
| 1095 | let application_closed = notify |
| 1096 | && self.streams.get(&id).is_some_and(|stream| { |
| 1097 | matches!( |
| 1098 | stream.state, |
| 1099 | RecvSendState::Closed | RecvSendState::Writable |
| 1100 | ) |
| 1101 | }); |
| 1102 | if application_closed { |
| 1103 | self.unregister_stream(id, false); |
| 1104 | } else if let Some(stream) = self.streams.get_mut(&id) { |
| 1105 | // RFC 6525 section 5.2.7 H4 resets the outgoing SSN. Do not |
| 1106 | // alter `state`: a protocol reset must never undo an |
| 1107 | // application-level `Stream::finish()`. |
| 1108 | stream.sequence_number = 0; |
| 1109 | } |
| 1110 | } |
| 1111 | if notify { |
| 1112 | self.events |
| 1113 | .push_back(Event::Stream(StreamEvent::ResetComplete { id })); |
| 1114 | } |
| 1115 | } |
| 1116 | } |
| 1117 | |
| 1118 | /// Report a terminal reset failure and keep the stream id quarantined. |
| 1119 | fn emit_reset_failed( |
| 1120 | &mut self, |
| 1121 | ids: impl IntoIterator<Item = StreamId>, |
| 1122 | reason: StreamResetError, |
| 1123 | ) { |
| 1124 | for id in ids { |
| 1125 | self.pending_reset_completions.take_one(id); |
| 1126 | self.failed_reset_streams.insert(id); |
| 1127 | self.events |
| 1128 | .push_back(Event::Stream(StreamEvent::ResetFailed { id, reason })); |
| 1129 | } |
| 1130 | } |
| 1131 | |
| 1132 | /// Returns true if the given stream ID appears in any pending outgoing |
| 1133 | /// RE-CONFIG that has not yet been acknowledged by the remote peer. |
| 1134 | fn has_pending_reset_for_stream(&self, stream_id: StreamId) -> bool { |
| 1135 | self.reconfigs.values().any(|c| { |
| 1136 | c.param_a |
| 1137 | .iter() |
| 1138 | .chain(c.param_b.iter()) |
| 1139 | .find_map(|p| p.as_any().downcast_ref::<ParamOutgoingResetRequest>()) |
| 1140 | .is_some_and(|p| Self::reset_request_affects_stream(p, stream_id)) |
| 1141 | }) |
| 1142 | } |
| 1143 | |
| 1144 | fn reset_request_affects_stream( |
| 1145 | request: &ParamOutgoingResetRequest, |
| 1146 | stream_id: StreamId, |
| 1147 | ) -> bool { |
| 1148 | request.stream_identifiers.is_empty() || request.stream_identifiers.contains(&stream_id) |
| 1149 | } |
| 1150 | |
| 1151 | /// Whether the current generation of a stream id is unsafe to send or reuse. |
| 1152 | fn stream_reset_blocked(&self, stream_id: StreamId) -> bool { |
| 1153 | self.pending_reset_completions.contains(&stream_id) |
| 1154 | || self.failed_reset_streams.contains(&stream_id) |
| 1155 | || self.retiring_streams.contains_key(&stream_id) |
| 1156 | || self.stream_reset_in_progress(stream_id) |
| 1157 | } |
| 1158 | |
| 1159 | /// Whether a request currently prevents assigning new SSNs for this stream. |
| 1160 | fn stream_reset_in_progress(&self, stream_id: StreamId) -> bool { |
| 1161 | self.pending_reset_streams.contains(&stream_id) |
| 1162 | || self.has_pending_reset_for_stream(stream_id) |
| 1163 | } |
| 1164 | |
| 1165 | fn reconfig_stream_ids(c: &ChunkReconfig) -> Vec<StreamId> { |
| 1166 | c.param_a |
| 1167 | .iter() |
| 1168 | .chain(c.param_b.iter()) |
| 1169 | .find_map(|p| p.as_any().downcast_ref::<ParamOutgoingResetRequest>()) |
| 1170 | .map(|p| p.stream_identifiers.clone()) |
| 1171 | .unwrap_or_default() |
| 1172 | } |
| 1173 | |
| 1174 | fn packet_reconfig_request_rsn(packet: &Packet) -> Option<u32> { |
| 1175 | packet.chunks.iter().find_map(|chunk| { |
| 1176 | let reconfig = chunk.as_any().downcast_ref::<ChunkReconfig>()?; |
| 1177 | reconfig |
| 1178 | .param_a |
| 1179 | .iter() |
| 1180 | .chain(reconfig.param_b.iter()) |
| 1181 | .find_map(|param| { |
| 1182 | param |
| 1183 | .as_any() |
| 1184 | .downcast_ref::<ParamOutgoingResetRequest>() |
| 1185 | .map(|request| request.reconfig_request_sequence_number) |
| 1186 | }) |
| 1187 | }) |
| 1188 | } |
| 1189 | |
| 1190 | fn reconfig_with_sender_last_tsn(c: &ChunkReconfig, sender_last_tsn: u32) -> ChunkReconfig { |
| 1191 | fn update( |
| 1192 | param: &Option<Box<dyn Param + Send + Sync>>, |
| 1193 | sender_last_tsn: u32, |
| 1194 | ) -> Option<Box<dyn Param + Send + Sync>> { |
| 1195 | param.as_ref().map(|param| { |
| 1196 | if let Some(request) = param.as_any().downcast_ref::<ParamOutgoingResetRequest>() { |
| 1197 | let mut request = request.clone(); |
| 1198 | request.sender_last_tsn = sender_last_tsn; |
| 1199 | Box::new(request) as Box<dyn Param + Send + Sync> |
| 1200 | } else { |
| 1201 | param.clone() |
| 1202 | } |
| 1203 | }) |
| 1204 | } |
| 1205 | |
| 1206 | ChunkReconfig { |
| 1207 | param_a: update(&c.param_a, sender_last_tsn), |
| 1208 | param_b: update(&c.param_b, sender_last_tsn), |
| 1209 | } |
| 1210 | } |
| 1211 | |
| 1212 | fn reconfig_has_pending_data(&self, rsn: u32) -> bool { |
| 1213 | self.reconfigs.get(&rsn).is_some_and(|c| { |
| 1214 | let stream_ids = Self::reconfig_stream_ids(c); |
| 1215 | if stream_ids.is_empty() { |
| 1216 | !self.pending_queue.is_empty() |
| 1217 | } else { |
| 1218 | stream_ids |
| 1219 | .iter() |
| 1220 | .any(|id| self.pending_queue.contains_stream(*id)) |
| 1221 | } |
| 1222 | }) |
| 1223 | } |
| 1224 | |
| 1225 | fn refresh_unsent_reconfig(&mut self, rsn: u32) -> Option<Packet> { |
| 1226 | let sender_last_tsn = self.my_next_tsn.wrapping_sub(1); |
| 1227 | let reconfig = self |
| 1228 | .reconfigs |
| 1229 | .get(&rsn) |
| 1230 | .map(|c| Self::reconfig_with_sender_last_tsn(c, sender_last_tsn))?; |
| 1231 | self.reconfigs.insert(rsn, reconfig.clone()); |
| 1232 | Some(self.create_packet(vec![Box::new(reconfig)])) |
| 1233 | } |
| 1234 | |
| 1235 | /// Finish the one request associated with the running Re-configuration Timer. |
| 1236 | fn finish_reconfig( |
| 1237 | &mut self, |
| 1238 | rsn: u32, |
| 1239 | outcome: core::result::Result<(), StreamResetError>, |
| 1240 | ) -> bool { |
| 1241 | if self.active_reconfig != Some(rsn) { |
| 1242 | return false; |
| 1243 | } |
| 1244 | |
| 1245 | let fallback_ids = self |
| 1246 | .reconfigs |
| 1247 | .remove(&rsn) |
| 1248 | .map(|c| Self::reconfig_stream_ids(&c)) |
| 1249 | .unwrap_or_default(); |
| 1250 | let ids = self |
| 1251 | .reconfig_reset_streams |
| 1252 | .remove(&rsn) |
| 1253 | .unwrap_or(fallback_ids); |
| 1254 | self.active_reconfig = None; |
| 1255 | self.will_retransmit_reconfig = false; |
| 1256 | self.timers.stop(Timer::Reconfig); |
| 1257 | |
| 1258 | match outcome { |
| 1259 | Ok(()) => self.emit_reset_complete(ids), |
| 1260 | Err(reason) => self.emit_reset_failed(ids, reason), |
| 1261 | } |
| 1262 | |
| 1263 | // A buffered request or locally queued reset can now become active. |
| 1264 | self.awake_write_loop(); |
| 1265 | true |
| 1266 | } |
| 1267 | |
| 1268 | /// unregister_stream un-registers a stream from the association |
| 1269 | /// The caller should hold the association write lock. |
| 1270 | fn unregister_stream(&mut self, stream_identifier: StreamId, emit_stream_finished: bool) { |
| 1271 | if let Some(mut s) = self.streams.remove(&stream_identifier) { |
| 1272 | debug!("[{}] unregister_stream {}", self.side, stream_identifier); |
| 1273 | s.state = RecvSendState::Closed; |
| 1274 | if emit_stream_finished { |
| 1275 | self.queue_stream_finished(stream_identifier, true); |
| 1276 | } |
| 1277 | } |
| 1278 | } |
| 1279 | |
| 1280 | fn queue_stream_finished(&mut self, stream_identifier: StreamId, ready: bool) { |
| 1281 | self.events.push_back(Event::Stream(StreamEvent::Finished { |
| 1282 | id: stream_identifier, |
| 1283 | })); |
| 1284 | // Every Finished owes a terminal reset event, even if the peer |
| 1285 | // re-creates the stream id before the handshake completes. |
| 1286 | self.pending_reset_completions.insert(stream_identifier); |
| 1287 | if ready { |
| 1288 | self.ready_stream_finishes.insert(stream_identifier); |
| 1289 | } |
| 1290 | } |
| 1291 | |
| 1292 | /// Queue the end of an incoming stream generation while preserving any |
| 1293 | /// already-acknowledged DATA until the application drains it. |
| 1294 | fn retire_stream(&mut self, stream_identifier: StreamId, sender_last_tsn: u32) { |
| 1295 | if let Some(boundaries) = self.retiring_streams.get_mut(&stream_identifier) { |
| 1296 | boundaries.push_back(sender_last_tsn); |
| 1297 | if let Some(updates) = self.deferred_forward_tsns.get_mut(&stream_identifier) { |
| 1298 | for update in updates { |
| 1299 | if update.generation_boundary.is_none() { |
| 1300 | update.generation_boundary = Some(sender_last_tsn); |
| 1301 | } |
| 1302 | } |
| 1303 | } |
| 1304 | self.queue_stream_finished(stream_identifier, false); |
| 1305 | return; |
| 1306 | } |
| 1307 | |
| 1308 | let has_readable_data = self.streams.get(&stream_identifier).is_some_and(|stream| { |
| 1309 | matches!( |
| 1310 | stream.state, |
| 1311 | RecvSendState::Readable | RecvSendState::ReadWritable |
| 1312 | ) && stream.get_num_bytes_in_reassembly_queue() != 0 |
| 1313 | }); |
| 1314 | |
| 1315 | if !has_readable_data { |
| 1316 | self.deferred_forward_tsns.remove(&stream_identifier); |
| 1317 | self.unregister_stream(stream_identifier, true); |
| 1318 | return; |
| 1319 | } |
| 1320 | |
| 1321 | if let Some(updates) = self.deferred_forward_tsns.get_mut(&stream_identifier) { |
| 1322 | for update in updates { |
| 1323 | if update.generation_boundary.is_none() { |
| 1324 | update.generation_boundary = Some(sender_last_tsn); |
| 1325 | } |
| 1326 | } |
| 1327 | } |
| 1328 | let mut boundaries = VecDeque::new(); |
| 1329 | boundaries.push_back(sender_last_tsn); |
| 1330 | self.retiring_streams.insert(stream_identifier, boundaries); |
| 1331 | self.queue_stream_finished(stream_identifier, false); |
| 1332 | } |
| 1333 | |
| 1334 | /// Drop a drained old generation and release DATA held for its successor. |
| 1335 | pub(crate) fn finish_retiring_stream(&mut self, stream_identifier: StreamId) -> Result<()> { |
| 1336 | if !self.retiring_streams.contains_key(&stream_identifier) { |
| 1337 | return Ok(()); |
| 1338 | } |
| 1339 | |
| 1340 | let inherited_state = self |
| 1341 | .streams |
| 1342 | .get(&stream_identifier) |
| 1343 | .map(|stream| stream.state) |
| 1344 | .unwrap_or(RecvSendState::ReadWritable); |
| 1345 | if let Some(mut stream) = self.streams.remove(&stream_identifier) { |
| 1346 | stream.state = RecvSendState::Closed; |
| 1347 | } |
| 1348 | |
| 1349 | loop { |
| 1350 | let (completed_boundary, next_boundary) = { |
| 1351 | let boundaries = self |
| 1352 | .retiring_streams |
| 1353 | .get_mut(&stream_identifier) |
| 1354 | .expect("retiring stream must retain a boundary"); |
| 1355 | let completed = boundaries |
| 1356 | .pop_front() |
| 1357 | .expect("retiring stream must have a current generation"); |
| 1358 | (completed, boundaries.front().copied()) |
| 1359 | }; |
| 1360 | self.discard_deferred_forward_tsns_through(stream_identifier, completed_boundary); |
| 1361 | self.ready_stream_finishes.insert(stream_identifier); |
| 1362 | |
| 1363 | let Some(next_boundary) = next_boundary else { |
| 1364 | self.retiring_streams.remove(&stream_identifier); |
| 1365 | let result = self.release_deferred_reset_data(); |
| 1366 | self.inherit_stream_state(stream_identifier, inherited_state); |
| 1367 | return result; |
| 1368 | }; |
| 1369 | |
| 1370 | self.release_deferred_generation_data( |
| 1371 | stream_identifier, |
| 1372 | completed_boundary, |
| 1373 | next_boundary, |
| 1374 | )?; |
| 1375 | self.inherit_stream_state(stream_identifier, inherited_state); |
| 1376 | |
| 1377 | let has_unread_data = self |
| 1378 | .streams |
| 1379 | .get(&stream_identifier) |
| 1380 | .is_some_and(|stream| stream.get_num_bytes_in_reassembly_queue() != 0); |
| 1381 | if has_unread_data { |
| 1382 | return Ok(()); |
| 1383 | } |
| 1384 | |
| 1385 | // This generation contained no deliverable DATA. Its already-queued |
| 1386 | // Finished event can become ready immediately, and any forwarding |
| 1387 | // state scoped to it must not leak into the following generation. |
| 1388 | if let Some(mut stream) = self.streams.remove(&stream_identifier) { |
| 1389 | stream.state = RecvSendState::Closed; |
| 1390 | } |
| 1391 | self.discard_deferred_forward_tsns_through(stream_identifier, next_boundary); |
| 1392 | } |
| 1393 | } |
| 1394 | |
| 1395 | fn inherit_stream_state( |
| 1396 | &mut self, |
| 1397 | stream_identifier: StreamId, |
| 1398 | inherited_state: RecvSendState, |
| 1399 | ) { |
| 1400 | if let Some(stream) = self.streams.get_mut(&stream_identifier) { |
| 1401 | stream.state = ((stream.state as u8) & (inherited_state as u8)).into(); |
| 1402 | } |
| 1403 | } |
| 1404 | |
| 1405 | /// Close every retained incoming generation without replaying its DATA. |
| 1406 | /// The current StreamState is kept so an application-owned write half and |
| 1407 | /// its configuration survive a read-side stop. |
| 1408 | pub(crate) fn discard_retiring_streams(&mut self, stream_identifier: StreamId) { |
| 1409 | let Some(boundaries) = self.retiring_streams.remove(&stream_identifier) else { |
| 1410 | return; |
| 1411 | }; |
| 1412 | |
| 1413 | for _ in boundaries { |
| 1414 | self.ready_stream_finishes.insert(stream_identifier); |
| 1415 | } |
| 1416 | let discarded_tsns: Vec<u32> = self |
| 1417 | .deferred_reset_data |
| 1418 | .iter() |
| 1419 | .filter_map(|(tsn, chunk)| { |
| 1420 | (chunk.stream_identifier == stream_identifier).then_some(*tsn) |
| 1421 | }) |
| 1422 | .collect(); |
| 1423 | for tsn in discarded_tsns { |
| 1424 | self.payload_queue.mark_as_acked(tsn); |
| 1425 | } |
| 1426 | self.deferred_reset_data |
| 1427 | .retain(|_, chunk| chunk.stream_identifier != stream_identifier); |
| 1428 | self.deferred_forward_tsns.remove(&stream_identifier); |
| 1429 | self.events.retain(|event| { |
| 1430 | !matches!( |
| 1431 | event, |
| 1432 | Event::Stream( |
| 1433 | StreamEvent::Opened { id } | StreamEvent::Readable { id } |
| 1434 | ) if *id == stream_identifier |
| 1435 | ) |
| 1436 | }); |
| 1437 | |
| 1438 | if let Some(stream) = self.streams.get_mut(&stream_identifier) { |
| 1439 | stream.reassembly_queue = |
| 1440 | ReassemblyQueue::new(stream_identifier, self.max_receive_message_size); |
| 1441 | } |
| 1442 | if !self.stream_reset_blocked(stream_identifier) { |
| 1443 | self.unregister_stream(stream_identifier, false); |
| 1444 | } |
| 1445 | } |
| 1446 | |
| 1447 | /// set_state atomically sets the state of the Association. |
| 1448 | fn set_state(&mut self, new_state: AssociationState) { |
| 1449 | if new_state != self.state { |
| 1450 | debug!( |
| 1451 | "[{}] state change: '{}' => '{}'", |
| 1452 | self.side, self.state, new_state, |
| 1453 | ); |
| 1454 | } |
| 1455 | self.state = new_state; |
| 1456 | } |
| 1457 | |
| 1458 | /// state atomically returns the state of the Association. |
| 1459 | pub(crate) fn state(&self) -> AssociationState { |
| 1460 | self.state |
| 1461 | } |
| 1462 | |
| 1463 | /// Apply common remote-side parameters from an INIT or INIT-ACK chunk. |
| 1464 | /// |
| 1465 | /// Sets `peer_last_tsn`, `rwnd`, `ssthresh`, and `use_forward_tsn` based |
| 1466 | /// on the remote peer's initial TSN, advertised window, and supported |
| 1467 | /// extensions. Used by `handle_init`, `handle_init_ack`, and |
| 1468 | /// `new_with_out_of_band_init`. |
| 1469 | fn apply_remote_init_params( |
| 1470 | &mut self, |
| 1471 | initial_tsn: u32, |
| 1472 | advertised_receiver_window_credit: u32, |
| 1473 | params: &[Box<dyn Param + Send + Sync>], |
| 1474 | context: &str, |
| 1475 | ) { |
| 1476 | // RFC 4960 ยง13.2: peer_last_tsn is the peer's initial TSN minus one. |
| 1477 | self.peer_last_tsn = initial_tsn.wrapping_sub(1); |
| 1478 | // RFC 6525 ยง4: the peer's first request sequence number is its initial |
| 1479 | // TSN, so A4's initial response value is one less than that. |
| 1480 | self.peer_last_reconfig_rsn = initial_tsn.wrapping_sub(1); |
| 1481 | self.peer_reconfig_rsn_initialized = true; |
| 1482 | |
| 1483 | self.rwnd = advertised_receiver_window_credit; |
| 1484 | debug!("[{}] initial rwnd={}", self.side, self.rwnd); |
| 1485 | |
| 1486 | // RFC 4960 Sec 7.2.1 |
| 1487 | // o The initial value of ssthresh MAY be arbitrarily high (for |
| 1488 | // example, implementations MAY use the size of the receiver |
| 1489 | // advertised window). |
| 1490 | self.ssthresh = self.rwnd; |
| 1491 | trace!( |
| 1492 | "[{}] updated cwnd={} ssthresh={} inflight={} ({})", |
| 1493 | self.side, |
| 1494 | self.cwnd, |
| 1495 | self.ssthresh, |
| 1496 | self.inflight_queue.get_num_bytes(), |
| 1497 | context, |
| 1498 | ); |
| 1499 | |
| 1500 | for param in params { |
| 1501 | if let Some(v) = param.as_any().downcast_ref::<ParamSupportedExtensions>() { |
| 1502 | for t in &v.chunk_types { |
| 1503 | if *t == CT_FORWARD_TSN { |
| 1504 | debug!("[{}] use ForwardTSN (on {})", self.side, context); |
| 1505 | self.use_forward_tsn = true; |
| 1506 | } |
| 1507 | } |
| 1508 | } |
| 1509 | } |
| 1510 | if !self.use_forward_tsn { |
| 1511 | warn!("[{}] not using ForwardTSN (on {})", self.side, context); |
| 1512 | } |
| 1513 | } |
| 1514 | |
| 1515 | /// caller must hold self.lock |
| 1516 | fn send_init(&mut self) -> Result<()> { |
| 1517 | if let Some(stored_init) = &self.stored_init { |
| 1518 | debug!("[{}] sending INIT", self.side); |
| 1519 | |
| 1520 | let outbound = Packet { |
| 1521 | common_header: CommonHeader { |
| 1522 | source_port: self.source_port, |
| 1523 | destination_port: self.destination_port, |
| 1524 | verification_tag: self.peer_verification_tag, |
| 1525 | }, |
| 1526 | chunks: vec![Box::new(stored_init.clone())], |
| 1527 | }; |
| 1528 | |
| 1529 | self.control_queue.push_back(outbound); |
| 1530 | self.awake_write_loop(); |
| 1531 | |
| 1532 | Ok(()) |
| 1533 | } else { |
| 1534 | Err(Error::ErrInitNotStoredToSend) |
| 1535 | } |
| 1536 | } |
| 1537 | |
| 1538 | /// caller must hold self.lock |
| 1539 | fn send_cookie_echo(&mut self) -> Result<()> { |
| 1540 | if let Some(stored_cookie_echo) = &self.stored_cookie_echo { |
| 1541 | debug!("[{}] sending COOKIE-ECHO", self.side); |
| 1542 | |
| 1543 | let outbound = Packet { |
| 1544 | common_header: CommonHeader { |
| 1545 | source_port: self.source_port, |
| 1546 | destination_port: self.destination_port, |
| 1547 | verification_tag: self.peer_verification_tag, |
| 1548 | }, |
| 1549 | chunks: vec![Box::new(stored_cookie_echo.clone())], |
| 1550 | }; |
| 1551 | |
| 1552 | self.control_queue.push_back(outbound); |
| 1553 | self.awake_write_loop(); |
| 1554 | |
| 1555 | Ok(()) |
| 1556 | } else { |
| 1557 | Err(Error::ErrCookieEchoNotStoredToSend) |
| 1558 | } |
| 1559 | } |
| 1560 | |
| 1561 | /// handle_inbound parses incoming raw packets |
| 1562 | fn handle_inbound(&mut self, p: Packet, now: Instant) -> Result<()> { |
| 1563 | if let Err(err) = p.check_packet() { |
| 1564 | warn!("[{}] failed validating packet {}", self.side, err); |
| 1565 | return Ok(()); |
| 1566 | } |
| 1567 | |
| 1568 | self.handle_chunk_start(); |
| 1569 | |
| 1570 | for c in &p.chunks { |
| 1571 | self.handle_chunk(&p, c, now)?; |
| 1572 | } |
| 1573 | |
| 1574 | self.handle_chunk_end(now); |
| 1575 | |
| 1576 | Ok(()) |
| 1577 | } |
| 1578 | |
| 1579 | fn handle_chunk_start(&mut self) { |
| 1580 | self.delayed_ack_triggered = false; |
| 1581 | self.immediate_ack_triggered = false; |
| 1582 | } |
| 1583 | |
| 1584 | fn handle_chunk_end(&mut self, now: Instant) { |
| 1585 | if self.immediate_ack_triggered { |
| 1586 | self.ack_state = AckState::Immediate; |
| 1587 | self.timers.stop(Timer::Ack); |
| 1588 | self.awake_write_loop(); |
| 1589 | } else if self.delayed_ack_triggered { |
| 1590 | // Will send delayed ack in the next ack timeout |
| 1591 | self.ack_state = AckState::Delay; |
| 1592 | self.timers.start(Timer::Ack, now, ACK_INTERVAL); |
| 1593 | } |
| 1594 | } |
| 1595 | |
| 1596 | #[allow(clippy::borrowed_box)] |
| 1597 | fn handle_chunk( |
| 1598 | &mut self, |
| 1599 | p: &Packet, |
| 1600 | chunk: &Box<dyn Chunk + Send + Sync>, |
| 1601 | now: Instant, |
| 1602 | ) -> Result<()> { |
| 1603 | chunk.check()?; |
| 1604 | let chunk_any = chunk.as_any(); |
| 1605 | let packets = if let Some(c) = chunk_any.downcast_ref::<ChunkInit>() { |
| 1606 | if c.is_ack { |
| 1607 | self.handle_init_ack(p, c, now)? |
| 1608 | } else { |
| 1609 | self.handle_init(p, c)? |
| 1610 | } |
| 1611 | } else if let Some(c) = chunk_any.downcast_ref::<ChunkAbort>() { |
| 1612 | let mut err_str = String::new(); |
| 1613 | for e in &c.error_causes { |
| 1614 | if matches!(e.code, USER_INITIATED_ABORT) { |
| 1615 | debug!("User initiated abort received"); |
| 1616 | let _ = self.close(); |
| 1617 | return Ok(()); |
| 1618 | } |
| 1619 | err_str += &format!("({})", e); |
| 1620 | } |
| 1621 | return Err(Error::ErrAbortChunk(err_str)); |
| 1622 | } else if let Some(c) = chunk_any.downcast_ref::<ChunkError>() { |
| 1623 | let mut err_str = String::new(); |
| 1624 | for e in &c.error_causes { |
| 1625 | err_str += &format!("({})", e); |
| 1626 | } |
| 1627 | return Err(Error::ErrAbortChunk(err_str)); |
| 1628 | } else if let Some(c) = chunk_any.downcast_ref::<ChunkHeartbeat>() { |
| 1629 | self.handle_heartbeat(c)? |
| 1630 | } else if let Some(c) = chunk_any.downcast_ref::<ChunkCookieEcho>() { |
| 1631 | self.handle_cookie_echo(c)? |
| 1632 | } else if chunk_any.downcast_ref::<ChunkCookieAck>().is_some() { |
| 1633 | self.handle_cookie_ack()? |
| 1634 | } else if let Some(c) = chunk_any.downcast_ref::<ChunkPayloadData>() { |
| 1635 | self.handle_data(c)? |
| 1636 | } else if let Some(c) = chunk_any.downcast_ref::<ChunkSelectiveAck>() { |
| 1637 | self.handle_sack(c, now)? |
| 1638 | } else if let Some(c) = chunk_any.downcast_ref::<ChunkReconfig>() { |
| 1639 | self.handle_reconfig(c)? |
| 1640 | } else if let Some(c) = chunk_any.downcast_ref::<ChunkForwardTsn>() { |
| 1641 | self.handle_forward_tsn(c)? |
| 1642 | } else if let Some(c) = chunk_any.downcast_ref::<ChunkIForwardTsn>() { |
| 1643 | self.handle_i_forward_tsn(c)? |
| 1644 | } else if let Some(c) = chunk_any.downcast_ref::<ChunkShutdown>() { |
| 1645 | self.handle_shutdown(c)? |
| 1646 | } else if let Some(c) = chunk_any.downcast_ref::<ChunkShutdownAck>() { |
| 1647 | self.handle_shutdown_ack(c)? |
| 1648 | } else if let Some(c) = chunk_any.downcast_ref::<ChunkShutdownComplete>() { |
| 1649 | self.handle_shutdown_complete(c)? |
| 1650 | } else { |
| 1651 | return Err(Error::ErrChunkTypeUnhandled); |
| 1652 | }; |
| 1653 | |
| 1654 | if !packets.is_empty() { |
| 1655 | let mut buf: VecDeque<_> = packets.into_iter().collect(); |
| 1656 | self.control_queue.append(&mut buf); |
| 1657 | self.awake_write_loop(); |
| 1658 | } |
| 1659 | |
| 1660 | Ok(()) |
| 1661 | } |
| 1662 | |
| 1663 | fn handle_init(&mut self, p: &Packet, i: &ChunkInit) -> Result<Vec<Packet>> { |
| 1664 | let state = self.state(); |
| 1665 | debug!("[{}] chunkInit received in state '{}'", self.side, state); |
| 1666 | |
| 1667 | // https://tools.ietf.org/html/rfc4960#section-5.2.1 |
| 1668 | // Upon receipt of an INIT in the COOKIE-WAIT state, an endpoint MUST |
| 1669 | // respond with an INIT ACK using the same parameters it sent in its |
| 1670 | // original INIT chunk (including its Initiate Tag, unchanged). When |
| 1671 | // responding, the endpoint MUST send the INIT ACK back to the same |
| 1672 | // address that the original INIT (sent by this endpoint) was sent. |
| 1673 | |
| 1674 | if state != AssociationState::Closed |
| 1675 | && state != AssociationState::CookieWait |
| 1676 | && state != AssociationState::CookieEchoed |
| 1677 | { |
| 1678 | // 5.2.2. Unexpected INIT in States Other than CLOSED, COOKIE-ECHOED, |
| 1679 | // COOKIE-WAIT, and SHUTDOWN-ACK-SENT |
| 1680 | return Err(Error::ErrHandleInitState); |
| 1681 | } |
| 1682 | |
| 1683 | // Should we be setting any of these permanently until we've ACKed further? |
| 1684 | self.my_max_num_inbound_streams = |
| 1685 | core::cmp::min(i.num_inbound_streams, self.my_max_num_inbound_streams); |
| 1686 | self.my_max_num_outbound_streams = |
| 1687 | core::cmp::min(i.num_outbound_streams, self.my_max_num_outbound_streams); |
| 1688 | self.peer_verification_tag = i.initiate_tag; |
| 1689 | self.source_port = p.common_header.destination_port; |
| 1690 | self.destination_port = p.common_header.source_port; |
| 1691 | |
| 1692 | self.apply_remote_init_params( |
| 1693 | i.initial_tsn, |
| 1694 | i.advertised_receiver_window_credit, |
| 1695 | &i.params, |
| 1696 | "init", |
| 1697 | ); |
| 1698 | |
| 1699 | let mut outbound = Packet { |
| 1700 | common_header: CommonHeader { |
| 1701 | verification_tag: self.peer_verification_tag, |
| 1702 | source_port: self.source_port, |
| 1703 | destination_port: self.destination_port, |
| 1704 | }, |
| 1705 | chunks: vec![], |
| 1706 | }; |
| 1707 | |
| 1708 | let mut init_ack = ChunkInit { |
| 1709 | is_ack: true, |
| 1710 | initial_tsn: self.my_next_tsn, |
| 1711 | num_outbound_streams: self.my_max_num_outbound_streams, |
| 1712 | num_inbound_streams: self.my_max_num_inbound_streams, |
| 1713 | initiate_tag: self.my_verification_tag, |
| 1714 | advertised_receiver_window_credit: self.max_receive_buffer_size, |
| 1715 | ..Default::default() |
| 1716 | }; |
| 1717 | |
| 1718 | if self.my_cookie.is_none() { |
| 1719 | self.my_cookie = Some(ParamStateCookie::new()); |
| 1720 | } |
| 1721 | |
| 1722 | if let Some(my_cookie) = &self.my_cookie { |
| 1723 | init_ack.params = vec![Box::new(my_cookie.clone())]; |
| 1724 | } |
| 1725 | |
| 1726 | init_ack.set_supported_extensions(); |
| 1727 | |
| 1728 | outbound.chunks = vec![Box::new(init_ack)]; |
| 1729 | |
| 1730 | Ok(vec![outbound]) |
| 1731 | } |
| 1732 | |
| 1733 | fn handle_init_ack( |
| 1734 | &mut self, |
| 1735 | p: &Packet, |
| 1736 | i: &ChunkInitAck, |
| 1737 | now: Instant, |
| 1738 | ) -> Result<Vec<Packet>> { |
| 1739 | let state = self.state(); |
| 1740 | debug!("[{}] chunkInitAck received in state '{}'", self.side, state); |
| 1741 | if state != AssociationState::CookieWait { |
| 1742 | // RFC 4960 |
| 1743 | // 5.2.3. Unexpected INIT ACK |
| 1744 | // If an INIT ACK is received by an endpoint in any state other than the |
| 1745 | // COOKIE-WAIT state, the endpoint should discard the INIT ACK chunk. |
| 1746 | // An unexpected INIT ACK usually indicates the processing of an old or |
| 1747 | // duplicated INIT chunk. |
| 1748 | return Ok(vec![]); |
| 1749 | } |
| 1750 | |
| 1751 | self.my_max_num_inbound_streams = |
| 1752 | core::cmp::min(i.num_inbound_streams, self.my_max_num_inbound_streams); |
| 1753 | self.my_max_num_outbound_streams = |
| 1754 | core::cmp::min(i.num_outbound_streams, self.my_max_num_outbound_streams); |
| 1755 | self.peer_verification_tag = i.initiate_tag; |
| 1756 | if self.source_port != p.common_header.destination_port |
| 1757 | || self.destination_port != p.common_header.source_port |
| 1758 | { |
| 1759 | warn!("[{}] handle_init_ack: port mismatch", self.side); |
| 1760 | return Ok(vec![]); |
| 1761 | } |
| 1762 | |
| 1763 | self.apply_remote_init_params( |
| 1764 | i.initial_tsn, |
| 1765 | i.advertised_receiver_window_credit, |
| 1766 | &i.params, |
| 1767 | "initAck", |
| 1768 | ); |
| 1769 | |
| 1770 | self.timers.stop(Timer::T1Init); |
| 1771 | self.stored_init = None; |
| 1772 | |
| 1773 | let cookie_param = i |
| 1774 | .params |
| 1775 | .iter() |
| 1776 | .find_map(|param| param.as_any().downcast_ref::<ParamStateCookie>()); |
| 1777 | |
| 1778 | if let Some(v) = cookie_param { |
| 1779 | self.stored_cookie_echo = Some(ChunkCookieEcho { |
| 1780 | cookie: v.cookie.clone(), |
| 1781 | }); |
| 1782 | |
| 1783 | self.send_cookie_echo()?; |
| 1784 | |
| 1785 | self.timers |
| 1786 | .start(Timer::T1Cookie, now, self.rto_mgr.get_rto()); |
| 1787 | |
| 1788 | self.set_state(AssociationState::CookieEchoed); |
| 1789 | |
| 1790 | Ok(vec![]) |
| 1791 | } else { |
| 1792 | Err(Error::ErrInitAckNoCookie) |
| 1793 | } |
| 1794 | } |
| 1795 | |
| 1796 | fn handle_heartbeat(&self, c: &ChunkHeartbeat) -> Result<Vec<Packet>> { |
| 1797 | trace!("[{}] chunkHeartbeat", self.side); |
| 1798 | if let Some(p) = c.params.first() { |
| 1799 | if let Some(hbi) = p.as_any().downcast_ref::<ParamHeartbeatInfo>() { |
| 1800 | return Ok(vec![Packet { |
| 1801 | common_header: CommonHeader { |
| 1802 | verification_tag: self.peer_verification_tag, |
| 1803 | source_port: self.source_port, |
| 1804 | destination_port: self.destination_port, |
| 1805 | }, |
| 1806 | chunks: vec![Box::new(ChunkHeartbeatAck { |
| 1807 | params: vec![Box::new(ParamHeartbeatInfo { |
| 1808 | heartbeat_information: hbi.heartbeat_information.clone(), |
| 1809 | })], |
| 1810 | })], |
| 1811 | }]); |
| 1812 | } else { |
| 1813 | warn!( |
| 1814 | "[{}] failed to handle Heartbeat, no ParamHeartbeatInfo", |
| 1815 | self.side, |
| 1816 | ); |
| 1817 | } |
| 1818 | } |
| 1819 | |
| 1820 | Ok(vec![]) |
| 1821 | } |
| 1822 | |
| 1823 | fn handle_cookie_echo(&mut self, c: &ChunkCookieEcho) -> Result<Vec<Packet>> { |
| 1824 | let state = self.state(); |
| 1825 | debug!("[{}] COOKIE-ECHO received in state '{}'", self.side, state); |
| 1826 | |
| 1827 | if let Some(my_cookie) = &self.my_cookie { |
| 1828 | match state { |
| 1829 | AssociationState::Established => { |
| 1830 | if my_cookie.cookie != c.cookie { |
| 1831 | return Ok(vec![]); |
| 1832 | } |
| 1833 | } |
| 1834 | AssociationState::Closed |
| 1835 | | AssociationState::CookieWait |
| 1836 | | AssociationState::CookieEchoed => { |
| 1837 | if my_cookie.cookie != c.cookie { |
| 1838 | return Ok(vec![]); |
| 1839 | } |
| 1840 | |
| 1841 | self.timers.stop(Timer::T1Init); |
| 1842 | self.stored_init = None; |
| 1843 | |
| 1844 | self.timers.stop(Timer::T1Cookie); |
| 1845 | self.stored_cookie_echo = None; |
| 1846 | |
| 1847 | self.events.push_back(Event::Connected); |
| 1848 | self.set_state(AssociationState::Established); |
| 1849 | self.handshake_completed = true; |
| 1850 | } |
| 1851 | _ => return Ok(vec![]), |
| 1852 | }; |
| 1853 | } else { |
| 1854 | debug!("[{}] COOKIE-ECHO received before initialization", self.side); |
| 1855 | return Ok(vec![]); |
| 1856 | } |
| 1857 | |
| 1858 | Ok(vec![Packet { |
| 1859 | common_header: CommonHeader { |
| 1860 | verification_tag: self.peer_verification_tag, |
| 1861 | source_port: self.source_port, |
| 1862 | destination_port: self.destination_port, |
| 1863 | }, |
| 1864 | chunks: vec![Box::new(ChunkCookieAck {})], |
| 1865 | }]) |
| 1866 | } |
| 1867 | |
| 1868 | fn handle_cookie_ack(&mut self) -> Result<Vec<Packet>> { |
| 1869 | let state = self.state(); |
| 1870 | debug!("[{}] COOKIE-ACK received in state '{}'", self.side, state); |
| 1871 | if state != AssociationState::CookieEchoed { |
| 1872 | // RFC 4960 |
| 1873 | // 5.2.5. Handle Duplicate COOKIE-ACK. |
| 1874 | // At any state other than COOKIE-ECHOED, an endpoint should silently |
| 1875 | // discard a received COOKIE ACK chunk. |
| 1876 | return Ok(vec![]); |
| 1877 | } |
| 1878 | |
| 1879 | self.timers.stop(Timer::T1Cookie); |
| 1880 | self.stored_cookie_echo = None; |
| 1881 | |
| 1882 | self.events.push_back(Event::Connected); |
| 1883 | self.set_state(AssociationState::Established); |
| 1884 | self.handshake_completed = true; |
| 1885 | |
| 1886 | Ok(vec![]) |
| 1887 | } |
| 1888 | |
| 1889 | fn data_is_above_pending_reset(&self, d: &ChunkPayloadData) -> bool { |
| 1890 | self.retiring_streams.contains_key(&d.stream_identifier) |
| 1891 | || self.reconfig_requests.values().any(|request| { |
| 1892 | Self::reset_request_affects_stream(request, d.stream_identifier) |
| 1893 | && sna32gt(d.tsn, request.sender_last_tsn) |
| 1894 | }) |
| 1895 | } |
| 1896 | |
| 1897 | fn current_reset_boundary(&self, stream_identifier: StreamId) -> Option<u32> { |
| 1898 | self.retiring_streams |
| 1899 | .get(&stream_identifier) |
| 1900 | .and_then(|boundaries| boundaries.front().copied()) |
| 1901 | } |
| 1902 | |
| 1903 | fn pending_reset_boundary(&self, stream_identifier: StreamId) -> Option<u32> { |
| 1904 | self.reconfig_requests |
| 1905 | .values() |
| 1906 | .find(|request| Self::reset_request_affects_stream(request, stream_identifier)) |
| 1907 | .map(|request| request.sender_last_tsn) |
| 1908 | } |
| 1909 | |
| 1910 | fn forward_tsn_applies_to_current_generation(&self, stream_identifier: StreamId) -> bool { |
| 1911 | let current_boundary = self.current_reset_boundary(stream_identifier); |
| 1912 | current_boundary.is_none() |
| 1913 | || current_boundary == self.pending_reset_boundary(stream_identifier) |
| 1914 | } |
| 1915 | |
| 1916 | fn apply_or_defer_ordered_forward_tsn( |
| 1917 | &mut self, |
| 1918 | stream_identifier: StreamId, |
| 1919 | last_ssn: u16, |
| 1920 | new_cumulative_tsn: u32, |
| 1921 | ) { |
| 1922 | if !self.forward_tsn_applies_to_current_generation(stream_identifier) { |
| 1923 | let generation_boundary = self.pending_reset_boundary(stream_identifier); |
| 1924 | let kind = { |
| 1925 | let updates = self |
| 1926 | .deferred_forward_tsns |
| 1927 | .entry(stream_identifier) |
| 1928 | .or_default(); |
| 1929 | let previous = updates.iter_mut().find(|update| { |
| 1930 | if update.generation_boundary != generation_boundary { |
| 1931 | return false; |
| 1932 | } |
| 1933 | let DeferredForwardTsnKind::Ordered { |
| 1934 | last_ssn: previous_ssn, |
| 1935 | .. |
| 1936 | } = update.kind |
| 1937 | else { |
| 1938 | return false; |
| 1939 | }; |
| 1940 | previous_ssn == last_ssn |
| 1941 | || (generation_boundary.is_some() && sna16lt(previous_ssn, last_ssn)) |
| 1942 | }); |
| 1943 | |
| 1944 | if let Some(previous) = previous { |
| 1945 | let DeferredForwardTsnKind::Ordered { |
| 1946 | last_ssn: previous_ssn, |
| 1947 | new_cumulative_tsn: previous_cumulative_tsn, |
| 1948 | } = &mut previous.kind |
| 1949 | else { |
| 1950 | unreachable!(); |
| 1951 | }; |
| 1952 | if sna16lt(*previous_ssn, last_ssn) { |
| 1953 | *previous_ssn = last_ssn; |
| 1954 | } |
| 1955 | if sna32lt(*previous_cumulative_tsn, new_cumulative_tsn) { |
| 1956 | *previous_cumulative_tsn = new_cumulative_tsn; |
| 1957 | } |
| 1958 | previous.kind |
| 1959 | } else { |
| 1960 | let kind = DeferredForwardTsnKind::Ordered { |
| 1961 | last_ssn, |
| 1962 | new_cumulative_tsn, |
| 1963 | }; |
| 1964 | updates.push_back(DeferredForwardTsn { |
| 1965 | generation_boundary, |
| 1966 | kind, |
| 1967 | }); |
| 1968 | kind |
| 1969 | } |
| 1970 | }; |
| 1971 | self.prune_deferred_reset_data_for_forward_tsn( |
| 1972 | stream_identifier, |
| 1973 | generation_boundary, |
| 1974 | kind, |
| 1975 | ); |
| 1976 | return; |
| 1977 | } |
| 1978 | |
| 1979 | let became_readable = self |
| 1980 | .streams |
| 1981 | .get_mut(&stream_identifier) |
| 1982 | .is_some_and(|stream| { |
| 1983 | let was_readable = stream.reassembly_queue.is_readable(); |
| 1984 | stream.reassembly_queue.forward_tsn_for_ordered(last_ssn); |
| 1985 | !was_readable && stream.reassembly_queue.is_readable() |
| 1986 | }); |
| 1987 | if became_readable { |
| 1988 | self.events.push_back(Event::Stream(StreamEvent::Readable { |
| 1989 | id: stream_identifier, |
| 1990 | })); |
| 1991 | } |
| 1992 | } |
| 1993 | |
| 1994 | fn apply_or_defer_unordered_forward_tsn( |
| 1995 | &mut self, |
| 1996 | stream_identifier: StreamId, |
| 1997 | new_cumulative_tsn: u32, |
| 1998 | ) { |
| 1999 | if !self.forward_tsn_applies_to_current_generation(stream_identifier) { |
| 2000 | let generation_boundary = self.pending_reset_boundary(stream_identifier); |
| 2001 | let effective_cumulative_tsn = { |
| 2002 | let updates = self |
| 2003 | .deferred_forward_tsns |
| 2004 | .entry(stream_identifier) |
| 2005 | .or_default(); |
| 2006 | if let Some(previous_cumulative_tsn) = updates.iter_mut().find_map(|update| { |
| 2007 | if update.generation_boundary != generation_boundary { |
| 2008 | return None; |
| 2009 | } |
| 2010 | match &mut update.kind { |
| 2011 | DeferredForwardTsnKind::Unordered { new_cumulative_tsn } => { |
| 2012 | Some(new_cumulative_tsn) |
| 2013 | } |
| 2014 | DeferredForwardTsnKind::Ordered { .. } => None, |
| 2015 | } |
| 2016 | }) { |
| 2017 | if sna32lt(*previous_cumulative_tsn, new_cumulative_tsn) { |
| 2018 | *previous_cumulative_tsn = new_cumulative_tsn; |
| 2019 | } |
| 2020 | *previous_cumulative_tsn |
| 2021 | } else { |
| 2022 | updates.push_back(DeferredForwardTsn { |
| 2023 | generation_boundary, |
| 2024 | kind: DeferredForwardTsnKind::Unordered { new_cumulative_tsn }, |
| 2025 | }); |
| 2026 | new_cumulative_tsn |
| 2027 | } |
| 2028 | }; |
| 2029 | self.prune_deferred_reset_data_for_forward_tsn( |
| 2030 | stream_identifier, |
| 2031 | generation_boundary, |
| 2032 | DeferredForwardTsnKind::Unordered { |
| 2033 | new_cumulative_tsn: effective_cumulative_tsn, |
| 2034 | }, |
| 2035 | ); |
| 2036 | return; |
| 2037 | } |
| 2038 | |
| 2039 | if let Some(stream) = self.streams.get_mut(&stream_identifier) { |
| 2040 | stream |
| 2041 | .reassembly_queue |
| 2042 | .forward_tsn_for_unordered(new_cumulative_tsn); |
| 2043 | } |
| 2044 | } |
| 2045 | |
| 2046 | fn apply_unordered_forward_tsn_to_current_streams(&mut self, new_cumulative_tsn: u32) { |
| 2047 | let stream_ids: Vec<StreamId> = self.streams.keys().copied().collect(); |
| 2048 | for stream_identifier in stream_ids { |
| 2049 | self.apply_or_defer_unordered_forward_tsn(stream_identifier, new_cumulative_tsn); |
| 2050 | } |
| 2051 | } |
| 2052 | |
| 2053 | fn deferred_forward_tsn_application( |
| 2054 | &self, |
| 2055 | stream_identifier: StreamId, |
| 2056 | update: DeferredForwardTsn, |
| 2057 | ) -> Option<(DeferredForwardTsnKind, bool)> { |
| 2058 | match update.kind { |
| 2059 | DeferredForwardTsnKind::Ordered { |
| 2060 | last_ssn, |
| 2061 | new_cumulative_tsn, |
| 2062 | } => self.streams.get(&stream_identifier).and_then(|stream| { |
| 2063 | stream |
| 2064 | .reassembly_queue |
| 2065 | .applicable_forward_tsn_for_ordered_bounded( |
| 2066 | last_ssn, |
| 2067 | new_cumulative_tsn, |
| 2068 | self.peer_last_tsn, |
| 2069 | ) |
| 2070 | .map(|applicable_last_ssn| { |
| 2071 | ( |
| 2072 | DeferredForwardTsnKind::Ordered { |
| 2073 | last_ssn: applicable_last_ssn, |
| 2074 | new_cumulative_tsn, |
| 2075 | }, |
| 2076 | applicable_last_ssn == last_ssn, |
| 2077 | ) |
| 2078 | }) |
| 2079 | }), |
| 2080 | DeferredForwardTsnKind::Unordered { .. } => Some((update.kind, true)), |
| 2081 | } |
| 2082 | } |
| 2083 | |
| 2084 | fn apply_deferred_forward_tsns(&mut self, stream_identifier: StreamId) { |
| 2085 | if !self.streams.contains_key(&stream_identifier) { |
| 2086 | return; |
| 2087 | } |
| 2088 | |
| 2089 | let boundary = self.current_reset_boundary(stream_identifier); |
| 2090 | let Some(updates) = self.deferred_forward_tsns.get(&stream_identifier) else { |
| 2091 | return; |
| 2092 | }; |
| 2093 | let mut removals = Vec::with_capacity(updates.len()); |
| 2094 | let mut ready = vec![]; |
| 2095 | for update in updates { |
| 2096 | let application = if update.generation_boundary == boundary { |
| 2097 | self.deferred_forward_tsn_application(stream_identifier, *update) |
| 2098 | } else { |
| 2099 | None |
| 2100 | }; |
| 2101 | let remove = application.is_some_and(|(_, remove)| remove); |
| 2102 | removals.push(remove); |
| 2103 | if let Some((kind, _)) = application { |
| 2104 | ready.push(kind); |
| 2105 | } |
| 2106 | } |
| 2107 | |
| 2108 | if ready.is_empty() { |
| 2109 | return; |
| 2110 | } |
| 2111 | |
| 2112 | let mut index = 0; |
| 2113 | if let Some(updates) = self.deferred_forward_tsns.get_mut(&stream_identifier) { |
| 2114 | updates.retain(|_| { |
| 2115 | let retain = !removals[index]; |
| 2116 | index += 1; |
| 2117 | retain |
| 2118 | }); |
| 2119 | } |
| 2120 | if self |
| 2121 | .deferred_forward_tsns |
| 2122 | .get(&stream_identifier) |
| 2123 | .is_some_and(VecDeque::is_empty) |
| 2124 | { |
| 2125 | self.deferred_forward_tsns.remove(&stream_identifier); |
| 2126 | } |
| 2127 | |
| 2128 | let became_readable = if let Some(stream) = self.streams.get_mut(&stream_identifier) { |
| 2129 | let was_readable = stream.reassembly_queue.is_readable(); |
| 2130 | for kind in ready { |
| 2131 | match kind { |
| 2132 | DeferredForwardTsnKind::Ordered { |
| 2133 | last_ssn, |
| 2134 | new_cumulative_tsn, |
| 2135 | } => { |
| 2136 | stream |
| 2137 | .reassembly_queue |
| 2138 | .forward_tsn_for_ordered_bounded(last_ssn, new_cumulative_tsn); |
| 2139 | } |
| 2140 | DeferredForwardTsnKind::Unordered { new_cumulative_tsn } => { |
| 2141 | stream |
| 2142 | .reassembly_queue |
| 2143 | .forward_tsn_for_unordered(new_cumulative_tsn); |
| 2144 | } |
| 2145 | } |
| 2146 | } |
| 2147 | !was_readable && stream.reassembly_queue.is_readable() |
| 2148 | } else { |
| 2149 | false |
| 2150 | }; |
| 2151 | if became_readable { |
| 2152 | self.events.push_back(Event::Stream(StreamEvent::Readable { |
| 2153 | id: stream_identifier, |
| 2154 | })); |
| 2155 | } |
| 2156 | } |
| 2157 | |
| 2158 | fn discard_deferred_forward_tsns_through( |
| 2159 | &mut self, |
| 2160 | stream_identifier: StreamId, |
| 2161 | boundary: u32, |
| 2162 | ) { |
| 2163 | if let Some(updates) = self.deferred_forward_tsns.get_mut(&stream_identifier) { |
| 2164 | updates.retain(|update| update.generation_boundary != Some(boundary)); |
| 2165 | if updates.is_empty() { |
| 2166 | self.deferred_forward_tsns.remove(&stream_identifier); |
| 2167 | } |
| 2168 | } |
| 2169 | } |
| 2170 | |
| 2171 | /// Release receive-window credit for incomplete successor messages that a |
| 2172 | /// deferred Forward-TSN has abandoned. Complete messages remain available |
| 2173 | /// when that stream generation is eventually exposed to the application. |
| 2174 | fn prune_deferred_reset_data_for_forward_tsn( |
| 2175 | &mut self, |
| 2176 | stream_identifier: StreamId, |
| 2177 | generation_boundary: Option<u32>, |
| 2178 | kind: DeferredForwardTsnKind, |
| 2179 | ) { |
| 2180 | let mut chunks: Vec<ChunkPayloadData> = self |
| 2181 | .deferred_reset_data |
| 2182 | .values() |
| 2183 | .filter(|chunk| { |
| 2184 | chunk.stream_identifier == stream_identifier |
| 2185 | && self |
| 2186 | .deferred_generation_bounds(stream_identifier, chunk.tsn) |
| 2187 | .is_some_and(|(_, upper)| upper == generation_boundary) |
| 2188 | }) |
| 2189 | .cloned() |
| 2190 | .collect(); |
| 2191 | if chunks.is_empty() { |
| 2192 | return; |
| 2193 | } |
| 2194 | |
| 2195 | chunks.sort_unstable_by_key(|chunk| { |
| 2196 | self.deferred_generation_bounds(stream_identifier, chunk.tsn) |
| 2197 | .map(|(lower, _)| chunk.tsn.wrapping_sub(lower)) |
| 2198 | .unwrap_or_default() |
| 2199 | }); |
| 2200 | let targeted_tsns: FxHashSet<u32> = chunks.iter().map(|chunk| chunk.tsn).collect(); |
| 2201 | |
| 2202 | let mut queue = ReassemblyQueue::new(stream_identifier, self.max_receive_message_size); |
| 2203 | for chunk in chunks { |
| 2204 | // These chunks were validated before being retained. If rebuilding |
| 2205 | // their generation unexpectedly fails, keep the data conservatively. |
| 2206 | if queue.push(chunk).is_err() { |
| 2207 | return; |
| 2208 | } |
| 2209 | } |
| 2210 | match kind { |
| 2211 | DeferredForwardTsnKind::Ordered { |
| 2212 | last_ssn, |
| 2213 | new_cumulative_tsn, |
| 2214 | } => { |
| 2215 | queue.forward_tsn_for_ordered_bounded(last_ssn, new_cumulative_tsn); |
| 2216 | } |
| 2217 | DeferredForwardTsnKind::Unordered { new_cumulative_tsn } => { |
| 2218 | queue.forward_tsn_for_unordered(new_cumulative_tsn); |
| 2219 | } |
| 2220 | } |
| 2221 | |
| 2222 | let retained_tsns: FxHashSet<u32> = queue |
| 2223 | .ordered |
| 2224 | .iter() |
| 2225 | .chain(&queue.unordered) |
| 2226 | .flat_map(|chunks| chunks.chunks.iter()) |
| 2227 | .chain(queue.unordered_chunks.iter()) |
| 2228 | .map(|chunk| chunk.tsn) |
| 2229 | .collect(); |
| 2230 | self.deferred_reset_data |
| 2231 | .retain(|tsn, _| !targeted_tsns.contains(tsn) || retained_tsns.contains(tsn)); |
| 2232 | } |
| 2233 | |
| 2234 | fn deliver_deferred_reset_data(&mut self, d: &ChunkPayloadData) -> Result<()> { |
| 2235 | if self.get_or_create_stream(d.stream_identifier).is_none() { |
| 2236 | debug!("[{}] discard {}", self.side, d.stream_sequence_number); |
| 2237 | return Ok(()); |
| 2238 | } |
| 2239 | |
| 2240 | if let Some(stream) = self.streams.get_mut(&d.stream_identifier) { |
| 2241 | let queued = stream.handle_data(d)?; |
| 2242 | self.events.push_back(Event::DatagramReceived); |
| 2243 | if queued && stream.reassembly_queue.is_readable() { |
| 2244 | self.events.push_back(Event::Stream(StreamEvent::Readable { |
| 2245 | id: d.stream_identifier, |
| 2246 | })); |
| 2247 | } |
| 2248 | } |
| 2249 | |
| 2250 | Ok(()) |
| 2251 | } |
| 2252 | |
| 2253 | fn release_deferred_generation_data( |
| 2254 | &mut self, |
| 2255 | stream_identifier: StreamId, |
| 2256 | previous_boundary: u32, |
| 2257 | boundary: u32, |
| 2258 | ) -> Result<()> { |
| 2259 | let mut ready: Vec<ChunkPayloadData> = self |
| 2260 | .deferred_reset_data |
| 2261 | .values() |
| 2262 | .filter(|chunk| { |
| 2263 | chunk.stream_identifier == stream_identifier |
| 2264 | && sna32gt(chunk.tsn, previous_boundary) |
| 2265 | && sna32lte(chunk.tsn, boundary) |
| 2266 | }) |
| 2267 | .cloned() |
| 2268 | .collect(); |
| 2269 | ready.sort_unstable_by_key(|chunk| chunk.tsn.wrapping_sub(previous_boundary)); |
| 2270 | |
| 2271 | for chunk in ready { |
| 2272 | self.deliver_deferred_reset_data(&chunk)?; |
| 2273 | self.deferred_reset_data.remove(&chunk.tsn); |
| 2274 | } |
| 2275 | self.apply_deferred_forward_tsns(stream_identifier); |
| 2276 | |
| 2277 | Ok(()) |
| 2278 | } |
| 2279 | |
| 2280 | fn deferred_generation_bounds( |
| 2281 | &self, |
| 2282 | stream_identifier: StreamId, |
| 2283 | tsn: u32, |
| 2284 | ) -> Option<(u32, Option<u32>)> { |
| 2285 | let retirement_boundaries = self.retiring_streams.get(&stream_identifier); |
| 2286 | let pending_boundary = self |
| 2287 | .reconfig_requests |
| 2288 | .values() |
| 2289 | .find(|request| Self::reset_request_affects_stream(request, stream_identifier)) |
| 2290 | .map(|request| request.sender_last_tsn); |
| 2291 | |
| 2292 | let mut boundaries: Vec<u32> = retirement_boundaries |
| 2293 | .into_iter() |
| 2294 | .flat_map(|boundaries| boundaries.iter().copied()) |
| 2295 | .collect(); |
| 2296 | if let Some(pending_boundary) = pending_boundary { |
| 2297 | if boundaries.last().copied() != Some(pending_boundary) { |
| 2298 | boundaries.push(pending_boundary); |
| 2299 | } |
| 2300 | } |
| 2301 | |
| 2302 | let mut lower = *boundaries.first()?; |
| 2303 | for upper in boundaries.into_iter().skip(1) { |
| 2304 | if sna32lte(tsn, upper) { |
| 2305 | return Some((lower, Some(upper))); |
| 2306 | } |
| 2307 | lower = upper; |
| 2308 | } |
| 2309 | Some((lower, None)) |
| 2310 | } |
| 2311 | |
| 2312 | fn validate_deferred_reset_data(&self, d: &ChunkPayloadData) -> Result<()> { |
| 2313 | let Some((lower, upper)) = self.deferred_generation_bounds(d.stream_identifier, d.tsn) |
| 2314 | else { |
| 2315 | return Ok(()); |
| 2316 | }; |
| 2317 | |
| 2318 | let mut chunks: Vec<ChunkPayloadData> = self |
| 2319 | .deferred_reset_data |
| 2320 | .values() |
| 2321 | .filter(|chunk| { |
| 2322 | chunk.stream_identifier == d.stream_identifier |
| 2323 | && sna32gt(chunk.tsn, lower) |
| 2324 | && upper.is_none_or(|upper| sna32lte(chunk.tsn, upper)) |
| 2325 | }) |
| 2326 | .cloned() |
| 2327 | .collect(); |
| 2328 | chunks.sort_unstable_by_key(|chunk| chunk.tsn.wrapping_sub(lower)); |
| 2329 | |
| 2330 | let mut validation_queue = |
| 2331 | ReassemblyQueue::new(d.stream_identifier, self.max_receive_message_size); |
| 2332 | for chunk in chunks { |
| 2333 | validation_queue.push(chunk)?; |
| 2334 | } |
| 2335 | validation_queue.push(d.clone())?; |
| 2336 | Ok(()) |
| 2337 | } |
| 2338 | |
| 2339 | fn release_deferred_reset_data(&mut self) -> Result<()> { |
| 2340 | let ready: Vec<ChunkPayloadData> = self |
| 2341 | .deferred_reset_data |
| 2342 | .values() |
| 2343 | .filter(|chunk| !self.data_is_above_pending_reset(chunk)) |
| 2344 | .cloned() |
| 2345 | .collect(); |
| 2346 | let stream_ids: FxHashSet<StreamId> = |
| 2347 | ready.iter().map(|chunk| chunk.stream_identifier).collect(); |
| 2348 | |
| 2349 | for chunk in ready { |
| 2350 | self.deliver_deferred_reset_data(&chunk)?; |
| 2351 | self.deferred_reset_data.remove(&chunk.tsn); |
| 2352 | } |
| 2353 | for stream_identifier in stream_ids { |
| 2354 | self.apply_deferred_forward_tsns(stream_identifier); |
| 2355 | } |
| 2356 | |
| 2357 | Ok(()) |
| 2358 | } |
| 2359 | |
| 2360 | fn handle_data(&mut self, d: &ChunkPayloadData) -> Result<Vec<Packet>> { |
| 2361 | trace!( |
| 2362 | "[{}] DATA: tsn={} immediateSack={} len={}", |
| 2363 | self.side, |
| 2364 | d.tsn, |
| 2365 | d.immediate_sack, |
| 2366 | d.user_data.len() |
| 2367 | ); |
| 2368 | self.stats.inc_datas(); |
| 2369 | |
| 2370 | let can_push = self.payload_queue.can_push(d, self.peer_last_tsn); |
| 2371 | if can_push { |
| 2372 | self.check_receive_limits(d)?; |
| 2373 | } |
| 2374 | // A small fragment may otherwise retain a whole large packet through |
| 2375 | // Bytes::slice. Compact only when enforcing the hard resource policy; |
| 2376 | // subsequent queue clones share this one bounded payload allocation. |
| 2377 | let mut compact; |
| 2378 | let d = if can_push && self.receive_limits.is_some() { |
| 2379 | compact = d.clone(); |
| 2380 | compact.user_data = Bytes::copy_from_slice(&d.user_data); |
| 2381 | &compact |
| 2382 | } else { |
| 2383 | d |
| 2384 | }; |
| 2385 | let mut stream_handle_data = false; |
| 2386 | let mut defer_stream_data = false; |
| 2387 | if can_push && self.data_is_above_pending_reset(d) { |
| 2388 | if self.get_my_receiver_window_credit() > 0 { |
| 2389 | defer_stream_data = true; |
| 2390 | } else if let Some(last_tsn) = self.payload_queue.get_last_tsn_received() { |
| 2391 | if sna32lt(d.tsn, *last_tsn) { |
| 2392 | debug!( |
| 2393 | "[{}] receive buffer full, but accepted deferred \ |
| 2394 | reset DATA as missing chunk tsn={} ssn={}", |
| 2395 | self.side, d.tsn, d.stream_sequence_number |
| 2396 | ); |
| 2397 | defer_stream_data = true; |
| 2398 | } |
| 2399 | } |
| 2400 | } else if can_push { |
| 2401 | if self.get_or_create_stream(d.stream_identifier).is_some() { |
| 2402 | if self.get_my_receiver_window_credit() > 0 { |
| 2403 | // Pass the new chunk to stream level as soon as it arrives |
| 2404 | stream_handle_data = true; |
| 2405 | } else { |
| 2406 | // Receive buffer is full |
| 2407 | if let Some(last_tsn) = self.payload_queue.get_last_tsn_received() { |
| 2408 | if sna32lt(d.tsn, *last_tsn) { |
| 2409 | debug!( |
| 2410 | "[{}] receive buffer full, but accepted \ |
| 2411 | as missing chunk tsn={} ssn={}", |
| 2412 | self.side, d.tsn, d.stream_sequence_number |
| 2413 | ); |
| 2414 | stream_handle_data = true; |
| 2415 | } |
| 2416 | } else { |
| 2417 | debug!( |
| 2418 | "[{}] receive buffer full. dropping DATA with tsn={} ssn={}", |
| 2419 | self.side, d.tsn, d.stream_sequence_number |
| 2420 | ); |
| 2421 | } |
| 2422 | } |
| 2423 | } else { |
| 2424 | // silently discard the data. (sender will retry on T3-rtx timeout) |
| 2425 | // see pion/sctp#30 |
| 2426 | debug!("[{}] discard {}", self.side, d.stream_sequence_number); |
| 2427 | return Ok(vec![]); |
| 2428 | } |
| 2429 | } |
| 2430 | |
| 2431 | let immediate_sack = d.immediate_sack; |
| 2432 | |
| 2433 | if defer_stream_data { |
| 2434 | self.validate_deferred_reset_data(d)?; |
| 2435 | if self.payload_queue.push(d.clone(), self.peer_last_tsn) { |
| 2436 | self.deferred_reset_data.insert(d.tsn, d.clone()); |
| 2437 | } |
| 2438 | } else if stream_handle_data { |
| 2439 | if let Some(s) = self.streams.get_mut(&d.stream_identifier) { |
| 2440 | let queued = s.handle_data(d)?; |
| 2441 | // Only commit to payload_queue after reassembly accepts the chunk |
| 2442 | self.payload_queue.push(d.clone(), self.peer_last_tsn); |
| 2443 | self.events.push_back(Event::DatagramReceived); |
| 2444 | if queued && s.reassembly_queue.is_readable() { |
| 2445 | self.events.push_back(Event::Stream(StreamEvent::Readable { |
| 2446 | id: d.stream_identifier, |
| 2447 | })) |
| 2448 | } |
| 2449 | } |
| 2450 | // A deferred skip may describe this successor generation. Queue |
| 2451 | // DATA first so a complete message reusing the skipped SSN remains |
| 2452 | // readable, while a genuinely abandoned SSN can still be advanced. |
| 2453 | self.apply_deferred_forward_tsns(d.stream_identifier); |
| 2454 | } |
| 2455 | |
| 2456 | self.handle_peer_last_tsn_and_acknowledgement(immediate_sack) |
| 2457 | } |
| 2458 | |
| 2459 | fn handle_sack(&mut self, d: &ChunkSelectiveAck, now: Instant) -> Result<Vec<Packet>> { |
| 2460 | trace!( |
| 2461 | "[{}] {}, SACK: cumTSN={} a_rwnd={}", |
| 2462 | self.side, |
| 2463 | self.cumulative_tsn_ack_point, |
| 2464 | d.cumulative_tsn_ack, |
| 2465 | d.advertised_receiver_window_credit |
| 2466 | ); |
| 2467 | let state = self.state(); |
| 2468 | if state != AssociationState::Established |
| 2469 | && state != AssociationState::ShutdownPending |
| 2470 | && state != AssociationState::ShutdownReceived |
| 2471 | { |
| 2472 | return Ok(vec![]); |
| 2473 | } |
| 2474 | |
| 2475 | self.stats.inc_sacks(); |
| 2476 | |
| 2477 | if sna32gt(self.cumulative_tsn_ack_point, d.cumulative_tsn_ack) { |
| 2478 | // RFC 4960 sec 6.2.1. Processing a Received SACK |
| 2479 | // D) |
| 2480 | // i) If Cumulative TSN Ack is less than the Cumulative TSN Ack |
| 2481 | // Point, then drop the SACK. Since Cumulative TSN Ack is |
| 2482 | // monotonically increasing, a SACK whose Cumulative TSN Ack is |
| 2483 | // less than the Cumulative TSN Ack Point indicates an out-of- |
| 2484 | // order SACK. |
| 2485 | |
| 2486 | debug!( |
| 2487 | "[{}] SACK Cumulative ACK {} is older than ACK point {}", |
| 2488 | self.side, d.cumulative_tsn_ack, self.cumulative_tsn_ack_point |
| 2489 | ); |
| 2490 | |
| 2491 | return Ok(vec![]); |
| 2492 | } |
| 2493 | |
| 2494 | // Process selective ack |
| 2495 | let (bytes_acked_per_stream, htna) = self.process_selective_ack(d, now)?; |
| 2496 | |
| 2497 | let mut total_bytes_acked = 0; |
| 2498 | for n_bytes_acked in bytes_acked_per_stream.values() { |
| 2499 | total_bytes_acked += *n_bytes_acked; |
| 2500 | } |
| 2501 | |
| 2502 | let mut cum_tsn_ack_point_advanced = false; |
| 2503 | if sna32lt(self.cumulative_tsn_ack_point, d.cumulative_tsn_ack) { |
| 2504 | trace!( |
| 2505 | "[{}] SACK: cumTSN advanced: {} -> {}", |
| 2506 | self.side, self.cumulative_tsn_ack_point, d.cumulative_tsn_ack |
| 2507 | ); |
| 2508 | |
| 2509 | self.cumulative_tsn_ack_point = d.cumulative_tsn_ack; |
| 2510 | cum_tsn_ack_point_advanced = true; |
| 2511 | self.on_cumulative_tsn_ack_point_advanced(total_bytes_acked, now); |
| 2512 | } |
| 2513 | |
| 2514 | for (si, n_bytes_acked) in &bytes_acked_per_stream { |
| 2515 | if let Some(s) = self.streams.get_mut(si) { |
| 2516 | if s.on_buffer_released(*n_bytes_acked) { |
| 2517 | self.events |
| 2518 | .push_back(Event::Stream(StreamEvent::BufferedAmountLow { id: *si })) |
| 2519 | } |
| 2520 | } |
| 2521 | } |
| 2522 | |
| 2523 | // New rwnd value |
| 2524 | // RFC 4960 sec 6.2.1. Processing a Received SACK |
| 2525 | // D) |
| 2526 | // ii) Set rwnd equal to the newly received a_rwnd minus the number |
| 2527 | // of bytes still outstanding after processing the Cumulative |
| 2528 | // TSN Ack and the Gap Ack Blocks. |
| 2529 | |
| 2530 | // bytes acked were already subtracted by markAsAcked() method |
| 2531 | let bytes_outstanding = self.inflight_queue.get_num_bytes() as u32; |
| 2532 | if bytes_outstanding >= d.advertised_receiver_window_credit { |
| 2533 | self.rwnd = 0; |
| 2534 | } else { |
| 2535 | self.rwnd = d.advertised_receiver_window_credit - bytes_outstanding; |
| 2536 | } |
| 2537 | |
| 2538 | self.process_fast_retransmission(d.cumulative_tsn_ack, htna, cum_tsn_ack_point_advanced)?; |
| 2539 | |
| 2540 | if self.use_forward_tsn { |
| 2541 | // RFC 3758 Sec 3.5 C1 |
| 2542 | if sna32lt( |
| 2543 | self.advanced_peer_tsn_ack_point, |
| 2544 | self.cumulative_tsn_ack_point, |
| 2545 | ) { |
| 2546 | self.advanced_peer_tsn_ack_point = self.cumulative_tsn_ack_point |
| 2547 | } |
| 2548 | |
| 2549 | // RFC 3758 Sec 3.5 C2 |
| 2550 | let mut i = self.advanced_peer_tsn_ack_point + 1; |
| 2551 | while let Some(c) = self.inflight_queue.get(i) { |
| 2552 | if !c.abandoned() { |
| 2553 | break; |
| 2554 | } |
| 2555 | self.advanced_peer_tsn_ack_point = i; |
| 2556 | i += 1; |
| 2557 | } |
| 2558 | |
| 2559 | // RFC 3758 Sec 3.5 C3 |
| 2560 | if sna32gt( |
| 2561 | self.advanced_peer_tsn_ack_point, |
| 2562 | self.cumulative_tsn_ack_point, |
| 2563 | ) { |
| 2564 | self.will_send_forward_tsn = true; |
| 2565 | debug!( |
| 2566 | "[{}] handleSack {}: sna32GT({}, {})", |
| 2567 | self.side, |
| 2568 | self.will_send_forward_tsn, |
| 2569 | self.advanced_peer_tsn_ack_point, |
| 2570 | self.cumulative_tsn_ack_point |
| 2571 | ); |
| 2572 | } |
| 2573 | self.awake_write_loop(); |
| 2574 | } |
| 2575 | |
| 2576 | self.postprocess_sack(state, cum_tsn_ack_point_advanced, now); |
| 2577 | |
| 2578 | Ok(vec![]) |
| 2579 | } |
| 2580 | |
| 2581 | fn handle_reconfig(&mut self, c: &ChunkReconfig) -> Result<Vec<Packet>> { |
| 2582 | trace!("[{}] handle_reconfig", self.side); |
| 2583 | |
| 2584 | let mut pp = vec![]; |
| 2585 | |
| 2586 | if let Some(param_a) = &c.param_a { |
| 2587 | self.handle_reconfig_param(param_a, &mut pp)?; |
| 2588 | } |
| 2589 | |
| 2590 | if let Some(param_b) = &c.param_b { |
| 2591 | self.handle_reconfig_param(param_b, &mut pp)?; |
| 2592 | } |
| 2593 | |
| 2594 | Ok(pp) |
| 2595 | } |
| 2596 | |
| 2597 | fn handle_forward_tsn(&mut self, c: &ChunkForwardTsn) -> Result<Vec<Packet>> { |
| 2598 | trace!("[{}] FwdTSN: {}", self.side, c); |
| 2599 | |
| 2600 | if !self.use_forward_tsn { |
| 2601 | warn!("[{}] received FwdTSN but not enabled", self.side); |
| 2602 | // Return an error chunk |
| 2603 | let cerr = ChunkError { |
| 2604 | error_causes: vec![ErrorCauseUnrecognizedChunkType::default()], |
| 2605 | }; |
| 2606 | |
| 2607 | let outbound = Packet { |
| 2608 | common_header: CommonHeader { |
| 2609 | verification_tag: self.peer_verification_tag, |
| 2610 | source_port: self.source_port, |
| 2611 | destination_port: self.destination_port, |
| 2612 | }, |
| 2613 | chunks: vec![Box::new(cerr)], |
| 2614 | }; |
| 2615 | return Ok(vec![outbound]); |
| 2616 | } |
| 2617 | |
| 2618 | // From RFC 3758 Sec 3.6: |
| 2619 | // Note, if the "New Cumulative TSN" value carried in the arrived |
| 2620 | // FORWARD TSN chunk is found to be behind or at the current cumulative |
| 2621 | // TSN point, the data receiver MUST treat this FORWARD TSN as out-of- |
| 2622 | // date and MUST NOT update its Cumulative TSN. The receiver SHOULD |
| 2623 | // send a SACK to its peer (the sender of the FORWARD TSN) since such a |
| 2624 | // duplicate may indicate the previous SACK was lost in the network. |
| 2625 | |
| 2626 | trace!( |
| 2627 | "[{}] should send ack? newCumTSN={} peer_last_tsn={}", |
| 2628 | self.side, c.new_cumulative_tsn, self.peer_last_tsn |
| 2629 | ); |
| 2630 | if sna32lte(c.new_cumulative_tsn, self.peer_last_tsn) { |
| 2631 | trace!("[{}] sending ack on Forward TSN", self.side); |
| 2632 | self.ack_state = AckState::Immediate; |
| 2633 | self.timers.stop(Timer::Ack); |
| 2634 | self.awake_write_loop(); |
| 2635 | return Ok(vec![]); |
| 2636 | } |
| 2637 | |
| 2638 | // From RFC 3758 Sec 3.6: |
| 2639 | // the receiver MUST perform the same TSN handling, including duplicate |
| 2640 | // detection, gap detection, SACK generation, cumulative TSN |
| 2641 | // advancement, etc. as defined in RFC 2960 [2]---with the following |
| 2642 | // exceptions and additions. |
| 2643 | |
| 2644 | // When a FORWARD TSN chunk arrives, the data receiver MUST first update |
| 2645 | // its cumulative TSN point to the value carried in the FORWARD TSN |
| 2646 | // chunk, |
| 2647 | |
| 2648 | // Advance to the peer's new cumulative tsn point, |
| 2649 | // dropping the chunks it abandoned along the way. |
| 2650 | self.payload_queue.pop_up_to(c.new_cumulative_tsn); |
| 2651 | self.peer_last_tsn = c.new_cumulative_tsn; |
| 2652 | |
| 2653 | // Report new peer_last_tsn value and abandoned largest SSN value to |
| 2654 | // corresponding streams so that the abandoned chunks can be removed |
| 2655 | // from the reassemblyQueue. |
| 2656 | for forwarded in &c.streams { |
| 2657 | self.apply_or_defer_ordered_forward_tsn( |
| 2658 | forwarded.identifier, |
| 2659 | forwarded.sequence, |
| 2660 | c.new_cumulative_tsn, |
| 2661 | ); |
| 2662 | } |
| 2663 | |
| 2664 | // TSN may be forwarded for unordered chunks. ForwardTSN chunk does not |
| 2665 | // report which stream identifier it skipped for unordered chunks. |
| 2666 | // Therefore, we need to broadcast this event to all existing streams for |
| 2667 | // unordered chunks. |
| 2668 | // See https://github.com/pion/sctp/issues/106 |
| 2669 | self.apply_unordered_forward_tsn_to_current_streams(c.new_cumulative_tsn); |
| 2670 | |
| 2671 | let mut reply = vec![]; |
| 2672 | self.reevaluate_pending_reset_requests(&mut reply)?; |
| 2673 | reply.extend(self.handle_peer_last_tsn_and_acknowledgement(false)?); |
| 2674 | Ok(reply) |
| 2675 | } |
| 2676 | |
| 2677 | /// Handle I-FORWARD-TSN (RFC 8260) โ identical to FORWARD-TSN but with |
| 2678 | /// 32-bit MID and explicit per-entry unordered flag. |
| 2679 | fn handle_i_forward_tsn(&mut self, c: &ChunkIForwardTsn) -> Result<Vec<Packet>> { |
| 2680 | trace!("[{}] I-FwdTSN: {}", self.side, c); |
| 2681 | |
| 2682 | if !self.use_forward_tsn { |
| 2683 | warn!("[{}] received I-FwdTSN but not enabled", self.side); |
| 2684 | let cerr = ChunkError { |
| 2685 | error_causes: vec![ErrorCauseUnrecognizedChunkType::default()], |
| 2686 | }; |
| 2687 | |
| 2688 | let outbound = Packet { |
| 2689 | common_header: CommonHeader { |
| 2690 | verification_tag: self.peer_verification_tag, |
| 2691 | source_port: self.source_port, |
| 2692 | destination_port: self.destination_port, |
| 2693 | }, |
| 2694 | chunks: vec![Box::new(cerr)], |
| 2695 | }; |
| 2696 | return Ok(vec![outbound]); |
| 2697 | } |
| 2698 | |
| 2699 | if sna32lte(c.new_cumulative_tsn, self.peer_last_tsn) { |
| 2700 | trace!("[{}] sending ack on I-Forward TSN", self.side); |
| 2701 | self.ack_state = AckState::Immediate; |
| 2702 | self.timers.stop(Timer::Ack); |
| 2703 | self.awake_write_loop(); |
| 2704 | return Ok(vec![]); |
| 2705 | } |
| 2706 | |
| 2707 | // Advance to the peer's new cumulative tsn point, |
| 2708 | // dropping the chunks it abandoned along the way. |
| 2709 | self.payload_queue.pop_up_to(c.new_cumulative_tsn); |
| 2710 | self.peer_last_tsn = c.new_cumulative_tsn; |
| 2711 | |
| 2712 | // Handle per-stream entries using the explicit unordered flag |
| 2713 | for forwarded in &c.streams { |
| 2714 | if forwarded.unordered { |
| 2715 | self.apply_or_defer_unordered_forward_tsn( |
| 2716 | forwarded.identifier, |
| 2717 | c.new_cumulative_tsn, |
| 2718 | ); |
| 2719 | } else { |
| 2720 | // MID maps to SSN for ordered streams; truncate to u16 |
| 2721 | self.apply_or_defer_ordered_forward_tsn( |
| 2722 | forwarded.identifier, |
| 2723 | forwarded.mid as u16, |
| 2724 | c.new_cumulative_tsn, |
| 2725 | ); |
| 2726 | } |
| 2727 | } |
| 2728 | |
| 2729 | // Broadcast to all unordered streams |
| 2730 | self.apply_unordered_forward_tsn_to_current_streams(c.new_cumulative_tsn); |
| 2731 | |
| 2732 | let mut reply = vec![]; |
| 2733 | self.reevaluate_pending_reset_requests(&mut reply)?; |
| 2734 | reply.extend(self.handle_peer_last_tsn_and_acknowledgement(false)?); |
| 2735 | Ok(reply) |
| 2736 | } |
| 2737 | |
| 2738 | fn handle_shutdown(&mut self, _: &ChunkShutdown) -> Result<Vec<Packet>> { |
| 2739 | let state = self.state(); |
| 2740 | |
| 2741 | if state == AssociationState::Established { |
| 2742 | if !self.inflight_queue.is_empty() { |
| 2743 | self.set_state(AssociationState::ShutdownReceived); |
| 2744 | } else { |
| 2745 | // No more outstanding, send shutdown ack. |
| 2746 | self.will_send_shutdown_ack = true; |
| 2747 | self.set_state(AssociationState::ShutdownAckSent); |
| 2748 | |
| 2749 | self.awake_write_loop(); |
| 2750 | } |
| 2751 | } else if state == AssociationState::ShutdownSent { |
| 2752 | // self.cumulative_tsn_ack_point = c.cumulative_tsn_ack |
| 2753 | |
| 2754 | self.will_send_shutdown_ack = true; |
| 2755 | self.set_state(AssociationState::ShutdownAckSent); |
| 2756 | |
| 2757 | self.awake_write_loop(); |
| 2758 | } |
| 2759 | |
| 2760 | Ok(vec![]) |
| 2761 | } |
| 2762 | |
| 2763 | fn handle_shutdown_ack(&mut self, _: &ChunkShutdownAck) -> Result<Vec<Packet>> { |
| 2764 | let state = self.state(); |
| 2765 | if state == AssociationState::ShutdownSent || state == AssociationState::ShutdownAckSent { |
| 2766 | self.timers.stop(Timer::T2Shutdown); |
| 2767 | self.will_send_shutdown_complete = true; |
| 2768 | |
| 2769 | self.awake_write_loop(); |
| 2770 | } |
| 2771 | |
| 2772 | Ok(vec![]) |
| 2773 | } |
| 2774 | |
| 2775 | fn handle_shutdown_complete(&mut self, _: &ChunkShutdownComplete) -> Result<Vec<Packet>> { |
| 2776 | let state = self.state(); |
| 2777 | if state == AssociationState::ShutdownAckSent { |
| 2778 | self.timers.stop(Timer::T2Shutdown); |
| 2779 | self.close()?; |
| 2780 | } |
| 2781 | |
| 2782 | Ok(vec![]) |
| 2783 | } |
| 2784 | |
| 2785 | fn reevaluate_pending_reset_requests(&mut self, reply: &mut Vec<Packet>) -> Result<()> { |
| 2786 | let requests: Vec<ParamOutgoingResetRequest> = |
| 2787 | self.reconfig_requests.values().cloned().collect(); |
| 2788 | for request in requests { |
| 2789 | let seq = request.reconfig_request_sequence_number; |
| 2790 | self.reset_streams_if_any(&request, false, reply)?; |
| 2791 | if !self.reconfig_requests.contains_key(&seq) { |
| 2792 | self.max_completed_reconfig_rsn = Some(seq); |
| 2793 | } |
| 2794 | } |
| 2795 | Ok(()) |
| 2796 | } |
| 2797 | |
| 2798 | /// A common routine for handle_data and handle_forward_tsn routines |
| 2799 | fn handle_peer_last_tsn_and_acknowledgement( |
| 2800 | &mut self, |
| 2801 | sack_immediately: bool, |
| 2802 | ) -> Result<Vec<Packet>> { |
| 2803 | let mut reply = vec![]; |
| 2804 | |
| 2805 | // Try to advance peer_last_tsn |
| 2806 | |
| 2807 | // From RFC 3758 Sec 3.6: |
| 2808 | // .. and then MUST further advance its cumulative TSN point locally |
| 2809 | // if possible |
| 2810 | // Meaning, if peer_last_tsn+1 points to a chunk that is received, |
| 2811 | // advance peer_last_tsn until peer_last_tsn+1 points to unreceived chunk. |
| 2812 | //debug!("[{}] peer_last_tsn = {}", self.side, self.peer_last_tsn); |
| 2813 | while let Some(chunk) = self.payload_queue.pop(self.peer_last_tsn.wrapping_add(1)) { |
| 2814 | self.peer_last_tsn = self.peer_last_tsn.wrapping_add(1); |
| 2815 | //debug!("[{}] peer_last_tsn = {}", self.side, self.peer_last_tsn); |
| 2816 | |
| 2817 | self.reevaluate_pending_reset_requests(&mut reply)?; |
| 2818 | |
| 2819 | if let Some(deferred) = self.deferred_reset_data.get(&chunk.tsn).cloned() { |
| 2820 | if !self.data_is_above_pending_reset(&deferred) { |
| 2821 | self.deliver_deferred_reset_data(&deferred)?; |
| 2822 | self.deferred_reset_data.remove(&chunk.tsn); |
| 2823 | } |
| 2824 | } |
| 2825 | } |
| 2826 | |
| 2827 | // A resolved TSN gap can disambiguate a deferred ordered skip whose |
| 2828 | // later SSN arrived first on a reset successor. |
| 2829 | let stream_ids: Vec<StreamId> = self.deferred_forward_tsns.keys().copied().collect(); |
| 2830 | for stream_identifier in stream_ids { |
| 2831 | self.apply_deferred_forward_tsns(stream_identifier); |
| 2832 | } |
| 2833 | |
| 2834 | let has_packet_loss = !self.payload_queue.is_empty(); |
| 2835 | if has_packet_loss { |
| 2836 | trace!( |
| 2837 | "[{}] packetloss: {}", |
| 2838 | self.side, |
| 2839 | self.payload_queue |
| 2840 | .get_gap_ack_blocks_string(self.peer_last_tsn) |
| 2841 | ); |
| 2842 | } |
| 2843 | |
| 2844 | if (self.ack_state != AckState::Immediate |
| 2845 | && !sack_immediately |
| 2846 | && !has_packet_loss |
| 2847 | && self.ack_mode == AckMode::Normal) |
| 2848 | || self.ack_mode == AckMode::AlwaysDelay |
| 2849 | { |
| 2850 | if self.ack_state == AckState::Idle { |
| 2851 | self.delayed_ack_triggered = true; |
| 2852 | } else { |
| 2853 | self.immediate_ack_triggered = true; |
| 2854 | } |
| 2855 | } else { |
| 2856 | self.immediate_ack_triggered = true; |
| 2857 | } |
| 2858 | |
| 2859 | Ok(reply) |
| 2860 | } |
| 2861 | |
| 2862 | #[allow(clippy::borrowed_box)] |
| 2863 | fn handle_reconfig_param( |
| 2864 | &mut self, |
| 2865 | raw: &Box<dyn Param + Send + Sync>, |
| 2866 | reply: &mut Vec<Packet>, |
| 2867 | ) -> Result<()> { |
| 2868 | if let Some(p) = raw.as_any().downcast_ref::<ParamOutgoingResetRequest>() { |
| 2869 | // RFC 6525 section 5.2.2 E1: the response sequence number in an |
| 2870 | // Outgoing Reset Request implicitly acknowledges our request. |
| 2871 | self.finish_reconfig(p.reconfig_response_sequence_number, Ok(())); |
| 2872 | |
| 2873 | let seq = p.reconfig_request_sequence_number; |
| 2874 | // Detect retransmission of a completed request. An InProgress request |
| 2875 | // is still in reconfig_requests, so we must let those through for |
| 2876 | // re-evaluation (the TSN may have advanced). |
| 2877 | if !self.reconfig_requests.contains_key(&seq) |
| 2878 | && self |
| 2879 | .max_completed_reconfig_rsn |
| 2880 | .is_some_and(|w| sna32lte(seq, w)) |
| 2881 | { |
| 2882 | // Retransmission of an already-completed request. Resend the response |
| 2883 | // but do NOT reprocess stream resets (stream IDs may have been reused). |
| 2884 | self.push_reconfig_response(reply, seq, ReconfigResult::SuccessPerformed); |
| 2885 | return Ok(()); |
| 2886 | } |
| 2887 | |
| 2888 | if !self.reconfig_requests.contains_key(&seq) { |
| 2889 | // RFC 6525 section 5.2.2 E4: only the next expected peer |
| 2890 | // request sequence number may start a new operation. A request |
| 2891 | // already deferred in reconfig_requests is a retransmission and |
| 2892 | // remains eligible for re-evaluation. |
| 2893 | if !self.reconfig_requests.is_empty() { |
| 2894 | self.push_reconfig_response( |
| 2895 | reply, |
| 2896 | seq, |
| 2897 | ReconfigResult::ErrorRequestAlreadyInProgress, |
| 2898 | ); |
| 2899 | return Ok(()); |
| 2900 | } |
| 2901 | |
| 2902 | if self.peer_reconfig_rsn_initialized |
| 2903 | && seq != self.peer_last_reconfig_rsn.wrapping_add(1) |
| 2904 | { |
| 2905 | self.push_reconfig_response(reply, seq, ReconfigResult::ErrorBadSequenceNumber); |
| 2906 | return Ok(()); |
| 2907 | } |
| 2908 | |
| 2909 | self.peer_last_reconfig_rsn = seq; |
| 2910 | self.peer_reconfig_rsn_initialized = true; |
| 2911 | self.reconfig_requests.insert(seq, p.clone()); |
| 2912 | } |
| 2913 | |
| 2914 | // Re-evaluate the original request for this RSN. This prevents a |
| 2915 | // retransmission with altered stream identifiers or TSN boundary |
| 2916 | // from changing an operation that is already in progress. |
| 2917 | let request = self.reconfig_requests.get(&seq).cloned().unwrap(); |
| 2918 | self.reset_streams_if_any(&request, true, reply)?; |
| 2919 | // Update watermark only after successful completion (request |
| 2920 | // removed from reconfig_requests by reset_streams_if_any). |
| 2921 | if !self.reconfig_requests.contains_key(&seq) { |
| 2922 | self.max_completed_reconfig_rsn = Some(seq); |
| 2923 | } |
| 2924 | Ok(()) |
| 2925 | } else if let Some(p) = raw.as_any().downcast_ref::<ParamReconfigResponse>() { |
| 2926 | let rsn = p.reconfig_response_sequence_number; |
| 2927 | // RFC 6525 section 5.2.7 H1: ignore responses unless this RSN owns |
| 2928 | // the running Re-configuration Timer. |
| 2929 | if self.active_reconfig != Some(rsn) { |
| 2930 | return Ok(()); |
| 2931 | } |
| 2932 | |
| 2933 | // In progress result means the peer has deferred the request, |
| 2934 | // not answered it. The request stays outstanding and its timer restarts |
| 2935 | // without counting toward the retransmission limit. |
| 2936 | if p.result == ReconfigResult::InProgress { |
| 2937 | if self.reconfigs.contains_key(&rsn) { |
| 2938 | self.will_retransmit_reconfig = false; |
| 2939 | // Pause without resetting the existing retry count. The outbound |
| 2940 | // poll starts it again and the next expiry is exempted by H2. |
| 2941 | self.timers.set(Timer::Reconfig, None); |
| 2942 | self.timers.suppress_error_count(Timer::Reconfig); |
| 2943 | self.awake_write_loop(); |
| 2944 | } |
| 2945 | return Ok(()); |
| 2946 | } |
| 2947 | |
| 2948 | let outcome = match p.result { |
| 2949 | ReconfigResult::SuccessNop | ReconfigResult::SuccessPerformed => Ok(()), |
| 2950 | ReconfigResult::Denied => Err(StreamResetError::Denied), |
| 2951 | _ => Err(StreamResetError::Failed), |
| 2952 | }; |
| 2953 | self.finish_reconfig(rsn, outcome); |
| 2954 | Ok(()) |
| 2955 | } else { |
| 2956 | Err(Error::ErrParameterType) |
| 2957 | } |
| 2958 | } |
| 2959 | |
| 2960 | fn push_reconfig_response( |
| 2961 | &self, |
| 2962 | reply: &mut Vec<Packet>, |
| 2963 | sequence_number: u32, |
| 2964 | result: ReconfigResult, |
| 2965 | ) { |
| 2966 | reply.push(self.create_packet(vec![Box::new(ChunkReconfig { |
| 2967 | param_a: Some(Box::new(ParamReconfigResponse { |
| 2968 | reconfig_response_sequence_number: sequence_number, |
| 2969 | result, |
| 2970 | })), |
| 2971 | param_b: None, |
| 2972 | })])); |
| 2973 | } |
| 2974 | |
| 2975 | fn process_selective_ack( |
| 2976 | &mut self, |
| 2977 | d: &ChunkSelectiveAck, |
| 2978 | now: Instant, |
| 2979 | ) -> Result<(HashMap<u16, i64>, u32)> { |
| 2980 | let mut bytes_acked_per_stream = HashMap::new(); |
| 2981 | |
| 2982 | // New ack point, so pop all ACKed packets from inflight_queue |
| 2983 | // We add 1 because the "currentAckPoint" has already been popped from the inflight queue |
| 2984 | // For the first SACK we take care of this by setting the ackpoint to cumAck - 1 |
| 2985 | let mut i = self.cumulative_tsn_ack_point + 1; |
| 2986 | //log::debug!("[{}] i={} d={}", self.name, i, d.cumulative_tsn_ack); |
| 2987 | while sna32lte(i, d.cumulative_tsn_ack) { |
| 2988 | if let Some(c) = self.inflight_queue.pop(i) { |
| 2989 | if !c.acked { |
| 2990 | // RFC 4096 sec 6.3.2. Retransmission Timer Rules |
| 2991 | // R3) Whenever a SACK is received that acknowledges the DATA chunk |
| 2992 | // with the earliest outstanding TSN for that address, restart the |
| 2993 | // T3-rtx timer for that address with its current RTO (if there is |
| 2994 | // still outstanding data on that address). |
| 2995 | if i == self.cumulative_tsn_ack_point + 1 { |
| 2996 | // T3 timer needs to be reset. Stop it for now. |
| 2997 | self.timers.stop(Timer::T3RTX); |
| 2998 | } |
| 2999 | |
| 3000 | let n_bytes_acked = c.user_data.len() as i64; |
| 3001 | |
| 3002 | // Sum the number of bytes acknowledged per stream |
| 3003 | if let Some(amount) = bytes_acked_per_stream.get_mut(&c.stream_identifier) { |
| 3004 | *amount += n_bytes_acked; |
| 3005 | } else { |
| 3006 | bytes_acked_per_stream.insert(c.stream_identifier, n_bytes_acked); |
| 3007 | } |
| 3008 | |
| 3009 | // RFC 4960 sec 6.3.1. RTO Calculation |
| 3010 | // C4) When data is in flight and when allowed by rule C5 below, a new |
| 3011 | // RTT measurement MUST be made each round trip. Furthermore, new |
| 3012 | // RTT measurements SHOULD be made no more than once per round trip |
| 3013 | // for a given destination transport address. |
| 3014 | // C5) Karn's algorithm: RTT measurements MUST NOT be made using |
| 3015 | // packets that were retransmitted (and thus for which it is |
| 3016 | // ambiguous whether the reply was for the first instance of the |
| 3017 | // chunk or for a later instance) |
| 3018 | if c.nsent == 1 && sna32gte(c.tsn, self.min_tsn2measure_rtt) { |
| 3019 | self.min_tsn2measure_rtt = self.my_next_tsn; |
| 3020 | if let Some(since) = &c.since { |
| 3021 | let rtt = now.duration_since(*since); |
| 3022 | let srtt = self.rto_mgr.set_new_rtt(rtt.as_millis() as u64); |
| 3023 | trace!( |
| 3024 | "[{}] SACK: measured-rtt={} srtt={} new-rto={}", |
| 3025 | self.side, |
| 3026 | rtt.as_millis(), |
| 3027 | srtt, |
| 3028 | self.rto_mgr.get_rto() |
| 3029 | ); |
| 3030 | } else { |
| 3031 | error!("[{}] invalid c.since", self.side); |
| 3032 | } |
| 3033 | } |
| 3034 | } |
| 3035 | |
| 3036 | if self.in_fast_recovery && c.tsn == self.fast_recover_exit_point { |
| 3037 | debug!("[{}] exit fast-recovery", self.side); |
| 3038 | self.in_fast_recovery = false; |
| 3039 | } |
| 3040 | } else { |
| 3041 | return Err(Error::ErrInflightQueueTsnPop); |
| 3042 | } |
| 3043 | |
| 3044 | i += 1; |
| 3045 | } |
| 3046 | |
| 3047 | let mut htna = d.cumulative_tsn_ack; |
| 3048 | |
| 3049 | // Mark selectively acknowledged chunks as "acked" |
| 3050 | for g in &d.gap_ack_blocks { |
| 3051 | for i in g.start..=g.end { |
| 3052 | let tsn = d.cumulative_tsn_ack + i as u32; |
| 3053 | |
| 3054 | let (is_existed, is_acked) = if let Some(c) = self.inflight_queue.get(tsn) { |
| 3055 | (true, c.acked) |
| 3056 | } else { |
| 3057 | (false, false) |
| 3058 | }; |
| 3059 | let n_bytes_acked = if is_existed && !is_acked { |
| 3060 | self.inflight_queue.mark_as_acked(tsn) as i64 |
| 3061 | } else { |
| 3062 | 0 |
| 3063 | }; |
| 3064 | |
| 3065 | if let Some(c) = self.inflight_queue.get(tsn) { |
| 3066 | if !is_acked { |
| 3067 | // Sum the number of bytes acknowledged per stream |
| 3068 | if let Some(amount) = bytes_acked_per_stream.get_mut(&c.stream_identifier) { |
| 3069 | *amount += n_bytes_acked; |
| 3070 | } else { |
| 3071 | bytes_acked_per_stream.insert(c.stream_identifier, n_bytes_acked); |
| 3072 | } |
| 3073 | |
| 3074 | trace!("[{}] tsn={} has been sacked", self.side, c.tsn); |
| 3075 | |
| 3076 | if c.nsent == 1 { |
| 3077 | self.min_tsn2measure_rtt = self.my_next_tsn; |
| 3078 | if let Some(since) = &c.since { |
| 3079 | let rtt = now.duration_since(*since); |
| 3080 | let srtt = self.rto_mgr.set_new_rtt(rtt.as_millis() as u64); |
| 3081 | trace!( |
| 3082 | "[{}] SACK: measured-rtt={} srtt={} new-rto={}", |
| 3083 | self.side, |
| 3084 | rtt.as_millis(), |
| 3085 | srtt, |
| 3086 | self.rto_mgr.get_rto() |
| 3087 | ); |
| 3088 | } else { |
| 3089 | error!("[{}] invalid c.since", self.side); |
| 3090 | } |
| 3091 | } |
| 3092 | |
| 3093 | if sna32lt(htna, tsn) { |
| 3094 | htna = tsn; |
| 3095 | } |
| 3096 | } |
| 3097 | } else { |
| 3098 | return Err(Error::ErrTsnRequestNotExist); |
| 3099 | } |
| 3100 | } |
| 3101 | } |
| 3102 | |
| 3103 | Ok((bytes_acked_per_stream, htna)) |
| 3104 | } |
| 3105 | |
| 3106 | fn on_cumulative_tsn_ack_point_advanced(&mut self, total_bytes_acked: i64, now: Instant) { |
| 3107 | // RFC 4096, sec 6.3.2. Retransmission Timer Rules |
| 3108 | // R2) Whenever all outstanding data sent to an address have been |
| 3109 | // acknowledged, turn off the T3-rtx timer of that address. |
| 3110 | if self.inflight_queue.is_empty() { |
| 3111 | trace!( |
| 3112 | "[{}] SACK: no more packet in-flight (pending={})", |
| 3113 | self.side, |
| 3114 | self.pending_queue.len() |
| 3115 | ); |
| 3116 | self.timers.stop(Timer::T3RTX); |
| 3117 | } else { |
| 3118 | trace!("[{}] T3-rtx timer start (pt2)", self.side); |
| 3119 | self.timers |
| 3120 | .restart_if_stale(Timer::T3RTX, now, self.rto_mgr.get_rto()); |
| 3121 | } |
| 3122 | |
| 3123 | // Update congestion control parameters |
| 3124 | if self.cwnd <= self.ssthresh { |
| 3125 | // RFC 4096, sec 7.2.1. Slow-Start |
| 3126 | // o When cwnd is less than or equal to ssthresh, an SCTP endpoint MUST |
| 3127 | // use the slow-start algorithm to increase cwnd only if the current |
| 3128 | // congestion window is being fully utilized, an incoming SACK |
| 3129 | // advances the Cumulative TSN Ack Point, and the data sender is not |
| 3130 | // in Fast Recovery. Only when these three conditions are met can |
| 3131 | // the cwnd be increased; otherwise, the cwnd MUST not be increased. |
| 3132 | // If these conditions are met, then cwnd MUST be increased by, at |
| 3133 | // most, the lesser of 1) the total size of the previously |
| 3134 | // outstanding DATA chunk(s) acknowledged, and 2) the destination's |
| 3135 | // path MTU. |
| 3136 | if !self.in_fast_recovery && !self.pending_queue.is_empty() { |
| 3137 | self.cwnd += core::cmp::min(total_bytes_acked as u32, self.cwnd); // TCP way |
| 3138 | // self.cwnd += min32(uint32(total_bytes_acked), self.mtu) // SCTP way (slow) |
| 3139 | trace!( |
| 3140 | "[{}] updated cwnd={} ssthresh={} acked={} (SS)", |
| 3141 | self.side, self.cwnd, self.ssthresh, total_bytes_acked |
| 3142 | ); |
| 3143 | } else { |
| 3144 | trace!( |
| 3145 | "[{}] cwnd did not grow: cwnd={} ssthresh={} acked={} FR={} pending={}", |
| 3146 | self.side, |
| 3147 | self.cwnd, |
| 3148 | self.ssthresh, |
| 3149 | total_bytes_acked, |
| 3150 | self.in_fast_recovery, |
| 3151 | self.pending_queue.len() |
| 3152 | ); |
| 3153 | } |
| 3154 | } else { |
| 3155 | // RFC 4096, sec 7.2.2. Congestion Avoidance |
| 3156 | // o Whenever cwnd is greater than ssthresh, upon each SACK arrival |
| 3157 | // that advances the Cumulative TSN Ack Point, increase |
| 3158 | // partial_bytes_acked by the total number of bytes of all new chunks |
| 3159 | // acknowledged in that SACK including chunks acknowledged by the new |
| 3160 | // Cumulative TSN Ack and by Gap Ack Blocks. |
| 3161 | self.partial_bytes_acked += total_bytes_acked as u32; |
| 3162 | |
| 3163 | // o When partial_bytes_acked is equal to or greater than cwnd and |
| 3164 | // before the arrival of the SACK the sender had cwnd or more bytes |
| 3165 | // of data outstanding (i.e., before arrival of the SACK, flight size |
| 3166 | // was greater than or equal to cwnd), increase cwnd by MTU, and |
| 3167 | // reset partial_bytes_acked to (partial_bytes_acked - cwnd). |
| 3168 | if self.partial_bytes_acked >= self.cwnd && !self.pending_queue.is_empty() { |
| 3169 | self.partial_bytes_acked -= self.cwnd; |
| 3170 | self.cwnd += self.mtu; |
| 3171 | trace!( |
| 3172 | "[{}] updated cwnd={} ssthresh={} acked={} (CA)", |
| 3173 | self.side, self.cwnd, self.ssthresh, total_bytes_acked |
| 3174 | ); |
| 3175 | } |
| 3176 | } |
| 3177 | } |
| 3178 | |
| 3179 | fn process_fast_retransmission( |
| 3180 | &mut self, |
| 3181 | cum_tsn_ack_point: u32, |
| 3182 | htna: u32, |
| 3183 | cum_tsn_ack_point_advanced: bool, |
| 3184 | ) -> Result<()> { |
| 3185 | // HTNA algorithm - RFC 4960 Sec 7.2.4 |
| 3186 | // Increment missIndicator of each chunks that the SACK reported missing |
| 3187 | // when either of the following is met: |
| 3188 | // a) Not in fast-recovery |
| 3189 | // miss indications are incremented only for missing TSNs prior to the |
| 3190 | // highest TSN newly acknowledged in the SACK. |
| 3191 | // b) In fast-recovery AND the Cumulative TSN Ack Point advanced |
| 3192 | // the miss indications are incremented for all TSNs reported missing |
| 3193 | // in the SACK. |
| 3194 | if !self.in_fast_recovery || cum_tsn_ack_point_advanced { |
| 3195 | let max_tsn = if !self.in_fast_recovery { |
| 3196 | // a) increment only for missing TSNs prior to the HTNA |
| 3197 | htna |
| 3198 | } else { |
| 3199 | // b) increment for all TSNs reported missing |
| 3200 | cum_tsn_ack_point + (self.inflight_queue.len() as u32) + 1 |
| 3201 | }; |
| 3202 | |
| 3203 | let mut tsn = cum_tsn_ack_point + 1; |
| 3204 | while sna32lt(tsn, max_tsn) { |
| 3205 | if let Some(c) = self.inflight_queue.get_mut(tsn) { |
| 3206 | if !c.acked && !c.abandoned() && c.miss_indicator < 3 { |
| 3207 | c.miss_indicator += 1; |
| 3208 | if c.miss_indicator == 3 && !self.in_fast_recovery { |
| 3209 | // 2) If not in Fast Recovery, adjust the ssthresh and cwnd of the |
| 3210 | // destination address(es) to which the missing DATA chunks were |
| 3211 | // last sent, according to the formula described in Section 7.2.3. |
| 3212 | self.in_fast_recovery = true; |
| 3213 | self.fast_recover_exit_point = htna; |
| 3214 | self.ssthresh = core::cmp::max(self.cwnd / 2, 4 * self.mtu); |
| 3215 | self.cwnd = self.ssthresh; |
| 3216 | self.partial_bytes_acked = 0; |
| 3217 | self.will_retransmit_fast = true; |
| 3218 | |
| 3219 | trace!( |
| 3220 | "[{}] updated cwnd={} ssthresh={} inflight={} (FR)", |
| 3221 | self.side, |
| 3222 | self.cwnd, |
| 3223 | self.ssthresh, |
| 3224 | self.inflight_queue.get_num_bytes() |
| 3225 | ); |
| 3226 | } |
| 3227 | } |
| 3228 | } else { |
| 3229 | return Err(Error::ErrTsnRequestNotExist); |
| 3230 | } |
| 3231 | |
| 3232 | tsn += 1; |
| 3233 | } |
| 3234 | } |
| 3235 | |
| 3236 | if self.in_fast_recovery && cum_tsn_ack_point_advanced { |
| 3237 | self.will_retransmit_fast = true; |
| 3238 | } |
| 3239 | |
| 3240 | Ok(()) |
| 3241 | } |
| 3242 | |
| 3243 | /// The caller must hold the lock. This method was only added because the |
| 3244 | /// linter was complaining about the "cognitive complexity" of handle_sack. |
| 3245 | fn postprocess_sack( |
| 3246 | &mut self, |
| 3247 | state: AssociationState, |
| 3248 | mut should_awake_write_loop: bool, |
| 3249 | now: Instant, |
| 3250 | ) { |
| 3251 | if !self.inflight_queue.is_empty() { |
| 3252 | // Start timer. (noop if already started) |
| 3253 | trace!("[{}] T3-rtx timer start (pt3)", self.side); |
| 3254 | self.timers |
| 3255 | .restart_if_stale(Timer::T3RTX, now, self.rto_mgr.get_rto()); |
| 3256 | } else if state == AssociationState::ShutdownPending { |
| 3257 | // No more outstanding, send shutdown. |
| 3258 | should_awake_write_loop = true; |
| 3259 | self.will_send_shutdown = true; |
| 3260 | self.set_state(AssociationState::ShutdownSent); |
| 3261 | } else if state == AssociationState::ShutdownReceived { |
| 3262 | // No more outstanding, send shutdown ack. |
| 3263 | should_awake_write_loop = true; |
| 3264 | self.will_send_shutdown_ack = true; |
| 3265 | self.set_state(AssociationState::ShutdownAckSent); |
| 3266 | } |
| 3267 | |
| 3268 | if should_awake_write_loop { |
| 3269 | self.awake_write_loop(); |
| 3270 | } |
| 3271 | } |
| 3272 | |
| 3273 | fn reset_streams_if_any( |
| 3274 | &mut self, |
| 3275 | p: &ParamOutgoingResetRequest, |
| 3276 | from_wire: bool, |
| 3277 | reply: &mut Vec<Packet>, |
| 3278 | ) -> Result<()> { |
| 3279 | let mut result = ReconfigResult::SuccessPerformed; |
| 3280 | let mut sis_to_reset = vec![]; |
| 3281 | |
| 3282 | let performed = sna32lte(p.sender_last_tsn, self.peer_last_tsn); |
| 3283 | if performed { |
| 3284 | debug!( |
| 3285 | "[{}] resetStream(): senderLastTSN={} <= peer_last_tsn={}", |
| 3286 | self.side, p.sender_last_tsn, self.peer_last_tsn |
| 3287 | ); |
| 3288 | let mut stream_ids = if p.stream_identifiers.is_empty() { |
| 3289 | self.streams.keys().copied().collect::<Vec<_>>() |
| 3290 | } else { |
| 3291 | p.stream_identifiers.clone() |
| 3292 | }; |
| 3293 | stream_ids.sort_unstable(); |
| 3294 | stream_ids.dedup(); |
| 3295 | |
| 3296 | for id in stream_ids { |
| 3297 | if self.streams.contains_key(&id) { |
| 3298 | sis_to_reset.push(id); |
| 3299 | self.retire_stream(id, p.sender_last_tsn); |
| 3300 | } |
| 3301 | } |
| 3302 | self.reconfig_requests |
| 3303 | .remove(&p.reconfig_request_sequence_number); |
| 3304 | } else { |
| 3305 | debug!( |
| 3306 | "[{}] resetStream(): senderLastTSN={} > peer_last_tsn={}", |
| 3307 | self.side, p.sender_last_tsn, self.peer_last_tsn |
| 3308 | ); |
| 3309 | result = ReconfigResult::InProgress; |
| 3310 | } |
| 3311 | |
| 3312 | // RFC 8831 section 6.7 closes a bidirectional WebRTC data channel by |
| 3313 | // resetting the corresponding outgoing stream after its incoming half |
| 3314 | // is reset. `Stream` models that paired channel, so perform the |
| 3315 | // reciprocal reset automatically. |
| 3316 | // |
| 3317 | // The reciprocal also keeps each unregistered id pending in |
| 3318 | // `reconfigs`, which is what defers `StreamEvent::ResetComplete` |
| 3319 | // until the reciprocal is acknowledged (handle_reconfig_param) or |
| 3320 | // abandoned (on_retransmission_failure). |
| 3321 | if !sis_to_reset.is_empty() { |
| 3322 | let rsn = self.generate_next_rsn(); |
| 3323 | let tsn = self.my_next_tsn.wrapping_sub(1); |
| 3324 | let reset_all = p.stream_identifiers.is_empty(); |
| 3325 | self.reconfig_reset_streams |
| 3326 | .insert(rsn, sis_to_reset.clone()); |
| 3327 | |
| 3328 | let c = ChunkReconfig { |
| 3329 | param_a: Some(Box::new(ParamOutgoingResetRequest { |
| 3330 | reconfig_request_sequence_number: rsn, |
| 3331 | reconfig_response_sequence_number: p.reconfig_request_sequence_number, |
| 3332 | sender_last_tsn: tsn, |
| 3333 | stream_identifiers: if reset_all { vec![] } else { sis_to_reset }, |
| 3334 | })), |
| 3335 | ..Default::default() |
| 3336 | }; |
| 3337 | |
| 3338 | // Store before queueing. It becomes the active retransmission entry only |
| 3339 | // when gather_outbound actually serializes this packet. |
| 3340 | self.reconfigs.insert(rsn, c.clone()); |
| 3341 | |
| 3342 | let p = self.create_packet(vec![Box::new(c)]); |
| 3343 | reply.push(p); |
| 3344 | } |
| 3345 | |
| 3346 | if performed { |
| 3347 | self.release_deferred_reset_data()?; |
| 3348 | } |
| 3349 | |
| 3350 | // Respond to every request that arrived, fresh, retransmitted |
| 3351 | // or deferred by arrival of in-flight data. |
| 3352 | // |
| 3353 | // Intermediate re-evaluations that still fail due to in-flight data |
| 3354 | // stay silent rather than repeating "In progress" on every advance. |
| 3355 | if from_wire || performed { |
| 3356 | let packet = self.create_packet(vec![Box::new(ChunkReconfig { |
| 3357 | param_a: Some(Box::new(ParamReconfigResponse { |
| 3358 | reconfig_response_sequence_number: p.reconfig_request_sequence_number, |
| 3359 | result, |
| 3360 | })), |
| 3361 | param_b: None, |
| 3362 | })]); |
| 3363 | |
| 3364 | debug!("[{}] RESET RESPONSE: {}", self.side, packet); |
| 3365 | |
| 3366 | reply.push(packet); |
| 3367 | } |
| 3368 | |
| 3369 | Ok(()) |
| 3370 | } |
| 3371 | |
| 3372 | /// create_packet wraps chunks in a packet. |
| 3373 | /// The caller should hold the read lock. |
| 3374 | pub(crate) fn create_packet(&self, chunks: Vec<Box<dyn Chunk + Send + Sync>>) -> Packet { |
| 3375 | Packet { |
| 3376 | common_header: CommonHeader { |
| 3377 | verification_tag: self.peer_verification_tag, |
| 3378 | source_port: self.source_port, |
| 3379 | destination_port: self.destination_port, |
| 3380 | }, |
| 3381 | chunks, |
| 3382 | } |
| 3383 | } |
| 3384 | |
| 3385 | /// create_stream creates a stream. The caller should hold the lock |
| 3386 | /// and check no stream exists for this id. |
| 3387 | fn create_stream( |
| 3388 | &mut self, |
| 3389 | stream_identifier: StreamId, |
| 3390 | accept: bool, |
| 3391 | default_payload_type: PayloadProtocolIdentifier, |
| 3392 | ) -> Option<Stream<'_>> { |
| 3393 | if self |
| 3394 | .receive_limits |
| 3395 | .is_some_and(|limits| self.streams.len() >= limits.max_streams()) |
| 3396 | { |
| 3397 | return None; |
| 3398 | } |
| 3399 | let s = StreamState::new( |
| 3400 | self.side, |
| 3401 | stream_identifier, |
| 3402 | self.max_payload_size, |
| 3403 | self.max_receive_message_size, |
| 3404 | default_payload_type, |
| 3405 | ); |
| 3406 | |
| 3407 | if accept { |
| 3408 | self.stream_queue.push_back(stream_identifier); |
| 3409 | self.events.push_back(Event::Stream(StreamEvent::Opened { |
| 3410 | id: stream_identifier, |
| 3411 | })); |
| 3412 | } |
| 3413 | |
| 3414 | self.streams.insert(stream_identifier, s); |
| 3415 | |
| 3416 | Some(Stream { |
| 3417 | stream_identifier, |
| 3418 | association: self, |
| 3419 | }) |
| 3420 | } |
| 3421 | |
| 3422 | /// get_or_create_stream gets or creates a stream. The caller should hold the lock. |
| 3423 | fn get_or_create_stream(&mut self, stream_identifier: StreamId) -> Option<Stream<'_>> { |
| 3424 | if self.streams.contains_key(&stream_identifier) { |
| 3425 | Some(Stream { |
| 3426 | stream_identifier, |
| 3427 | association: self, |
| 3428 | }) |
| 3429 | } else { |
| 3430 | self.create_stream( |
| 3431 | stream_identifier, |
| 3432 | true, |
| 3433 | PayloadProtocolIdentifier::default(), |
| 3434 | ) |
| 3435 | } |
| 3436 | } |
| 3437 | |
| 3438 | /// Count each retained TSN once, including DATA already delivered to the |
| 3439 | /// application but still held behind a gap in the association payload queue. |
| 3440 | fn retained_receive_data(&self) -> (usize, usize) { |
| 3441 | let mut bytes = self.payload_queue.get_num_bytes(); |
| 3442 | let mut chunks = self.payload_queue.len(); |
| 3443 | for stream in self.streams.values() { |
| 3444 | for chunk in stream.reassembly_queue.chunks() { |
| 3445 | if self.payload_queue.get(chunk.tsn).is_none() { |
| 3446 | bytes = bytes.saturating_add(chunk.user_data.len()); |
| 3447 | chunks = chunks.saturating_add(1); |
| 3448 | } |
| 3449 | } |
| 3450 | } |
| 3451 | for chunk in self.deferred_reset_data.values() { |
| 3452 | if self.payload_queue.get(chunk.tsn).is_none() { |
| 3453 | bytes = bytes.saturating_add(chunk.user_data.len()); |
| 3454 | chunks = chunks.saturating_add(1); |
| 3455 | } |
| 3456 | } |
| 3457 | (bytes, chunks) |
| 3458 | } |
| 3459 | |
| 3460 | fn check_receive_limits(&self, data: &ChunkPayloadData) -> Result<()> { |
| 3461 | let Some(limits) = self.receive_limits else { |
| 3462 | return Ok(()); |
| 3463 | }; |
| 3464 | let (bytes, chunks) = self.retained_receive_data(); |
| 3465 | if data.user_data.len() > (limits.max_buffered_bytes() as usize).saturating_sub(bytes) |
| 3466 | || chunks >= limits.max_buffered_chunks() |
| 3467 | || (!self.streams.contains_key(&data.stream_identifier) |
| 3468 | && self.streams.len() >= limits.max_streams()) |
| 3469 | { |
| 3470 | return Err(Error::ErrReceiveLimitExceeded); |
| 3471 | } |
| 3472 | Ok(()) |
| 3473 | } |
| 3474 | |
| 3475 | pub(crate) fn get_my_receiver_window_credit(&self) -> u32 { |
| 3476 | if let Some(limits) = self.receive_limits { |
| 3477 | let (bytes, chunks) = self.retained_receive_data(); |
| 3478 | if chunks >= limits.max_buffered_chunks() { |
| 3479 | return 0; |
| 3480 | } |
| 3481 | return self.max_receive_buffer_size.saturating_sub(bytes as u32); |
| 3482 | } |
| 3483 | let mut bytes_queued = 0; |
| 3484 | for s in self.streams.values() { |
| 3485 | bytes_queued += s.get_num_bytes_in_reassembly_queue() as u32; |
| 3486 | } |
| 3487 | for chunk in self.deferred_reset_data.values() { |
| 3488 | bytes_queued += chunk.user_data.len() as u32; |
| 3489 | } |
| 3490 | |
| 3491 | self.max_receive_buffer_size.saturating_sub(bytes_queued) |
| 3492 | } |
| 3493 | |
| 3494 | /// gather_outbound gathers outgoing packets. The returned bool value set to |
| 3495 | /// false means the association should be closed down after the final send. |
| 3496 | fn gather_outbound(&mut self, now: Instant) -> (Vec<Bytes>, bool) { |
| 3497 | let mut raw_packets = self.gather_outbound_control_packets(vec![], now); |
| 3498 | |
| 3499 | let state = self.state(); |
| 3500 | match state { |
| 3501 | AssociationState::Established => { |
| 3502 | raw_packets = self.gather_data_packets_to_retransmit(raw_packets, now); |
| 3503 | raw_packets = self.gather_outbound_data_and_reconfig_packets(raw_packets, now); |
| 3504 | // A queued reciprocal reset may have been held back until the |
| 3505 | // DATA it covers received TSNs above. Give it another chance in |
| 3506 | // this same poll so DATA and RE-CONFIG can leave together. |
| 3507 | raw_packets = self.gather_outbound_control_packets(raw_packets, now); |
| 3508 | raw_packets = self.gather_outbound_fast_retransmission_packets(raw_packets, now); |
| 3509 | raw_packets = self.gather_outbound_sack_packets(raw_packets); |
| 3510 | raw_packets = self.gather_outbound_forward_tsn_packets(raw_packets); |
| 3511 | (raw_packets, true) |
| 3512 | } |
| 3513 | AssociationState::ShutdownPending |
| 3514 | | AssociationState::ShutdownSent |
| 3515 | | AssociationState::ShutdownReceived => { |
| 3516 | raw_packets = self.gather_data_packets_to_retransmit(raw_packets, now); |
| 3517 | raw_packets = self.gather_outbound_fast_retransmission_packets(raw_packets, now); |
| 3518 | raw_packets = self.gather_outbound_sack_packets(raw_packets); |
| 3519 | self.gather_outbound_shutdown_packets(raw_packets, now) |
| 3520 | } |
| 3521 | AssociationState::ShutdownAckSent => { |
| 3522 | self.gather_outbound_shutdown_packets(raw_packets, now) |
| 3523 | } |
| 3524 | _ => (raw_packets, true), |
| 3525 | } |
| 3526 | } |
| 3527 | |
| 3528 | fn gather_outbound_control_packets( |
| 3529 | &mut self, |
| 3530 | mut raw_packets: Vec<Bytes>, |
| 3531 | now: Instant, |
| 3532 | ) -> Vec<Bytes> { |
| 3533 | if !self.control_queue.is_empty() { |
| 3534 | let mut buffered = VecDeque::new(); |
| 3535 | let mut request_blocked = false; |
| 3536 | let queued = core::mem::take(&mut self.control_queue); |
| 3537 | for p in queued { |
| 3538 | let outgoing_rsn = Self::packet_reconfig_request_rsn(&p); |
| 3539 | |
| 3540 | if outgoing_rsn.is_some() && request_blocked { |
| 3541 | buffered.push_back(p); |
| 3542 | continue; |
| 3543 | } |
| 3544 | |
| 3545 | if outgoing_rsn.is_some_and(|rsn| { |
| 3546 | self.active_reconfig.is_some() && self.active_reconfig != Some(rsn) |
| 3547 | }) { |
| 3548 | request_blocked = true; |
| 3549 | buffered.push_back(p); |
| 3550 | continue; |
| 3551 | } |
| 3552 | |
| 3553 | // The reset boundary cannot be fixed until all already-queued |
| 3554 | // DATA for the affected streams has received a TSN. |
| 3555 | if outgoing_rsn.is_some_and(|rsn| self.reconfig_has_pending_data(rsn)) { |
| 3556 | request_blocked = true; |
| 3557 | buffered.push_back(p); |
| 3558 | continue; |
| 3559 | } |
| 3560 | |
| 3561 | let p = if let Some(rsn) = outgoing_rsn { |
| 3562 | self.refresh_unsent_reconfig(rsn).unwrap_or(p) |
| 3563 | } else { |
| 3564 | p |
| 3565 | }; |
| 3566 | |
| 3567 | match p.marshal() { |
| 3568 | Ok(raw) => { |
| 3569 | raw_packets.push(raw); |
| 3570 | if let Some(rsn) = outgoing_rsn { |
| 3571 | self.active_reconfig = Some(rsn); |
| 3572 | } |
| 3573 | } |
| 3574 | Err(_) => { |
| 3575 | warn!("[{}] failed to serialize a control packet", self.side); |
| 3576 | // A queued request owns reset state and must not be lost. |
| 3577 | if outgoing_rsn.is_some() { |
| 3578 | request_blocked = true; |
| 3579 | buffered.push_back(p); |
| 3580 | } |
| 3581 | } |
| 3582 | } |
| 3583 | } |
| 3584 | self.control_queue = buffered; |
| 3585 | } |
| 3586 | |
| 3587 | if self.active_reconfig.is_some() { |
| 3588 | self.timers |
| 3589 | .restart_if_stale(Timer::Reconfig, now, self.rto_mgr.get_rto()); |
| 3590 | } |
| 3591 | |
| 3592 | raw_packets |
| 3593 | } |
| 3594 | |
| 3595 | fn gather_data_packets_to_retransmit( |
| 3596 | &mut self, |
| 3597 | mut raw_packets: Vec<Bytes>, |
| 3598 | now: Instant, |
| 3599 | ) -> Vec<Bytes> { |
| 3600 | for p in &self.get_data_packets_to_retransmit(now) { |
| 3601 | if let Ok(raw) = p.marshal() { |
| 3602 | raw_packets.push(raw); |
| 3603 | } else { |
| 3604 | warn!( |
| 3605 | "[{}] failed to serialize a DATA packet to be retransmitted", |
| 3606 | self.side |
| 3607 | ); |
| 3608 | } |
| 3609 | } |
| 3610 | |
| 3611 | raw_packets |
| 3612 | } |
| 3613 | |
| 3614 | fn gather_outbound_data_and_reconfig_packets( |
| 3615 | &mut self, |
| 3616 | mut raw_packets: Vec<Bytes>, |
| 3617 | now: Instant, |
| 3618 | ) -> Vec<Bytes> { |
| 3619 | // Pop unsent data chunks from the pending queue to send as much as |
| 3620 | // cwnd and rwnd allow. |
| 3621 | let chunks = self.pop_pending_data_chunks_to_send(now); |
| 3622 | if !chunks.is_empty() { |
| 3623 | // Start timer. (noop if already started) |
| 3624 | trace!("[{}] T3-rtx timer start (pt1)", self.side); |
| 3625 | self.timers |
| 3626 | .restart_if_stale(Timer::T3RTX, now, self.rto_mgr.get_rto()); |
| 3627 | |
| 3628 | for p in &self.bundle_data_chunks_into_packets(chunks) { |
| 3629 | if let Ok(raw) = p.marshal() { |
| 3630 | raw_packets.push(raw); |
| 3631 | } else { |
| 3632 | warn!("[{}] failed to serialize a DATA packet", self.side); |
| 3633 | } |
| 3634 | } |
| 3635 | } |
| 3636 | |
| 3637 | if self.will_retransmit_reconfig { |
| 3638 | self.will_retransmit_reconfig = false; |
| 3639 | if let Some(rsn) = self.active_reconfig { |
| 3640 | debug!("[{}] retransmit RECONFIG rsn={}", self.side, rsn); |
| 3641 | if let Some(c) = self.reconfigs.get(&rsn) { |
| 3642 | let p = self.create_packet(vec![Box::new(c.clone())]); |
| 3643 | if let Ok(raw) = p.marshal() { |
| 3644 | raw_packets.push(raw); |
| 3645 | } else { |
| 3646 | warn!( |
| 3647 | "[{}] failed to serialize a RECONFIG packet to be retransmitted", |
| 3648 | self.side, |
| 3649 | ); |
| 3650 | } |
| 3651 | } |
| 3652 | } |
| 3653 | } |
| 3654 | |
| 3655 | // RFC 6525 section 5.1.1 permits only one request in flight. Keep |
| 3656 | // application resets separate from DATA until that request completes. |
| 3657 | if self.active_reconfig.is_none() |
| 3658 | && self.reconfigs.is_empty() |
| 3659 | && !self.pending_reset_streams.is_empty() |
| 3660 | { |
| 3661 | // DATA on an unrelated stream must not starve a reset, nor should |
| 3662 | // one busy stream hold back ready reset requests for other streams. |
| 3663 | let pending_queue = &self.pending_queue; |
| 3664 | let mut stream_ids = vec![]; |
| 3665 | self.pending_reset_streams.retain(|id| { |
| 3666 | if pending_queue.contains_stream(*id) { |
| 3667 | true |
| 3668 | } else { |
| 3669 | stream_ids.push(*id); |
| 3670 | false |
| 3671 | } |
| 3672 | }); |
| 3673 | |
| 3674 | if stream_ids.is_empty() { |
| 3675 | return raw_packets; |
| 3676 | } |
| 3677 | |
| 3678 | let rsn = self.generate_next_rsn(); |
| 3679 | let tsn = self.my_next_tsn.wrapping_sub(1); |
| 3680 | debug!( |
| 3681 | "[{}] sending RECONFIG: rsn={} tsn={} streams={:?}", |
| 3682 | self.side, rsn, tsn, stream_ids |
| 3683 | ); |
| 3684 | |
| 3685 | self.reconfig_reset_streams.insert(rsn, stream_ids.clone()); |
| 3686 | let c = ChunkReconfig { |
| 3687 | param_a: Some(Box::new(ParamOutgoingResetRequest { |
| 3688 | reconfig_request_sequence_number: rsn, |
| 3689 | reconfig_response_sequence_number: self.peer_last_reconfig_rsn, |
| 3690 | sender_last_tsn: tsn, |
| 3691 | stream_identifiers: stream_ids, |
| 3692 | })), |
| 3693 | ..Default::default() |
| 3694 | }; |
| 3695 | self.reconfigs.insert(rsn, c.clone()); |
| 3696 | |
| 3697 | let p = self.create_packet(vec![Box::new(c)]); |
| 3698 | if let Ok(raw) = p.marshal() { |
| 3699 | self.active_reconfig = Some(rsn); |
| 3700 | raw_packets.push(raw); |
| 3701 | } else { |
| 3702 | warn!( |
| 3703 | "[{}] failed to serialize a RECONFIG packet to be transmitted", |
| 3704 | self.side |
| 3705 | ); |
| 3706 | // Preserve the request for a later serialization attempt. |
| 3707 | self.control_queue.push_back(p); |
| 3708 | } |
| 3709 | } |
| 3710 | |
| 3711 | // The timer belongs to exactly the request selected above or while |
| 3712 | // draining the control queue. |
| 3713 | if self.active_reconfig.is_some() { |
| 3714 | self.timers |
| 3715 | .restart_if_stale(Timer::Reconfig, now, self.rto_mgr.get_rto()); |
| 3716 | } |
| 3717 | |
| 3718 | raw_packets |
| 3719 | } |
| 3720 | |
| 3721 | fn gather_outbound_fast_retransmission_packets( |
| 3722 | &mut self, |
| 3723 | mut raw_packets: Vec<Bytes>, |
| 3724 | now: Instant, |
| 3725 | ) -> Vec<Bytes> { |
| 3726 | if self.will_retransmit_fast { |
| 3727 | self.will_retransmit_fast = false; |
| 3728 | |
| 3729 | let mut to_fast_retrans: Vec<Box<dyn Chunk + Send + Sync>> = vec![]; |
| 3730 | let mut fast_retrans_size = COMMON_HEADER_SIZE; |
| 3731 | |
| 3732 | let mut i = 0; |
| 3733 | loop { |
| 3734 | let tsn = self.cumulative_tsn_ack_point + i + 1; |
| 3735 | if let Some(c) = self.inflight_queue.get_mut(tsn) { |
| 3736 | if c.acked || c.abandoned() || c.nsent > 1 || c.miss_indicator < 3 { |
| 3737 | i += 1; |
| 3738 | continue; |
| 3739 | } |
| 3740 | |
| 3741 | // RFC 4960 Sec 7.2.4 Fast Retransmit on Gap Reports |
| 3742 | // 3) Determine how many of the earliest (i.e., lowest TSN) DATA chunks |
| 3743 | // marked for retransmission will fit into a single packet, subject |
| 3744 | // to constraint of the path MTU of the destination transport |
| 3745 | // address to which the packet is being sent. Call this value K. |
| 3746 | // Retransmit those K DATA chunks in a single packet. When a Fast |
| 3747 | // Retransmit is being performed, the sender SHOULD ignore the value |
| 3748 | // of cwnd and SHOULD NOT delay retransmission for this single |
| 3749 | // packet. |
| 3750 | |
| 3751 | let data_chunk_size = DATA_CHUNK_HEADER_SIZE + c.user_data.len() as u32; |
| 3752 | if self.mtu < fast_retrans_size + data_chunk_size { |
| 3753 | break; |
| 3754 | } |
| 3755 | |
| 3756 | fast_retrans_size += data_chunk_size; |
| 3757 | self.stats.inc_fast_retrans(); |
| 3758 | c.nsent += 1; |
| 3759 | } else { |
| 3760 | break; // end of pending data |
| 3761 | } |
| 3762 | |
| 3763 | if let Some(c) = self.inflight_queue.get_mut(tsn) { |
| 3764 | Association::check_partial_reliability_status( |
| 3765 | c, |
| 3766 | now, |
| 3767 | self.use_forward_tsn, |
| 3768 | self.side, |
| 3769 | &self.streams, |
| 3770 | ); |
| 3771 | to_fast_retrans.push(Box::new(c.clone())); |
| 3772 | trace!( |
| 3773 | "[{}] fast-retransmit: tsn={} sent={} htna={}", |
| 3774 | self.side, c.tsn, c.nsent, self.fast_recover_exit_point |
| 3775 | ); |
| 3776 | } |
| 3777 | i += 1; |
| 3778 | } |
| 3779 | |
| 3780 | if !to_fast_retrans.is_empty() { |
| 3781 | if let Ok(raw) = self.create_packet(to_fast_retrans).marshal() { |
| 3782 | raw_packets.push(raw); |
| 3783 | } else { |
| 3784 | warn!( |
| 3785 | "[{}] failed to serialize a DATA packet to be fast-retransmitted", |
| 3786 | self.side |
| 3787 | ); |
| 3788 | } |
| 3789 | } |
| 3790 | } |
| 3791 | |
| 3792 | raw_packets |
| 3793 | } |
| 3794 | |
| 3795 | fn gather_outbound_sack_packets(&mut self, mut raw_packets: Vec<Bytes>) -> Vec<Bytes> { |
| 3796 | if self.ack_state == AckState::Immediate { |
| 3797 | self.ack_state = AckState::Idle; |
| 3798 | let sack = self.create_selective_ack_chunk(); |
| 3799 | trace!("[{}] sending SACK: {}", self.side, sack); |
| 3800 | if let Ok(raw) = self.create_packet(vec![Box::new(sack)]).marshal() { |
| 3801 | raw_packets.push(raw); |
| 3802 | } else { |
| 3803 | warn!("[{}] failed to serialize a SACK packet", self.side); |
| 3804 | } |
| 3805 | } |
| 3806 | |
| 3807 | raw_packets |
| 3808 | } |
| 3809 | |
| 3810 | fn gather_outbound_forward_tsn_packets(&mut self, mut raw_packets: Vec<Bytes>) -> Vec<Bytes> { |
| 3811 | /*log::debug!( |
| 3812 | "[{}] gatherOutboundForwardTSNPackets {}", |
| 3813 | self.name, |
| 3814 | self.will_send_forward_tsn |
| 3815 | );*/ |
| 3816 | if self.will_send_forward_tsn { |
| 3817 | self.will_send_forward_tsn = false; |
| 3818 | if sna32gt( |
| 3819 | self.advanced_peer_tsn_ack_point, |
| 3820 | self.cumulative_tsn_ack_point, |
| 3821 | ) { |
| 3822 | let fwd_tsn = self.create_forward_tsn(); |
| 3823 | if let Ok(raw) = self.create_packet(vec![Box::new(fwd_tsn)]).marshal() { |
| 3824 | raw_packets.push(raw); |
| 3825 | } else { |
| 3826 | warn!("[{}] failed to serialize a Forward TSN packet", self.side); |
| 3827 | } |
| 3828 | } |
| 3829 | } |
| 3830 | |
| 3831 | raw_packets |
| 3832 | } |
| 3833 | |
| 3834 | fn gather_outbound_shutdown_packets( |
| 3835 | &mut self, |
| 3836 | mut raw_packets: Vec<Bytes>, |
| 3837 | now: Instant, |
| 3838 | ) -> (Vec<Bytes>, bool) { |
| 3839 | let mut ok = true; |
| 3840 | |
| 3841 | if self.will_send_shutdown { |
| 3842 | self.will_send_shutdown = false; |
| 3843 | |
| 3844 | let shutdown = ChunkShutdown { |
| 3845 | cumulative_tsn_ack: self.cumulative_tsn_ack_point, |
| 3846 | }; |
| 3847 | |
| 3848 | if let Ok(raw) = self.create_packet(vec![Box::new(shutdown)]).marshal() { |
| 3849 | self.timers |
| 3850 | .start(Timer::T2Shutdown, now, self.rto_mgr.get_rto()); |
| 3851 | raw_packets.push(raw); |
| 3852 | } else { |
| 3853 | warn!("[{}] failed to serialize a Shutdown packet", self.side); |
| 3854 | } |
| 3855 | } else if self.will_send_shutdown_ack { |
| 3856 | self.will_send_shutdown_ack = false; |
| 3857 | |
| 3858 | let shutdown_ack = ChunkShutdownAck {}; |
| 3859 | |
| 3860 | if let Ok(raw) = self.create_packet(vec![Box::new(shutdown_ack)]).marshal() { |
| 3861 | self.timers |
| 3862 | .start(Timer::T2Shutdown, now, self.rto_mgr.get_rto()); |
| 3863 | raw_packets.push(raw); |
| 3864 | } else { |
| 3865 | warn!("[{}] failed to serialize a ShutdownAck packet", self.side); |
| 3866 | } |
| 3867 | } else if self.will_send_shutdown_complete { |
| 3868 | self.will_send_shutdown_complete = false; |
| 3869 | |
| 3870 | let shutdown_complete = ChunkShutdownComplete {}; |
| 3871 | |
| 3872 | if let Ok(raw) = self |
| 3873 | .create_packet(vec![Box::new(shutdown_complete)]) |
| 3874 | .marshal() |
| 3875 | { |
| 3876 | raw_packets.push(raw); |
| 3877 | ok = false; |
| 3878 | } else { |
| 3879 | warn!( |
| 3880 | "[{}] failed to serialize a ShutdownComplete packet", |
| 3881 | self.side |
| 3882 | ); |
| 3883 | } |
| 3884 | } |
| 3885 | |
| 3886 | (raw_packets, ok) |
| 3887 | } |
| 3888 | |
| 3889 | /// get_data_packets_to_retransmit is called when T3-rtx is timed |
| 3890 | /// out and retransmit outstanding data chunks that are not acked |
| 3891 | /// or abandoned yet. |
| 3892 | fn get_data_packets_to_retransmit(&mut self, now: Instant) -> Vec<Packet> { |
| 3893 | let awnd = core::cmp::min(self.cwnd, self.rwnd); |
| 3894 | let mut chunks = vec![]; |
| 3895 | let mut bytes_to_send = 0; |
| 3896 | let mut done = false; |
| 3897 | let mut i = 0; |
| 3898 | while !done { |
| 3899 | let tsn = self.cumulative_tsn_ack_point + i + 1; |
| 3900 | if let Some(c) = self.inflight_queue.get_mut(tsn) { |
| 3901 | if !c.retransmit { |
| 3902 | i += 1; |
| 3903 | continue; |
| 3904 | } |
| 3905 | |
| 3906 | if i == 0 && self.rwnd < c.user_data.len() as u32 { |
| 3907 | // Send it as a zero window probe |
| 3908 | done = true; |
| 3909 | } else if bytes_to_send + c.user_data.len() > awnd as usize { |
| 3910 | break; |
| 3911 | } |
| 3912 | |
| 3913 | // reset the retransmit flag not to retransmit again before the next |
| 3914 | // t3-rtx timer fires |
| 3915 | c.retransmit = false; |
| 3916 | bytes_to_send += c.user_data.len(); |
| 3917 | |
| 3918 | c.nsent += 1; |
| 3919 | } else { |
| 3920 | break; // end of pending data |
| 3921 | } |
| 3922 | |
| 3923 | if let Some(c) = self.inflight_queue.get_mut(tsn) { |
| 3924 | Association::check_partial_reliability_status( |
| 3925 | c, |
| 3926 | now, |
| 3927 | self.use_forward_tsn, |
| 3928 | self.side, |
| 3929 | &self.streams, |
| 3930 | ); |
| 3931 | |
| 3932 | trace!( |
| 3933 | "[{}] retransmitting tsn={} ssn={} sent={}", |
| 3934 | self.side, c.tsn, c.stream_sequence_number, c.nsent |
| 3935 | ); |
| 3936 | |
| 3937 | chunks.push(c.clone()); |
| 3938 | } |
| 3939 | i += 1; |
| 3940 | } |
| 3941 | |
| 3942 | self.bundle_data_chunks_into_packets(chunks) |
| 3943 | } |
| 3944 | |
| 3945 | /// pop_pending_data_chunks_to_send pops chunks from the pending queues as many as |
| 3946 | /// the cwnd and rwnd allows to send. |
| 3947 | fn pop_pending_data_chunks_to_send(&mut self, now: Instant) -> Vec<ChunkPayloadData> { |
| 3948 | let mut chunks = vec![]; |
| 3949 | if !self.pending_queue.is_empty() { |
| 3950 | // RFC 4960 sec 6.1. Transmission of DATA Chunks |
| 3951 | // A) At any given time, the data sender MUST NOT transmit new data to |
| 3952 | // any destination transport address if its peer's rwnd indicates |
| 3953 | // that the peer has no buffer space (i.e., rwnd is 0; see Section |
| 3954 | // 6.2.1). However, regardless of the value of rwnd (including if it |
| 3955 | // is 0), the data sender can always have one DATA chunk in flight to |
| 3956 | // the receiver if allowed by cwnd (see rule B, below). |
| 3957 | |
| 3958 | while let Some(c) = self.pending_queue.peek() { |
| 3959 | let (beginning_fragment, unordered, data_len) = |
| 3960 | (c.beginning_fragment, c.unordered, c.user_data.len()); |
| 3961 | |
| 3962 | if self.inflight_queue.get_num_bytes() + data_len > self.cwnd as usize { |
| 3963 | break; // would exceeds cwnd |
| 3964 | } |
| 3965 | |
| 3966 | if data_len > self.rwnd as usize { |
| 3967 | break; // no more rwnd |
| 3968 | } |
| 3969 | |
| 3970 | self.rwnd -= data_len as u32; |
| 3971 | |
| 3972 | if let Some(chunk) = self.move_pending_data_chunk_to_inflight_queue( |
| 3973 | beginning_fragment, |
| 3974 | unordered, |
| 3975 | now, |
| 3976 | ) { |
| 3977 | chunks.push(chunk); |
| 3978 | } |
| 3979 | } |
| 3980 | |
| 3981 | // the data sender can always have one DATA chunk in flight to the receiver |
| 3982 | if chunks.is_empty() && self.inflight_queue.is_empty() { |
| 3983 | // Send zero window probe |
| 3984 | if let Some(c) = self.pending_queue.peek() { |
| 3985 | let (beginning_fragment, unordered) = (c.beginning_fragment, c.unordered); |
| 3986 | |
| 3987 | if let Some(chunk) = self.move_pending_data_chunk_to_inflight_queue( |
| 3988 | beginning_fragment, |
| 3989 | unordered, |
| 3990 | now, |
| 3991 | ) { |
| 3992 | chunks.push(chunk); |
| 3993 | } |
| 3994 | } |
| 3995 | } |
| 3996 | } |
| 3997 | |
| 3998 | chunks |
| 3999 | } |
| 4000 | |
| 4001 | /// bundle_data_chunks_into_packets packs DATA chunks into packets. It tries to bundle |
| 4002 | /// DATA chunks into a packet so long as the resulting packet size does not exceed |
| 4003 | /// the path MTU. |
| 4004 | fn bundle_data_chunks_into_packets(&self, chunks: Vec<ChunkPayloadData>) -> Vec<Packet> { |
| 4005 | let mut packets = vec![]; |
| 4006 | let mut chunks_to_send = vec![]; |
| 4007 | let mut bytes_in_packet = COMMON_HEADER_SIZE; |
| 4008 | |
| 4009 | for c in chunks { |
| 4010 | // RFC 4960 sec 6.1. Transmission of DATA Chunks |
| 4011 | // Multiple DATA chunks committed for transmission MAY be bundled in a |
| 4012 | // single packet. Furthermore, DATA chunks being retransmitted MAY be |
| 4013 | // bundled with new DATA chunks, as long as the resulting packet size |
| 4014 | // does not exceed the path MTU. |
| 4015 | if bytes_in_packet + c.user_data.len() as u32 > self.mtu { |
| 4016 | packets.push(self.create_packet(chunks_to_send)); |
| 4017 | chunks_to_send = vec![]; |
| 4018 | bytes_in_packet = COMMON_HEADER_SIZE; |
| 4019 | } |
| 4020 | |
| 4021 | bytes_in_packet += DATA_CHUNK_HEADER_SIZE + c.user_data.len() as u32; |
| 4022 | chunks_to_send.push(Box::new(c)); |
| 4023 | } |
| 4024 | |
| 4025 | if !chunks_to_send.is_empty() { |
| 4026 | packets.push(self.create_packet(chunks_to_send)); |
| 4027 | } |
| 4028 | |
| 4029 | packets |
| 4030 | } |
| 4031 | |
| 4032 | /// generate_next_tsn returns the my_next_tsn and increases it. The caller should hold the lock. |
| 4033 | fn generate_next_tsn(&mut self) -> u32 { |
| 4034 | let tsn = self.my_next_tsn; |
| 4035 | self.my_next_tsn = self.my_next_tsn.wrapping_add(1); |
| 4036 | tsn |
| 4037 | } |
| 4038 | |
| 4039 | /// generate_next_rsn returns the my_next_rsn and increases it. The caller should hold the lock. |
| 4040 | fn generate_next_rsn(&mut self) -> u32 { |
| 4041 | let rsn = self.my_next_rsn; |
| 4042 | self.my_next_rsn = self.my_next_rsn.wrapping_add(1); |
| 4043 | rsn |
| 4044 | } |
| 4045 | |
| 4046 | fn check_partial_reliability_status( |
| 4047 | c: &mut ChunkPayloadData, |
| 4048 | now: Instant, |
| 4049 | use_forward_tsn: bool, |
| 4050 | side: Side, |
| 4051 | streams: &FxHashMap<u16, StreamState>, |
| 4052 | ) { |
| 4053 | if !use_forward_tsn { |
| 4054 | return; |
| 4055 | } |
| 4056 | |
| 4057 | // draft-ietf-rtcweb-data-protocol-09.txt section 6 |
| 4058 | // 6. Procedures |
| 4059 | // All Data Channel Establishment Protocol messages MUST be sent using |
| 4060 | // ordered delivery and reliable transmission. |
| 4061 | // |
| 4062 | if c.payload_type == PayloadProtocolIdentifier::Dcep { |
| 4063 | return; |
| 4064 | } |
| 4065 | |
| 4066 | // PR-SCTP |
| 4067 | if let Some(s) = streams.get(&c.stream_identifier) { |
| 4068 | let reliability_type: ReliabilityType = s.reliability_type; |
| 4069 | let reliability_value = s.reliability_value; |
| 4070 | |
| 4071 | if reliability_type == ReliabilityType::Rexmit { |
| 4072 | if c.nsent >= reliability_value { |
| 4073 | c.set_abandoned(true); |
| 4074 | trace!( |
| 4075 | "[{}] marked as abandoned: tsn={} ppi={} (remix: {})", |
| 4076 | side, c.tsn, c.payload_type, c.nsent |
| 4077 | ); |
| 4078 | } |
| 4079 | } else if reliability_type == ReliabilityType::Timed { |
| 4080 | if let Some(since) = &c.since { |
| 4081 | let elapsed = now.duration_since(*since); |
| 4082 | if elapsed.as_millis() as u32 >= reliability_value { |
| 4083 | c.set_abandoned(true); |
| 4084 | trace!( |
| 4085 | "[{}] marked as abandoned: tsn={} ppi={} (timed: {:?})", |
| 4086 | side, c.tsn, c.payload_type, elapsed |
| 4087 | ); |
| 4088 | } |
| 4089 | } else { |
| 4090 | error!("[{}] invalid c.since", side); |
| 4091 | } |
| 4092 | } |
| 4093 | } else { |
| 4094 | error!("[{}] stream {} not found)", side, c.stream_identifier); |
| 4095 | } |
| 4096 | } |
| 4097 | |
| 4098 | fn create_selective_ack_chunk(&mut self) -> ChunkSelectiveAck { |
| 4099 | ChunkSelectiveAck { |
| 4100 | cumulative_tsn_ack: self.peer_last_tsn, |
| 4101 | advertised_receiver_window_credit: self.get_my_receiver_window_credit(), |
| 4102 | gap_ack_blocks: self.payload_queue.get_gap_ack_blocks(self.peer_last_tsn), |
| 4103 | duplicate_tsn: self.payload_queue.pop_duplicates(), |
| 4104 | } |
| 4105 | } |
| 4106 | |
| 4107 | /// create_forward_tsn generates ForwardTSN chunk. |
| 4108 | /// This method is called when use_forward_tsn is set to true. |
| 4109 | fn create_forward_tsn(&self) -> ChunkForwardTsn { |
| 4110 | // RFC 3758 Sec 3.5 C4 |
| 4111 | let mut stream_map: HashMap<u16, u16> = HashMap::new(); // to report only once per SI |
| 4112 | let mut i = self.cumulative_tsn_ack_point + 1; |
| 4113 | while sna32lte(i, self.advanced_peer_tsn_ack_point) { |
| 4114 | if let Some(c) = self.inflight_queue.get(i) { |
| 4115 | if let Some(ssn) = stream_map.get(&c.stream_identifier) { |
| 4116 | if sna16lt(*ssn, c.stream_sequence_number) { |
| 4117 | // to report only once with greatest SSN |
| 4118 | stream_map.insert(c.stream_identifier, c.stream_sequence_number); |
| 4119 | } |
| 4120 | } else { |
| 4121 | stream_map.insert(c.stream_identifier, c.stream_sequence_number); |
| 4122 | } |
| 4123 | } else { |
| 4124 | break; |
| 4125 | } |
| 4126 | |
| 4127 | i += 1; |
| 4128 | } |
| 4129 | |
| 4130 | let mut fwd_tsn = ChunkForwardTsn { |
| 4131 | new_cumulative_tsn: self.advanced_peer_tsn_ack_point, |
| 4132 | streams: vec![], |
| 4133 | }; |
| 4134 | |
| 4135 | let mut stream_str = String::new(); |
| 4136 | for (si, ssn) in &stream_map { |
| 4137 | stream_str += format!("(si={} ssn={})", si, ssn).as_str(); |
| 4138 | fwd_tsn.streams.push(ChunkForwardTsnStream { |
| 4139 | identifier: *si, |
| 4140 | sequence: *ssn, |
| 4141 | }); |
| 4142 | } |
| 4143 | trace!( |
| 4144 | "[{}] building fwd_tsn: newCumulativeTSN={} cumTSN={} - {}", |
| 4145 | self.side, fwd_tsn.new_cumulative_tsn, self.cumulative_tsn_ack_point, stream_str |
| 4146 | ); |
| 4147 | |
| 4148 | fwd_tsn |
| 4149 | } |
| 4150 | |
| 4151 | /// Move the chunk peeked with self.pending_queue.peek() to the inflight_queue. |
| 4152 | fn move_pending_data_chunk_to_inflight_queue( |
| 4153 | &mut self, |
| 4154 | beginning_fragment: bool, |
| 4155 | unordered: bool, |
| 4156 | now: Instant, |
| 4157 | ) -> Option<ChunkPayloadData> { |
| 4158 | if let Some(mut c) = self.pending_queue.pop(beginning_fragment, unordered) { |
| 4159 | // Mark all fragements are in-flight now |
| 4160 | if c.ending_fragment { |
| 4161 | c.set_all_inflight(); |
| 4162 | } |
| 4163 | |
| 4164 | // Assign TSN |
| 4165 | c.tsn = self.generate_next_tsn(); |
| 4166 | |
| 4167 | c.since = Some(now); // use to calculate RTT and also for maxPacketLifeTime |
| 4168 | c.nsent = 1; // being sent for the first time |
| 4169 | |
| 4170 | Association::check_partial_reliability_status( |
| 4171 | &mut c, |
| 4172 | now, |
| 4173 | self.use_forward_tsn, |
| 4174 | self.side, |
| 4175 | &self.streams, |
| 4176 | ); |
| 4177 | |
| 4178 | trace!( |
| 4179 | "[{}] sending ppi={} tsn={} ssn={} sent={} len={} ({},{})", |
| 4180 | self.side, |
| 4181 | c.payload_type as u32, |
| 4182 | c.tsn, |
| 4183 | c.stream_sequence_number, |
| 4184 | c.nsent, |
| 4185 | c.user_data.len(), |
| 4186 | c.beginning_fragment, |
| 4187 | c.ending_fragment |
| 4188 | ); |
| 4189 | |
| 4190 | self.inflight_queue.push_no_check(c.clone()); |
| 4191 | |
| 4192 | Some(c) |
| 4193 | } else { |
| 4194 | error!("[{}] failed to pop from pending queue", self.side); |
| 4195 | None |
| 4196 | } |
| 4197 | } |
| 4198 | |
| 4199 | pub(crate) fn send_reset_request(&mut self, stream_identifier: StreamId) -> Result<()> { |
| 4200 | let state = self.state(); |
| 4201 | if state != AssociationState::Established { |
| 4202 | return Err(Error::ErrResetPacketInStateNotExist); |
| 4203 | } |
| 4204 | |
| 4205 | self.pending_reset_streams.push_back(stream_identifier); |
| 4206 | self.awake_write_loop(); |
| 4207 | |
| 4208 | Ok(()) |
| 4209 | } |
| 4210 | |
| 4211 | /// send_payload_data sends the data chunks. |
| 4212 | pub(crate) fn send_payload_data(&mut self, chunks: Vec<ChunkPayloadData>) -> Result<()> { |
| 4213 | let state = self.state(); |
| 4214 | if state != AssociationState::Established { |
| 4215 | return Err(Error::ErrPayloadDataStateNotExist); |
| 4216 | } |
| 4217 | |
| 4218 | // Push the chunks into the pending queue first. |
| 4219 | for c in chunks { |
| 4220 | self.pending_queue.push(c); |
| 4221 | } |
| 4222 | |
| 4223 | self.awake_write_loop(); |
| 4224 | Ok(()) |
| 4225 | } |
| 4226 | |
| 4227 | /// buffered_amount returns total amount (in bytes) of currently buffered user data. |
| 4228 | /// This is used only by testing. |
| 4229 | pub(crate) fn buffered_amount(&self) -> usize { |
| 4230 | self.pending_queue.get_num_bytes() + self.inflight_queue.get_num_bytes() |
| 4231 | } |
| 4232 | |
| 4233 | fn awake_write_loop(&self) { |
| 4234 | // No Op on Purpose |
| 4235 | } |
| 4236 | |
| 4237 | fn close_all_timers(&mut self) { |
| 4238 | // Close all retransmission & ack timers |
| 4239 | for timer in Timer::VALUES { |
| 4240 | self.timers.stop(timer); |
| 4241 | } |
| 4242 | } |
| 4243 | |
| 4244 | fn on_ack_timeout(&mut self) { |
| 4245 | trace!( |
| 4246 | "[{}] ack timed out (ack_state: {})", |
| 4247 | self.side, self.ack_state |
| 4248 | ); |
| 4249 | self.stats.inc_ack_timeouts(); |
| 4250 | self.ack_state = AckState::Immediate; |
| 4251 | self.awake_write_loop(); |
| 4252 | } |
| 4253 | |
| 4254 | fn on_retransmission_timeout(&mut self, timer_id: Timer, n_rtos: usize) { |
| 4255 | match timer_id { |
| 4256 | Timer::T1Init => { |
| 4257 | if let Err(err) = self.send_init() { |
| 4258 | debug!( |
| 4259 | "[{}] failed to retransmit init (n_rtos={}): {:?}", |
| 4260 | self.side, n_rtos, err |
| 4261 | ); |
| 4262 | } |
| 4263 | } |
| 4264 | |
| 4265 | Timer::T1Cookie => { |
| 4266 | if let Err(err) = self.send_cookie_echo() { |
| 4267 | debug!( |
| 4268 | "[{}] failed to retransmit cookie-echo (n_rtos={}): {:?}", |
| 4269 | self.side, n_rtos, err |
| 4270 | ); |
| 4271 | } |
| 4272 | } |
| 4273 | |
| 4274 | Timer::T2Shutdown => { |
| 4275 | debug!( |
| 4276 | "[{}] retransmission of shutdown timeout (n_rtos={})", |
| 4277 | self.side, n_rtos |
| 4278 | ); |
| 4279 | let state = self.state(); |
| 4280 | match state { |
| 4281 | AssociationState::ShutdownSent => { |
| 4282 | self.will_send_shutdown = true; |
| 4283 | self.awake_write_loop(); |
| 4284 | } |
| 4285 | AssociationState::ShutdownAckSent => { |
| 4286 | self.will_send_shutdown_ack = true; |
| 4287 | self.awake_write_loop(); |
| 4288 | } |
| 4289 | _ => {} |
| 4290 | } |
| 4291 | } |
| 4292 | |
| 4293 | Timer::T3RTX => { |
| 4294 | self.stats.inc_t3timeouts(); |
| 4295 | |
| 4296 | // RFC 4960 sec 6.3.3 |
| 4297 | // E1) For the destination address for which the timer expires, adjust |
| 4298 | // its ssthresh with rules defined in Section 7.2.3 and set the |
| 4299 | // cwnd <- MTU. |
| 4300 | // RFC 4960 sec 7.2.3 |
| 4301 | // When the T3-rtx timer expires on an address, SCTP should perform slow |
| 4302 | // start by: |
| 4303 | // ssthresh = max(cwnd/2, 4*MTU) |
| 4304 | // cwnd = 1*MTU |
| 4305 | |
| 4306 | self.ssthresh = core::cmp::max(self.cwnd / 2, 4 * self.mtu); |
| 4307 | self.cwnd = self.mtu; |
| 4308 | trace!( |
| 4309 | "[{}] updated cwnd={} ssthresh={} inflight={} (RTO)", |
| 4310 | self.side, |
| 4311 | self.cwnd, |
| 4312 | self.ssthresh, |
| 4313 | self.inflight_queue.get_num_bytes() |
| 4314 | ); |
| 4315 | |
| 4316 | // RFC 3758 sec 3.5 |
| 4317 | // A5) Any time the T3-rtx timer expires, on any destination, the sender |
| 4318 | // SHOULD try to advance the "Advanced.Peer.Ack.Point" by following |
| 4319 | // the procedures outlined in C2 - C5. |
| 4320 | if self.use_forward_tsn { |
| 4321 | // RFC 3758 Sec 3.5 C2 |
| 4322 | let mut i = self.advanced_peer_tsn_ack_point + 1; |
| 4323 | while let Some(c) = self.inflight_queue.get(i) { |
| 4324 | if !c.abandoned() { |
| 4325 | break; |
| 4326 | } |
| 4327 | self.advanced_peer_tsn_ack_point = i; |
| 4328 | i += 1; |
| 4329 | } |
| 4330 | |
| 4331 | // RFC 3758 Sec 3.5 C3 |
| 4332 | if sna32gt( |
| 4333 | self.advanced_peer_tsn_ack_point, |
| 4334 | self.cumulative_tsn_ack_point, |
| 4335 | ) { |
| 4336 | self.will_send_forward_tsn = true; |
| 4337 | debug!( |
| 4338 | "[{}] on_retransmission_timeout {}: sna32GT({}, {})", |
| 4339 | self.side, |
| 4340 | self.will_send_forward_tsn, |
| 4341 | self.advanced_peer_tsn_ack_point, |
| 4342 | self.cumulative_tsn_ack_point |
| 4343 | ); |
| 4344 | } |
| 4345 | } |
| 4346 | |
| 4347 | debug!( |
| 4348 | "[{}] T3-rtx timed out: n_rtos={} cwnd={} ssthresh={}", |
| 4349 | self.side, n_rtos, self.cwnd, self.ssthresh |
| 4350 | ); |
| 4351 | |
| 4352 | self.inflight_queue.mark_all_to_retrasmit(); |
| 4353 | self.awake_write_loop(); |
| 4354 | } |
| 4355 | |
| 4356 | Timer::Reconfig => { |
| 4357 | self.will_retransmit_reconfig = true; |
| 4358 | self.awake_write_loop(); |
| 4359 | } |
| 4360 | |
| 4361 | _ => {} |
| 4362 | } |
| 4363 | } |
| 4364 | |
| 4365 | fn on_retransmission_failure(&mut self, id: Timer) { |
| 4366 | match id { |
| 4367 | Timer::T1Init => { |
| 4368 | error!("[{}] retransmission failure: T1-init", self.side); |
| 4369 | self.error = Some(AssociationError::HandshakeFailed( |
| 4370 | Error::ErrHandshakeInitAck, |
| 4371 | )); |
| 4372 | } |
| 4373 | |
| 4374 | Timer::T1Cookie => { |
| 4375 | error!("[{}] retransmission failure: T1-cookie", self.side); |
| 4376 | self.error = Some(AssociationError::HandshakeFailed( |
| 4377 | Error::ErrHandshakeCookieEcho, |
| 4378 | )); |
| 4379 | } |
| 4380 | |
| 4381 | Timer::T2Shutdown => { |
| 4382 | error!("[{}] retransmission failure: T2-shutdown", self.side); |
| 4383 | } |
| 4384 | |
| 4385 | Timer::T3RTX => { |
| 4386 | // T3-rtx timer will not fail by design |
| 4387 | // Justifications: |
| 4388 | // * ICE would fail if the connectivity is lost |
| 4389 | // * WebRTC spec is not clear how this incident should be reported to ULP |
| 4390 | error!("[{}] retransmission failure: T3-rtx (DATA)", self.side); |
| 4391 | } |
| 4392 | |
| 4393 | Timer::Reconfig => { |
| 4394 | if let Some(rsn) = self.active_reconfig { |
| 4395 | error!( |
| 4396 | "[{}] retransmission failure: Reconfig rsn={}", |
| 4397 | self.side, rsn |
| 4398 | ); |
| 4399 | self.finish_reconfig(rsn, Err(StreamResetError::Failed)); |
| 4400 | } else { |
| 4401 | self.timers.stop(Timer::Reconfig); |
| 4402 | } |
| 4403 | } |
| 4404 | |
| 4405 | _ => {} |
| 4406 | } |
| 4407 | } |
| 4408 | |
| 4409 | /// Whether no timers are running |
| 4410 | #[cfg(test)] |
| 4411 | pub(crate) fn is_idle(&self) -> bool { |
| 4412 | Timer::VALUES |
| 4413 | .iter() |
| 4414 | //.filter(|&&t| t != Timer::KeepAlive && t != Timer::PushNewCid) |
| 4415 | .filter_map(|&t| Some((t, self.timers.get(t)?))) |
| 4416 | .min_by_key(|&(_, time)| time) |
| 4417 | //.map_or(true, |(timer, _)| timer == Timer::Idle) |
| 4418 | .is_none() |
| 4419 | } |
| 4420 | } |