Skip to content
File

Blob: firmware/vendor/sctp-proto/src/association/mod.rs

rust4421 lines
1use crate::association::state::{AckMode, AckState, AssociationState};
2use crate::association::stats::AssociationStats;
3use crate::chunk::Chunk;
4use crate::chunk::ErrorCauseUnrecognizedChunkType;
5use crate::chunk::USER_INITIATED_ABORT;
6use crate::chunk::chunk_abort::ChunkAbort;
7use crate::chunk::chunk_cookie_ack::ChunkCookieAck;
8use crate::chunk::chunk_cookie_echo::ChunkCookieEcho;
9use crate::chunk::chunk_error::ChunkError;
10use crate::chunk::chunk_forward_tsn::{ChunkForwardTsn, ChunkForwardTsnStream};
11use crate::chunk::chunk_heartbeat::ChunkHeartbeat;
12use crate::chunk::chunk_heartbeat_ack::ChunkHeartbeatAck;
13use crate::chunk::chunk_i_forward_tsn::ChunkIForwardTsn;
14use crate::chunk::chunk_init::{ChunkInit, ChunkInitAck};
15use crate::chunk::chunk_payload_data::{ChunkPayloadData, PayloadProtocolIdentifier};
16use crate::chunk::chunk_reconfig::ChunkReconfig;
17use crate::chunk::chunk_selective_ack::ChunkSelectiveAck;
18use crate::chunk::chunk_shutdown::ChunkShutdown;
19use crate::chunk::chunk_shutdown_ack::ChunkShutdownAck;
20use crate::chunk::chunk_shutdown_complete::ChunkShutdownComplete;
21use crate::chunk::chunk_type::CT_FORWARD_TSN;
22use crate::config::COMMON_HEADER_SIZE;
23use crate::config::DATA_CHUNK_HEADER_SIZE;
24use crate::config::DEFAULT_SCTP_PORT;
25use crate::config::{ServerConfig, TransportConfig};
26use crate::error::{Error, Result};
27use crate::packet::{CommonHeader, Packet};
28use crate::param::Param;
29use crate::param::param_heartbeat_info::ParamHeartbeatInfo;
30use crate::param::param_outgoing_reset_request::ParamOutgoingResetRequest;
31use crate::param::param_reconfig_response::{ParamReconfigResponse, ReconfigResult};
32use crate::param::param_state_cookie::ParamStateCookie;
33use crate::param::param_supported_extensions::ParamSupportedExtensions;
34use crate::queue::payload_queue::PayloadQueue;
35use crate::queue::pending_queue::PendingQueue;
36use crate::queue::reassembly_queue::ReassemblyQueue;
37use crate::shared::{AssociationEventInner, AssociationId, EndpointEvent, EndpointEventInner};
38use crate::util::{sna16lt, sna32gt, sna32gte, sna32lt, sna32lte};
39use crate::{AssociationEvent, Payload, Side, Transmit};
40use stream::{ReliabilityType, Stream, StreamEvent, StreamId, StreamResetError, StreamState};
41use timer::{ACK_INTERVAL, RtoManager, Timer, TimerTable};
42 
43use crate::association::stream::RecvSendState;
44use alloc::boxed::Box;
45use alloc::collections::VecDeque;
46use alloc::string::String;
47use alloc::sync::Arc;
48use alloc::vec;
49use alloc::vec::Vec;
50use bytes::Bytes;
51use core::net::{IpAddr, SocketAddr};
52use core::num::NonZeroU32;
53use core::str::FromStr;
54use core::time::Duration;
55use log::{debug, error, trace, warn};
56use rand::random;
57use rustc_hash::{FxHashMap, FxHashSet};
58use std::collections::HashMap;
59use std::time::Instant;
60use thiserror::Error;
61 
62pub(crate) mod state;
63pub(crate) mod stats;
64pub(crate) mod stream;
65mod timer;
66 
67#[cfg(test)]
68mod association_test;
69#[cfg(test)]
70mod receive_limits_test;
71 
72/// Reasons why an association might be lost
73#[non_exhaustive]
74#[derive(Debug, Error, Clone, PartialEq)]
75pub 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)]
105pub 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)]
132struct PendingResetCompletions(FxHashMap<StreamId, usize>);
133 
134impl 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)]
163enum 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)]
174struct 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)]
199pub 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 
321impl 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 
425impl 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}