Skip to content
File

Blob: firmware/vendor/sctp-proto/src/endpoint/endpoint_test.rs

rust3868 lines
1use alloc::borrow::ToOwned;
2use alloc::vec;
3use alloc::vec::Vec;
4use std::println;
5 
6use super::*;
7use crate::association::Event;
8use crate::config::generate_snap_token;
9use crate::error::{Error, Result};
10 
11use crate::association::state::{AckMode, AssociationState};
12use crate::association::stream::{ReliabilityType, Stream, StreamEvent, StreamResetError};
13use crate::chunk::chunk_abort::ChunkAbort;
14use crate::chunk::chunk_cookie_echo::ChunkCookieEcho;
15use crate::chunk::chunk_error::ChunkError;
16use crate::chunk::chunk_forward_tsn::ChunkForwardTsn;
17use crate::chunk::chunk_heartbeat::ChunkHeartbeat;
18use crate::chunk::chunk_init::ChunkInit;
19use crate::chunk::chunk_payload_data::{ChunkPayloadData, PayloadProtocolIdentifier};
20use crate::chunk::chunk_reconfig::ChunkReconfig;
21use crate::chunk::chunk_selective_ack::{ChunkSelectiveAck, GapAckBlock};
22use crate::chunk::chunk_shutdown::ChunkShutdown;
23use crate::chunk::chunk_shutdown_ack::ChunkShutdownAck;
24use crate::chunk::chunk_shutdown_complete::ChunkShutdownComplete;
25use crate::chunk::{ErrorCauseProtocolViolation, PROTOCOL_VIOLATION};
26use crate::packet::{CommonHeader, Packet};
27use crate::param::param_outgoing_reset_request::ParamOutgoingResetRequest;
28use crate::param::param_reconfig_response::ParamReconfigResponse;
29use assert_matches::assert_matches;
30use core::net::Ipv6Addr;
31use core::ops::RangeFrom;
32use core::str::FromStr;
33use core::time::Duration;
34use core::{cmp, mem};
35use log::{info, trace};
36use std::net::UdpSocket;
37use std::sync::{LazyLock, Mutex};
38use std::time::Instant;
39 
40pub static SERVER_PORTS: LazyLock<Mutex<RangeFrom<u16>>> = LazyLock::new(|| Mutex::new(4433..));
41pub static CLIENT_PORTS: LazyLock<Mutex<RangeFrom<u16>>> = LazyLock::new(|| Mutex::new(44433..));
42 
43fn min_opt<T: Ord>(x: Option<T>, y: Option<T>) -> Option<T> {
44 match (x, y) {
45 (Some(x), Some(y)) => Some(cmp::min(x, y)),
46 (Some(x), _) => Some(x),
47 (_, Some(y)) => Some(y),
48 _ => None,
49 }
50}
51 
52/// The maximum of datagrams TestEndpoint will produce via `poll_transmit`
53const MAX_DATAGRAMS: usize = 10;
54 
55fn split_transmit(transmit: Transmit) -> Vec<Transmit> {
56 let mut transmits = Vec::new();
57 if let Payload::RawEncode(contents) = transmit.payload {
58 for content in contents {
59 transmits.push(Transmit {
60 now: transmit.now,
61 remote: transmit.remote,
62 payload: Payload::RawEncode(vec![content]),
63 ecn: transmit.ecn,
64 local_ip: transmit.local_ip,
65 });
66 }
67 }
68 
69 transmits
70}
71 
72pub fn client_config() -> ClientConfig {
73 ClientConfig::new()
74}
75 
76pub fn server_config() -> ServerConfig {
77 ServerConfig::new()
78}
79 
80struct TestEndpoint {
81 endpoint: Endpoint,
82 addr: SocketAddr,
83 socket: Option<UdpSocket>,
84 timeout: Option<Instant>,
85 outbound: VecDeque<Transmit>,
86 delayed: VecDeque<Transmit>,
87 inbound: VecDeque<(Instant, Option<EcnCodepoint>, Bytes)>,
88 accepted: Option<AssociationHandle>,
89 associations: HashMap<AssociationHandle, Association>,
90 conn_events: HashMap<AssociationHandle, VecDeque<AssociationEvent>>,
91}
92 
93impl TestEndpoint {
94 fn new(endpoint: Endpoint, addr: SocketAddr) -> Self {
95 let socket = UdpSocket::bind(addr).expect("failed to bind UDP socket");
96 socket
97 .set_read_timeout(Some(Duration::new(0, 10_000_000)))
98 .unwrap();
99 
100 Self {
101 endpoint,
102 addr,
103 socket: Some(socket),
104 timeout: None,
105 outbound: VecDeque::new(),
106 delayed: VecDeque::new(),
107 inbound: VecDeque::new(),
108 accepted: None,
109 associations: HashMap::default(),
110 conn_events: HashMap::default(),
111 }
112 }
113 
114 pub fn drive(&mut self, now: Instant, remote: SocketAddr) {
115 if let Some(ref socket) = self.socket {
116 loop {
117 let mut buf = [0; 8192];
118 if socket.recv_from(&mut buf).is_err() {
119 break;
120 }
121 }
122 }
123 
124 while self.inbound.front().is_some_and(|x| x.0 <= now) {
125 let (recv_time, ecn, packet) = self.inbound.pop_front().unwrap();
126 if let Some((ch, event)) = self.endpoint.handle(recv_time, remote, None, ecn, packet) {
127 match event {
128 DatagramEvent::NewAssociation(conn) => {
129 self.associations.insert(ch, conn);
130 self.accepted = Some(ch);
131 }
132 DatagramEvent::AssociationEvent(event) => {
133 self.conn_events.entry(ch).or_default().push_back(event);
134 }
135 }
136 }
137 }
138 
139 while let Some(x) = self.poll_transmit() {
140 self.outbound.extend(split_transmit(x));
141 }
142 
143 let mut endpoint_events: Vec<(AssociationHandle, EndpointEvent)> = vec![];
144 for (ch, conn) in self.associations.iter_mut() {
145 if self.timeout.is_some_and(|x| x <= now) {
146 self.timeout = None;
147 conn.handle_timeout(now);
148 }
149 
150 for (_, mut events) in self.conn_events.drain() {
151 for event in events.drain(..) {
152 conn.handle_event(event);
153 }
154 }
155 
156 while let Some(event) = conn.poll_endpoint_event() {
157 endpoint_events.push((*ch, event));
158 }
159 
160 while let Some(x) = conn.poll_transmit(now) {
161 self.outbound.extend(split_transmit(x));
162 }
163 self.timeout = conn.poll_timeout();
164 }
165 
166 for (ch, event) in endpoint_events {
167 if let Some(event) = self.handle_event(ch, event) {
168 if let Some(conn) = self.associations.get_mut(&ch) {
169 conn.handle_event(event);
170 }
171 }
172 }
173 }
174 
175 pub fn next_wakeup(&self) -> Option<Instant> {
176 let next_inbound = self.inbound.front().map(|x| x.0);
177 min_opt(self.timeout, next_inbound)
178 }
179 
180 fn is_idle(&self) -> bool {
181 self.associations.values().all(|x| x.is_idle())
182 }
183 
184 pub fn delay_outbound(&mut self) {
185 assert!(self.delayed.is_empty());
186 mem::swap(&mut self.delayed, &mut self.outbound);
187 }
188 
189 pub fn finish_delay(&mut self) {
190 self.outbound.extend(self.delayed.drain(..));
191 }
192 
193 pub fn assert_accept(&mut self) -> AssociationHandle {
194 self.accepted.take().expect("server didn't connect")
195 }
196}
197 
198impl ::core::ops::Deref for TestEndpoint {
199 type Target = Endpoint;
200 fn deref(&self) -> &Endpoint {
201 &self.endpoint
202 }
203}
204 
205impl ::core::ops::DerefMut for TestEndpoint {
206 fn deref_mut(&mut self) -> &mut Endpoint {
207 &mut self.endpoint
208 }
209}
210 
211struct Pair {
212 server: TestEndpoint,
213 client: TestEndpoint,
214 time: Instant,
215 latency: Duration, // One-way
216}
217 
218impl Pair {
219 pub fn new(endpoint_config: Arc<EndpointConfig>, server_config: ServerConfig) -> Self {
220 let server = Endpoint::new(endpoint_config.clone(), Some(Arc::new(server_config)));
221 let client = Endpoint::new(endpoint_config, None);
222 
223 Pair::new_from_endpoint(client, server)
224 }
225 
226 pub fn new_from_endpoint(client: Endpoint, server: Endpoint) -> Self {
227 let server_addr = SocketAddr::new(
228 Ipv6Addr::LOCALHOST.into(),
229 SERVER_PORTS.lock().unwrap().next().unwrap(),
230 );
231 let client_addr = SocketAddr::new(
232 Ipv6Addr::LOCALHOST.into(),
233 CLIENT_PORTS.lock().unwrap().next().unwrap(),
234 );
235 Self {
236 server: TestEndpoint::new(server, server_addr),
237 client: TestEndpoint::new(client, client_addr),
238 time: Instant::now(),
239 latency: Duration::new(0, 0),
240 }
241 }
242 
243 /// Returns whether the association is not idle
244 pub fn step(&mut self) -> bool {
245 self.drive_client();
246 self.drive_server();
247 if self.client.is_idle() && self.server.is_idle() {
248 return false;
249 }
250 
251 let client_t = self.client.next_wakeup();
252 let server_t = self.server.next_wakeup();
253 match min_opt(client_t, server_t) {
254 Some(t) if Some(t) == client_t => {
255 if t != self.time {
256 self.time = self.time.max(t);
257 trace!("advancing to {:?} for client", self.time);
258 }
259 true
260 }
261 Some(t) if Some(t) == server_t => {
262 if t != self.time {
263 self.time = self.time.max(t);
264 trace!("advancing to {:?} for server", self.time);
265 }
266 true
267 }
268 Some(_) => unreachable!(),
269 None => false,
270 }
271 }
272 
273 /// Advance time until both associations are idle
274 pub fn drive(&mut self) {
275 while self.step() {}
276 }
277 
278 pub fn drive_client(&mut self) {
279 self.client.drive(self.time, self.server.addr);
280 for x in self.client.outbound.drain(..) {
281 if let Payload::RawEncode(contents) = x.payload {
282 for content in contents {
283 if let Some(ref socket) = self.client.socket {
284 socket.send_to(&content, x.remote).unwrap();
285 }
286 if self.server.addr == x.remote {
287 self.server
288 .inbound
289 .push_back((self.time + self.latency, x.ecn, content));
290 }
291 }
292 }
293 }
294 }
295 
296 pub fn drive_server(&mut self) {
297 self.server.drive(self.time, self.client.addr);
298 for x in self.server.outbound.drain(..) {
299 if let Payload::RawEncode(contents) = x.payload {
300 for content in contents {
301 if let Some(ref socket) = self.server.socket {
302 socket.send_to(&content, x.remote).unwrap();
303 }
304 if self.client.addr == x.remote {
305 self.client
306 .inbound
307 .push_back((self.time + self.latency, x.ecn, content));
308 }
309 }
310 }
311 }
312 }
313 
314 pub fn connect(&mut self) -> (AssociationHandle, AssociationHandle) {
315 self.connect_with(client_config())
316 }
317 
318 pub fn connect_with(&mut self, config: ClientConfig) -> (AssociationHandle, AssociationHandle) {
319 info!("connecting");
320 let client_ch = self.begin_connect(config);
321 self.drive();
322 let server_ch = self.server.assert_accept();
323 self.finish_connect(client_ch, server_ch);
324 (client_ch, server_ch)
325 }
326 
327 /// Just start connecting the client
328 pub fn begin_connect(&mut self, config: ClientConfig) -> AssociationHandle {
329 let (client_ch, client_conn) = self.client.connect(config, self.server.addr).unwrap();
330 self.client.associations.insert(client_ch, client_conn);
331 client_ch
332 }
333 
334 fn finish_connect(&mut self, client_ch: AssociationHandle, server_ch: AssociationHandle) {
335 assert_matches!(
336 self.client_conn_mut(client_ch).poll(),
337 Some(Event::Connected)
338 );
339 
340 assert_matches!(
341 self.server_conn_mut(server_ch).poll(),
342 Some(Event::Connected)
343 );
344 }
345 
346 pub fn client_conn_mut(&mut self, ch: AssociationHandle) -> &mut Association {
347 self.client.associations.get_mut(&ch).unwrap()
348 }
349 
350 pub fn client_stream(&mut self, ch: AssociationHandle, si: u16) -> Result<Stream<'_>> {
351 self.client_conn_mut(ch).stream(si)
352 }
353 
354 pub fn server_conn_mut(&mut self, ch: AssociationHandle) -> &mut Association {
355 self.server.associations.get_mut(&ch).unwrap()
356 }
357 
358 pub fn server_stream(&mut self, ch: AssociationHandle, si: u16) -> Result<Stream<'_>> {
359 self.server_conn_mut(ch).stream(si)
360 }
361}
362 
363impl Default for Pair {
364 fn default() -> Self {
365 Pair::new(Default::default(), server_config())
366 }
367}
368 
369fn create_association_pair(
370 ack_mode: AckMode,
371 recv_buf_size: u32,
372) -> Result<(Pair, AssociationHandle, AssociationHandle)> {
373 let mut pair = Pair::new(
374 Arc::new(EndpointConfig::default()),
375 ServerConfig {
376 transport: Arc::new(if recv_buf_size > 0 {
377 TransportConfig::default().with_max_receive_buffer_size(recv_buf_size)
378 } else {
379 TransportConfig::default()
380 }),
381 ..Default::default()
382 },
383 );
384 let (client_ch, server_ch) = pair.connect_with(ClientConfig {
385 transport: Arc::new(if recv_buf_size > 0 {
386 TransportConfig::default().with_max_receive_buffer_size(recv_buf_size)
387 } else {
388 TransportConfig::default()
389 }),
390 ..Default::default()
391 });
392 pair.client_conn_mut(client_ch).ack_mode = ack_mode;
393 pair.server_conn_mut(server_ch).ack_mode = ack_mode;
394 Ok((pair, client_ch, server_ch))
395}
396 
397fn establish_session_pair(
398 pair: &mut Pair,
399 client_ch: AssociationHandle,
400 server_ch: AssociationHandle,
401 si: u16,
402) -> Result<()> {
403 let hello_msg = Bytes::from_static(b"Hello");
404 let _ = pair
405 .client_conn_mut(client_ch)
406 .open_stream(si, PayloadProtocolIdentifier::Binary)?;
407 let _ = pair
408 .client_stream(client_ch, si)?
409 .write_sctp(&hello_msg, PayloadProtocolIdentifier::Dcep)?;
410 pair.drive();
411 
412 {
413 let s1 = pair.server_conn_mut(server_ch).accept_stream().unwrap();
414 if si != s1.stream_identifier {
415 return Err(Error::Other("si should match".to_owned()));
416 }
417 }
418 pair.drive();
419 
420 let mut buf = vec![0u8; 1024];
421 let chunks = pair.server_stream(server_ch, si)?.read_sctp()?.unwrap();
422 let n = chunks.read(&mut buf)?;
423 
424 if n != hello_msg.len() {
425 return Err(Error::Other("received data must by 3 bytes".to_owned()));
426 }
427 
428 if chunks.ppi != PayloadProtocolIdentifier::Dcep {
429 return Err(Error::Other("unexpected ppi".to_owned()));
430 }
431 
432 if buf[..n] != hello_msg {
433 return Err(Error::Other("received data mismatch".to_owned()));
434 }
435 pair.drive();
436 
437 Ok(())
438}
439 
440fn close_association_pair(
441 _pair: &mut Pair,
442 _client_ch: AssociationHandle,
443 _server_ch: AssociationHandle,
444 _si: u16,
445) {
446 /*TODO:
447 // Close client
448 tokio::spawn(async move {
449 client.close().await?;
450 let _ = handshake0ch_tx.send(()).await;
451 let _ = closed_rx0.recv().await;
452
453 Result::<()>::Ok(())
454 });
455
456 // Close server
457 tokio::spawn(async move {
458 server.close().await?;
459 let _ = handshake1ch_tx.send(()).await;
460 let _ = closed_rx1.recv().await;
461
462 Result::<()>::Ok(())
463 });
464 */
465}
466 
467#[test]
468fn test_assoc_reliable_simple() -> Result<()> {
469 //let _guard = subscribe();
470 
471 let si: u16 = 1;
472 let msg: Bytes = Bytes::from_static(b"ABC");
473 
474 let (mut pair, client_ch, server_ch) = create_association_pair(AckMode::NoDelay, 0)?;
475 
476 establish_session_pair(&mut pair, client_ch, server_ch, si)?;
477 
478 {
479 let a = pair.client_conn_mut(client_ch);
480 assert_eq!(0, a.buffered_amount(), "incorrect bufferedAmount");
481 }
482 
483 let n = pair
484 .client_stream(client_ch, si)?
485 .write_sctp(&msg, PayloadProtocolIdentifier::Binary)?;
486 assert_eq!(msg.len(), n, "unexpected length of received data");
487 {
488 let a = pair.client_conn_mut(client_ch);
489 assert_eq!(msg.len(), a.buffered_amount(), "incorrect bufferedAmount");
490 }
491 
492 pair.drive();
493 
494 let chunks = pair.server_stream(server_ch, si)?.read_sctp()?.unwrap();
495 let (n, ppi) = (chunks.len(), chunks.ppi);
496 assert_eq!(n, msg.len(), "unexpected length of received data");
497 assert_eq!(ppi, PayloadProtocolIdentifier::Binary, "unexpected ppi");
498 
499 {
500 let q = &pair
501 .client_conn_mut(client_ch)
502 .streams
503 .get(&si)
504 .unwrap()
505 .reassembly_queue;
506 assert!(!q.is_readable(), "should no longer be readable");
507 }
508 
509 {
510 let a = pair.client_conn_mut(client_ch);
511 assert_eq!(0, a.buffered_amount(), "incorrect bufferedAmount");
512 }
513 
514 close_association_pair(&mut pair, client_ch, server_ch, si);
515 
516 Ok(())
517}
518 
519#[test]
520fn test_assoc_reliable_ordered_reordered() -> Result<()> {
521 // let _guard = subscribe();
522 
523 let si: u16 = 2;
524 let mut sbuf = vec![0u8; 1000];
525 for (i, b) in sbuf.iter_mut().enumerate() {
526 *b = (i & 0xff) as u8;
527 }
528 
529 let (mut pair, client_ch, server_ch) = create_association_pair(AckMode::NoDelay, 0)?;
530 
531 establish_session_pair(&mut pair, client_ch, server_ch, si)?;
532 
533 {
534 let a = pair.client_conn_mut(client_ch);
535 assert_eq!(0, a.buffered_amount(), "incorrect bufferedAmount");
536 }
537 
538 sbuf[0..4].copy_from_slice(&0u32.to_be_bytes());
539 let n = pair.client_stream(client_ch, si)?.write_sctp(
540 &Bytes::from(sbuf.clone()),
541 PayloadProtocolIdentifier::Binary,
542 )?;
543 assert_eq!(sbuf.len(), n, "unexpected length of received data");
544 pair.client.drive(pair.time, pair.server.addr);
545 pair.client.delay_outbound(); // Delay it
546 
547 sbuf[0..4].copy_from_slice(&1u32.to_be_bytes());
548 let n = pair.client_stream(client_ch, si)?.write_sctp(
549 &Bytes::from(sbuf.clone()),
550 PayloadProtocolIdentifier::Binary,
551 )?;
552 assert_eq!(sbuf.len(), n, "unexpected length of received data");
553 pair.client.drive(pair.time, pair.server.addr);
554 pair.client.finish_delay(); // Reorder it
555 
556 pair.drive();
557 
558 let mut buf = vec![0u8; 2000];
559 
560 let chunks = pair.server_stream(server_ch, si)?.read_sctp()?.unwrap();
561 let (n, ppi) = (chunks.len(), chunks.ppi);
562 chunks.read(&mut buf)?;
563 assert_eq!(n, sbuf.len(), "unexpected length of received data");
564 assert_eq!(ppi, PayloadProtocolIdentifier::Binary, "unexpected ppi");
565 assert_eq!(
566 u32::from_be_bytes([buf[0], buf[1], buf[2], buf[3]]),
567 0,
568 "unexpected received data"
569 );
570 
571 let chunks = pair.server_stream(server_ch, si)?.read_sctp()?.unwrap();
572 let (n, ppi) = (chunks.len(), chunks.ppi);
573 chunks.read(&mut buf)?;
574 assert_eq!(n, sbuf.len(), "unexpected length of received data");
575 assert_eq!(ppi, PayloadProtocolIdentifier::Binary, "unexpected ppi");
576 assert_eq!(
577 u32::from_be_bytes([buf[0], buf[1], buf[2], buf[3]]),
578 1,
579 "unexpected received data"
580 );
581 
582 pair.drive();
583 
584 {
585 let q = &pair
586 .client_conn_mut(client_ch)
587 .streams
588 .get(&si)
589 .unwrap()
590 .reassembly_queue;
591 assert!(!q.is_readable(), "should no longer be readable");
592 }
593 
594 close_association_pair(&mut pair, client_ch, server_ch, si);
595 
596 Ok(())
597}
598 
599#[test]
600fn test_assoc_reliable_ordered_fragmented_then_defragmented() -> Result<()> {
601 //let _guard = subscribe();
602 
603 let si: u16 = 3;
604 let mut sbuf = vec![0u8; 1000];
605 for (i, b) in sbuf.iter_mut().enumerate() {
606 *b = (i & 0xff) as u8;
607 }
608 let mut sbufl = vec![0u8; 2000];
609 for (i, b) in sbufl.iter_mut().enumerate() {
610 *b = (i & 0xff) as u8;
611 }
612 
613 let (mut pair, client_ch, server_ch) = create_association_pair(AckMode::NoDelay, 0)?;
614 
615 establish_session_pair(&mut pair, client_ch, server_ch, si)?;
616 
617 pair.client_stream(client_ch, si)?.set_reliability_params(
618 false,
619 ReliabilityType::Reliable,
620 0,
621 )?;
622 pair.server_stream(server_ch, si)?.set_reliability_params(
623 false,
624 ReliabilityType::Reliable,
625 0,
626 )?;
627 
628 let n = pair.client_stream(client_ch, si)?.write_sctp(
629 &Bytes::from(sbufl.clone()),
630 PayloadProtocolIdentifier::Binary,
631 )?;
632 assert_eq!(sbufl.len(), n, "unexpected length of received data");
633 
634 pair.drive();
635 
636 let mut rbuf = vec![0u8; 2000];
637 let chunks = pair.server_stream(server_ch, si)?.read_sctp()?.unwrap();
638 let (n, ppi) = (chunks.len(), chunks.ppi);
639 chunks.read(&mut rbuf)?;
640 assert_eq!(n, sbufl.len(), "unexpected length of received data");
641 assert_eq!(&rbuf[..n], &sbufl, "unexpected received data");
642 assert_eq!(ppi, PayloadProtocolIdentifier::Binary, "unexpected ppi");
643 
644 pair.drive();
645 
646 {
647 let q = &pair
648 .client_conn_mut(client_ch)
649 .streams
650 .get(&si)
651 .unwrap()
652 .reassembly_queue;
653 assert!(!q.is_readable(), "should no longer be readable");
654 }
655 
656 close_association_pair(&mut pair, client_ch, server_ch, si);
657 
658 Ok(())
659}
660 
661#[test]
662fn test_assoc_reliable_unordered_fragmented_then_defragmented() -> Result<()> {
663 //let _guard = subscribe();
664 
665 let si: u16 = 4;
666 let mut sbuf = vec![0u8; 1000];
667 for (i, b) in sbuf.iter_mut().enumerate() {
668 *b = (i & 0xff) as u8;
669 }
670 let sbufl = vec![0u8; 2000];
671 for (i, b) in sbuf.iter_mut().enumerate() {
672 *b = (i & 0xff) as u8;
673 }
674 
675 let (mut pair, client_ch, server_ch) = create_association_pair(AckMode::NoDelay, 0)?;
676 
677 establish_session_pair(&mut pair, client_ch, server_ch, si)?;
678 
679 pair.client_stream(client_ch, si)?.set_reliability_params(
680 true,
681 ReliabilityType::Reliable,
682 0,
683 )?;
684 pair.server_stream(server_ch, si)?.set_reliability_params(
685 true,
686 ReliabilityType::Reliable,
687 0,
688 )?;
689 
690 let n = pair.client_stream(client_ch, si)?.write_sctp(
691 &Bytes::from(sbufl.clone()),
692 PayloadProtocolIdentifier::Binary,
693 )?;
694 assert_eq!(sbufl.len(), n, "unexpected length of received data");
695 
696 pair.drive();
697 
698 let mut rbuf = vec![0u8; 2000];
699 let chunks = pair.server_stream(server_ch, si)?.read_sctp()?.unwrap();
700 let (n, ppi) = (chunks.len(), chunks.ppi);
701 chunks.read(&mut rbuf)?;
702 assert_eq!(n, sbufl.len(), "unexpected length of received data");
703 assert_eq!(&rbuf[..n], &sbufl, "unexpected received data");
704 assert_eq!(ppi, PayloadProtocolIdentifier::Binary, "unexpected ppi");
705 
706 pair.drive();
707 
708 {
709 let q = &pair
710 .client_conn_mut(client_ch)
711 .streams
712 .get(&si)
713 .unwrap()
714 .reassembly_queue;
715 assert!(!q.is_readable(), "should no longer be readable");
716 }
717 
718 close_association_pair(&mut pair, client_ch, server_ch, si);
719 
720 Ok(())
721}
722 
723#[test]
724fn test_assoc_reliable_unordered_ordered() -> Result<()> {
725 //let _guard = subscribe();
726 
727 let si: u16 = 5;
728 let mut sbuf = vec![0u8; 1000];
729 for (i, b) in sbuf.iter_mut().enumerate() {
730 *b = (i & 0xff) as u8;
731 }
732 
733 let (mut pair, client_ch, server_ch) = create_association_pair(AckMode::NoDelay, 0)?;
734 
735 establish_session_pair(&mut pair, client_ch, server_ch, si)?;
736 
737 pair.client_stream(client_ch, si)?.set_reliability_params(
738 true,
739 ReliabilityType::Reliable,
740 0,
741 )?;
742 pair.server_stream(server_ch, si)?.set_reliability_params(
743 true,
744 ReliabilityType::Reliable,
745 0,
746 )?;
747 
748 sbuf[0..4].copy_from_slice(&0u32.to_be_bytes());
749 let n = pair.client_stream(client_ch, si)?.write_sctp(
750 &Bytes::from(sbuf.clone()),
751 PayloadProtocolIdentifier::Binary,
752 )?;
753 assert_eq!(sbuf.len(), n, "unexpected length of received data");
754 pair.client.drive(pair.time, pair.server.addr);
755 pair.client.delay_outbound(); // Delay it
756 
757 sbuf[0..4].copy_from_slice(&1u32.to_be_bytes());
758 let n = pair.client_stream(client_ch, si)?.write_sctp(
759 &Bytes::from(sbuf.clone()),
760 PayloadProtocolIdentifier::Binary,
761 )?;
762 assert_eq!(sbuf.len(), n, "unexpected length of received data");
763 pair.client.drive(pair.time, pair.server.addr);
764 pair.client.finish_delay(); // Reorder it
765 
766 pair.drive();
767 
768 let mut buf = vec![0u8; 2000];
769 
770 let chunks = pair.server_stream(server_ch, si)?.read_sctp()?.unwrap();
771 let (n, ppi) = (chunks.len(), chunks.ppi);
772 chunks.read(&mut buf)?;
773 assert_eq!(n, sbuf.len(), "unexpected length of received data");
774 assert_eq!(ppi, PayloadProtocolIdentifier::Binary, "unexpected ppi");
775 assert_eq!(
776 u32::from_be_bytes([buf[0], buf[1], buf[2], buf[3]]),
777 1,
778 "unexpected received data"
779 );
780 
781 let chunks = pair.server_stream(server_ch, si)?.read_sctp()?.unwrap();
782 let (n, ppi) = (chunks.len(), chunks.ppi);
783 chunks.read(&mut buf)?;
784 assert_eq!(n, sbuf.len(), "unexpected length of received data");
785 assert_eq!(ppi, PayloadProtocolIdentifier::Binary, "unexpected ppi");
786 assert_eq!(
787 u32::from_be_bytes([buf[0], buf[1], buf[2], buf[3]]),
788 0,
789 "unexpected received data"
790 );
791 
792 pair.drive();
793 
794 {
795 let q = &pair
796 .client_conn_mut(client_ch)
797 .streams
798 .get(&si)
799 .unwrap()
800 .reassembly_queue;
801 assert!(!q.is_readable(), "should no longer be readable");
802 }
803 
804 close_association_pair(&mut pair, client_ch, server_ch, si);
805 
806 Ok(())
807}
808 
809#[test]
810fn test_assoc_reliable_retransmission() -> Result<()> {
811 //let _guard = subscribe();
812 
813 let si: u16 = 6;
814 let msg1: Bytes = Bytes::from_static(b"ABC");
815 let msg2: Bytes = Bytes::from_static(b"DEFG");
816 
817 let (mut pair, client_ch, server_ch) = create_association_pair(AckMode::NoDelay, 0)?;
818 
819 {
820 let a = pair.client_conn_mut(client_ch);
821 a.rto_mgr.set_rto(100, true);
822 }
823 
824 establish_session_pair(&mut pair, client_ch, server_ch, si)?;
825 
826 let n = pair
827 .client_stream(client_ch, si)?
828 .write_sctp(&msg1, PayloadProtocolIdentifier::Binary)?;
829 assert_eq!(msg1.len(), n, "unexpected length of received data");
830 pair.drive_client(); // send data to server
831 pair.server.inbound.clear(); // Lose it
832 debug!("dropping packet");
833 
834 let n = pair
835 .client_stream(client_ch, si)?
836 .write_sctp(&msg2, PayloadProtocolIdentifier::Binary)?;
837 assert_eq!(msg2.len(), n, "unexpected length of received data");
838 
839 pair.drive();
840 
841 let mut buf = vec![0u8; 32];
842 
843 let chunks = pair.server_stream(server_ch, si)?.read_sctp()?.unwrap();
844 let (n, ppi) = (chunks.len(), chunks.ppi);
845 chunks.read(&mut buf)?;
846 assert_eq!(n, msg1.len(), "unexpected length of received data");
847 assert_eq!(&buf[..n], &msg1, "unexpected length of received data");
848 assert_eq!(ppi, PayloadProtocolIdentifier::Binary, "unexpected ppi");
849 
850 let chunks = pair.server_stream(server_ch, si)?.read_sctp()?.unwrap();
851 let (n, ppi) = (chunks.len(), chunks.ppi);
852 chunks.read(&mut buf)?;
853 assert_eq!(n, msg2.len(), "unexpected length of received data");
854 assert_eq!(&buf[..n], &msg2, "unexpected length of received data");
855 assert_eq!(ppi, PayloadProtocolIdentifier::Binary, "unexpected ppi");
856 
857 pair.drive();
858 
859 {
860 let q = &pair
861 .client_conn_mut(client_ch)
862 .streams
863 .get(&si)
864 .unwrap()
865 .reassembly_queue;
866 assert!(!q.is_readable(), "should no longer be readable");
867 }
868 
869 close_association_pair(&mut pair, client_ch, server_ch, si);
870 
871 Ok(())
872}
873 
874#[test]
875fn test_assoc_reliable_short_buffer() -> Result<()> {
876 //let _guard = subscribe();
877 
878 let si: u16 = 1;
879 let msg: Bytes = Bytes::from_static(b"Hello");
880 
881 let (mut pair, client_ch, server_ch) = create_association_pair(AckMode::NoDelay, 0)?;
882 
883 establish_session_pair(&mut pair, client_ch, server_ch, si)?;
884 
885 {
886 let a = pair.client_conn_mut(client_ch);
887 assert_eq!(0, a.buffered_amount(), "incorrect bufferedAmount");
888 }
889 
890 let n = pair
891 .client_stream(client_ch, si)?
892 .write_sctp(&msg, PayloadProtocolIdentifier::Binary)?;
893 assert_eq!(msg.len(), n, "unexpected length of received data");
894 {
895 let a = pair.client_conn_mut(client_ch);
896 assert_eq!(msg.len(), a.buffered_amount(), "incorrect bufferedAmount");
897 }
898 
899 pair.drive();
900 
901 let mut buf = vec![0u8; 3];
902 let chunks = pair.server_stream(server_ch, si)?.read_sctp()?.unwrap();
903 let result = chunks.read(&mut buf);
904 assert!(result.is_err(), "expected error to be io.ErrShortBuffer");
905 if let Err(err) = result {
906 assert_eq!(
907 Error::ErrShortBuffer,
908 err,
909 "expected error to be io.ErrShortBuffer"
910 );
911 }
912 
913 {
914 let q = &pair
915 .client_conn_mut(client_ch)
916 .streams
917 .get(&si)
918 .unwrap()
919 .reassembly_queue;
920 assert!(!q.is_readable(), "should no longer be readable");
921 }
922 
923 {
924 let a = pair.client_conn_mut(client_ch);
925 assert_eq!(0, a.buffered_amount(), "incorrect bufferedAmount");
926 }
927 
928 close_association_pair(&mut pair, client_ch, server_ch, si);
929 
930 Ok(())
931}
932 
933#[test]
934fn test_assoc_unreliable_rexmit_ordered_no_fragment() -> Result<()> {
935 //let _guard = subscribe();
936 
937 let si: u16 = 1;
938 let mut sbuf = vec![0u8; 1000];
939 for (i, b) in sbuf.iter_mut().enumerate() {
940 *b = (i & 0xff) as u8;
941 }
942 
943 let (mut pair, client_ch, server_ch) = create_association_pair(AckMode::NoDelay, 0)?;
944 
945 establish_session_pair(&mut pair, client_ch, server_ch, si)?;
946 
947 // When we set the reliability value to 0 [times], then it will cause
948 // the chunk to be abandoned immediately after the first transmission.
949 pair.client_stream(client_ch, si)?
950 .set_reliability_params(false, ReliabilityType::Rexmit, 0)?;
951 pair.server_stream(server_ch, si)?
952 .set_reliability_params(false, ReliabilityType::Rexmit, 0)?; // doesn't matter
953 
954 //br.drop_next_nwrites(0, 1).await; // drop the first packet (second one should be sacked)
955 
956 sbuf[0..4].copy_from_slice(&0u32.to_be_bytes());
957 let n = pair.client_stream(client_ch, si)?.write_sctp(
958 &Bytes::from(sbuf.clone()),
959 PayloadProtocolIdentifier::Binary,
960 )?;
961 assert_eq!(sbuf.len(), n, "unexpected length of received data");
962 pair.drive_client(); // send data to server
963 pair.server.inbound.clear(); // Lose it
964 debug!("dropping packet");
965 
966 sbuf[0..4].copy_from_slice(&1u32.to_be_bytes());
967 let n = pair.client_stream(client_ch, si)?.write_sctp(
968 &Bytes::from(sbuf.clone()),
969 PayloadProtocolIdentifier::Binary,
970 )?;
971 assert_eq!(sbuf.len(), n, "unexpected length of received data");
972 
973 debug!("flush_buffers");
974 pair.drive();
975 
976 let mut buf = vec![0u8; 2000];
977 
978 debug!("read_sctp");
979 let chunks = pair.server_stream(server_ch, si)?.read_sctp()?.unwrap();
980 let (n, ppi) = (chunks.len(), chunks.ppi);
981 chunks.read(&mut buf)?;
982 assert_eq!(n, sbuf.len(), "unexpected length of received data");
983 assert_eq!(ppi, PayloadProtocolIdentifier::Binary, "unexpected ppi");
984 assert_eq!(
985 u32::from_be_bytes([buf[0], buf[1], buf[2], buf[3]]),
986 1,
987 "unexpected received data"
988 );
989 
990 debug!("process");
991 pair.drive();
992 
993 {
994 let q = &pair
995 .client_conn_mut(client_ch)
996 .streams
997 .get(&si)
998 .unwrap()
999 .reassembly_queue;
1000 assert!(!q.is_readable(), "should no longer be readable");
1001 }
1002 
1003 close_association_pair(&mut pair, client_ch, server_ch, si);
1004 
1005 Ok(())
1006}
1007 
1008#[test]
1009fn test_assoc_unreliable_rexmit_ordered_fragment() -> Result<()> {
1010 //let _guard = subscribe();
1011 
1012 let si: u16 = 1;
1013 let mut sbuf = vec![0u8; 2000];
1014 for (i, b) in sbuf.iter_mut().enumerate() {
1015 *b = (i & 0xff) as u8;
1016 }
1017 
1018 let (mut pair, client_ch, server_ch) = create_association_pair(AckMode::NoDelay, 0)?;
1019 
1020 establish_session_pair(&mut pair, client_ch, server_ch, si)?;
1021 
1022 {
1023 // lock RTO value at 100 [msec]
1024 let a = pair.client_conn_mut(client_ch);
1025 a.rto_mgr.set_rto(100, true);
1026 }
1027 // When we set the reliability value to 0 [times], then it will cause
1028 // the chunk to be abandoned immediately after the first transmission.
1029 pair.client_stream(client_ch, si)?
1030 .set_reliability_params(false, ReliabilityType::Rexmit, 0)?;
1031 pair.server_stream(server_ch, si)?
1032 .set_reliability_params(false, ReliabilityType::Rexmit, 0)?; // doesn't matter
1033 
1034 //br.drop_next_nwrites(0, 1).await; // drop the first packet (second one should be sacked)
1035 
1036 sbuf[0..4].copy_from_slice(&0u32.to_be_bytes());
1037 let n = pair.client_stream(client_ch, si)?.write_sctp(
1038 &Bytes::from(sbuf.clone()),
1039 PayloadProtocolIdentifier::Binary,
1040 )?;
1041 assert_eq!(sbuf.len(), n, "unexpected length of received data");
1042 pair.drive_client(); // send data to server
1043 pair.server.inbound.clear(); // Lose it
1044 
1045 sbuf[0..4].copy_from_slice(&1u32.to_be_bytes());
1046 let n = pair.client_stream(client_ch, si)?.write_sctp(
1047 &Bytes::from(sbuf.clone()),
1048 PayloadProtocolIdentifier::Binary,
1049 )?;
1050 assert_eq!(sbuf.len(), n, "unexpected length of received data");
1051 
1052 //log::debug!("flush_buffers");
1053 pair.drive();
1054 
1055 let mut buf = vec![0u8; 2000];
1056 
1057 //log::debug!("read_sctp");
1058 let chunks = pair.server_stream(server_ch, si)?.read_sctp()?.unwrap();
1059 let (n, ppi) = (chunks.len(), chunks.ppi);
1060 chunks.read(&mut buf)?;
1061 assert_eq!(n, sbuf.len(), "unexpected length of received data");
1062 assert_eq!(ppi, PayloadProtocolIdentifier::Binary, "unexpected ppi");
1063 assert_eq!(
1064 u32::from_be_bytes([buf[0], buf[1], buf[2], buf[3]]),
1065 1,
1066 "unexpected received data"
1067 );
1068 
1069 //log::debug!("process");
1070 pair.drive();
1071 
1072 {
1073 let q = &pair
1074 .client_conn_mut(client_ch)
1075 .streams
1076 .get(&si)
1077 .unwrap()
1078 .reassembly_queue;
1079 assert!(!q.is_readable(), "should no longer be readable");
1080 }
1081 
1082 close_association_pair(&mut pair, client_ch, server_ch, si);
1083 
1084 Ok(())
1085}
1086 
1087#[test]
1088fn test_assoc_unreliable_rexmit_unordered_no_fragment() -> Result<()> {
1089 //let _guard = subscribe();
1090 
1091 let si: u16 = 2;
1092 let mut sbuf = vec![0u8; 1000];
1093 for (i, b) in sbuf.iter_mut().enumerate() {
1094 *b = (i & 0xff) as u8;
1095 }
1096 
1097 let (mut pair, client_ch, server_ch) = create_association_pair(AckMode::NoDelay, 0)?;
1098 
1099 establish_session_pair(&mut pair, client_ch, server_ch, si)?;
1100 
1101 // When we set the reliability value to 0 [times], then it will cause
1102 // the chunk to be abandoned immediately after the first transmission.
1103 pair.client_stream(client_ch, si)?
1104 .set_reliability_params(true, ReliabilityType::Rexmit, 0)?;
1105 pair.server_stream(server_ch, si)?
1106 .set_reliability_params(true, ReliabilityType::Rexmit, 0)?; // doesn't matter
1107 
1108 //br.drop_next_nwrites(0, 1).await; // drop the first packet (second one should be sacked)
1109 
1110 sbuf[0..4].copy_from_slice(&0u32.to_be_bytes());
1111 let n = pair.client_stream(client_ch, si)?.write_sctp(
1112 &Bytes::from(sbuf.clone()),
1113 PayloadProtocolIdentifier::Binary,
1114 )?;
1115 assert_eq!(sbuf.len(), n, "unexpected length of received data");
1116 pair.drive_client(); // send data to server
1117 pair.server.inbound.clear(); // Lose it
1118 
1119 sbuf[0..4].copy_from_slice(&1u32.to_be_bytes());
1120 let n = pair.client_stream(client_ch, si)?.write_sctp(
1121 &Bytes::from(sbuf.clone()),
1122 PayloadProtocolIdentifier::Binary,
1123 )?;
1124 assert_eq!(sbuf.len(), n, "unexpected length of received data");
1125 
1126 //log::debug!("flush_buffers");
1127 pair.drive();
1128 
1129 let mut buf = vec![0u8; 2000];
1130 
1131 //log::debug!("read_sctp");
1132 let chunks = pair.server_stream(server_ch, si)?.read_sctp()?.unwrap();
1133 let (n, ppi) = (chunks.len(), chunks.ppi);
1134 chunks.read(&mut buf)?;
1135 assert_eq!(n, sbuf.len(), "unexpected length of received data");
1136 assert_eq!(ppi, PayloadProtocolIdentifier::Binary, "unexpected ppi");
1137 assert_eq!(
1138 u32::from_be_bytes([buf[0], buf[1], buf[2], buf[3]]),
1139 1,
1140 "unexpected received data"
1141 );
1142 
1143 //log::debug!("process");
1144 pair.drive();
1145 
1146 {
1147 let q = &pair
1148 .client_conn_mut(client_ch)
1149 .streams
1150 .get(&si)
1151 .unwrap()
1152 .reassembly_queue;
1153 assert!(!q.is_readable(), "should no longer be readable");
1154 }
1155 
1156 close_association_pair(&mut pair, client_ch, server_ch, si);
1157 
1158 Ok(())
1159}
1160 
1161#[test]
1162fn test_assoc_unreliable_rexmit_unordered_fragment() -> Result<()> {
1163 //let _guard = subscribe();
1164 
1165 let si: u16 = 1;
1166 let mut sbuf = vec![0u8; 2000];
1167 for (i, b) in sbuf.iter_mut().enumerate() {
1168 *b = (i & 0xff) as u8;
1169 }
1170 
1171 let (mut pair, client_ch, server_ch) = create_association_pair(AckMode::NoDelay, 0)?;
1172 
1173 establish_session_pair(&mut pair, client_ch, server_ch, si)?;
1174 
1175 // When we set the reliability value to 0 [times], then it will cause
1176 // the chunk to be abandoned immediately after the first transmission.
1177 pair.client_stream(client_ch, si)?
1178 .set_reliability_params(true, ReliabilityType::Rexmit, 0)?;
1179 pair.server_stream(server_ch, si)?
1180 .set_reliability_params(true, ReliabilityType::Rexmit, 0)?; // doesn't matter
1181 
1182 sbuf[0..4].copy_from_slice(&0u32.to_be_bytes());
1183 let n = pair.client_stream(client_ch, si)?.write_sctp(
1184 &Bytes::from(sbuf.clone()),
1185 PayloadProtocolIdentifier::Binary,
1186 )?;
1187 assert_eq!(sbuf.len(), n, "unexpected length of received data");
1188 pair.client.drive(pair.time, pair.server.addr);
1189 pair.client.outbound.clear();
1190 //debug!("outbound len={}", pair.client.outbound.len());
1191 
1192 sbuf[0..4].copy_from_slice(&1u32.to_be_bytes());
1193 let n = pair.client_stream(client_ch, si)?.write_sctp(
1194 &Bytes::from(sbuf.clone()),
1195 PayloadProtocolIdentifier::Binary,
1196 )?;
1197 assert_eq!(sbuf.len(), n, "unexpected length of received data");
1198 
1199 pair.drive();
1200 
1201 let mut buf = vec![0u8; 2000];
1202 
1203 //log::debug!("read_sctp");
1204 let chunks = pair.server_stream(server_ch, si)?.read_sctp()?.unwrap();
1205 let (n, ppi) = (chunks.len(), chunks.ppi);
1206 chunks.read(&mut buf)?;
1207 assert_eq!(n, sbuf.len(), "unexpected length of received data");
1208 assert_eq!(ppi, PayloadProtocolIdentifier::Binary, "unexpected ppi");
1209 assert_eq!(
1210 u32::from_be_bytes([buf[0], buf[1], buf[2], buf[3]]),
1211 1,
1212 "unexpected received data"
1213 );
1214 
1215 //log::debug!("process");
1216 pair.drive();
1217 
1218 {
1219 let q = &pair
1220 .client_conn_mut(client_ch)
1221 .streams
1222 .get(&si)
1223 .unwrap()
1224 .reassembly_queue;
1225 assert!(!q.is_readable(), "should no longer be readable");
1226 assert_eq!(
1227 0,
1228 q.unordered.len(),
1229 "should be nothing in the unordered queue"
1230 );
1231 assert_eq!(
1232 0,
1233 q.unordered_chunks.len(),
1234 "should be nothing in the unorderedChunks list"
1235 );
1236 }
1237 
1238 close_association_pair(&mut pair, client_ch, server_ch, si);
1239 
1240 Ok(())
1241}
1242 
1243#[test]
1244fn test_assoc_unreliable_rexmit_timed_ordered() -> Result<()> {
1245 //let _guard = subscribe();
1246 
1247 let si: u16 = 3;
1248 let mut sbuf = vec![0u8; 1000];
1249 for (i, b) in sbuf.iter_mut().enumerate() {
1250 *b = (i & 0xff) as u8;
1251 }
1252 
1253 let (mut pair, client_ch, server_ch) = create_association_pair(AckMode::NoDelay, 0)?;
1254 
1255 establish_session_pair(&mut pair, client_ch, server_ch, si)?;
1256 
1257 // When we set the reliability value to 0 [times], then it will cause
1258 // the chunk to be abandoned immediately after the first transmission.
1259 pair.client_stream(client_ch, si)?
1260 .set_reliability_params(false, ReliabilityType::Timed, 0)?;
1261 pair.server_stream(server_ch, si)?
1262 .set_reliability_params(false, ReliabilityType::Timed, 0)?; // doesn't matter
1263 
1264 //br.drop_next_nwrites(0, 1).await; // drop the first packet (second one should be sacked)
1265 
1266 sbuf[0..4].copy_from_slice(&0u32.to_be_bytes());
1267 let n = pair.client_stream(client_ch, si)?.write_sctp(
1268 &Bytes::from(sbuf.clone()),
1269 PayloadProtocolIdentifier::Binary,
1270 )?;
1271 assert_eq!(sbuf.len(), n, "unexpected length of received data");
1272 pair.client.drive(pair.time, pair.server.addr);
1273 pair.client.outbound.clear();
1274 
1275 sbuf[0..4].copy_from_slice(&1u32.to_be_bytes());
1276 let n = pair.client_stream(client_ch, si)?.write_sctp(
1277 &Bytes::from(sbuf.clone()),
1278 PayloadProtocolIdentifier::Binary,
1279 )?;
1280 assert_eq!(sbuf.len(), n, "unexpected length of received data");
1281 
1282 //log::debug!("flush_buffers");
1283 pair.drive();
1284 
1285 let mut buf = vec![0u8; 2000];
1286 
1287 //log::debug!("read_sctp");
1288 let chunks = pair.server_stream(server_ch, si)?.read_sctp()?.unwrap();
1289 let (n, ppi) = (chunks.len(), chunks.ppi);
1290 chunks.read(&mut buf)?;
1291 assert_eq!(n, sbuf.len(), "unexpected length of received data");
1292 assert_eq!(ppi, PayloadProtocolIdentifier::Binary, "unexpected ppi");
1293 assert_eq!(
1294 u32::from_be_bytes([buf[0], buf[1], buf[2], buf[3]]),
1295 1,
1296 "unexpected received data"
1297 );
1298 
1299 //log::debug!("process");
1300 pair.drive();
1301 
1302 {
1303 let q = &pair
1304 .client_conn_mut(client_ch)
1305 .streams
1306 .get(&si)
1307 .unwrap()
1308 .reassembly_queue;
1309 assert!(!q.is_readable(), "should no longer be readable");
1310 }
1311 
1312 close_association_pair(&mut pair, client_ch, server_ch, si);
1313 
1314 Ok(())
1315}
1316 
1317#[test]
1318fn test_assoc_unreliable_rexmit_timed_unordered() -> Result<()> {
1319 //let _guard = subscribe();
1320 
1321 let si: u16 = 3;
1322 let mut sbuf = vec![0u8; 1000];
1323 for (i, b) in sbuf.iter_mut().enumerate() {
1324 *b = (i & 0xff) as u8;
1325 }
1326 
1327 let (mut pair, client_ch, server_ch) = create_association_pair(AckMode::NoDelay, 0)?;
1328 
1329 establish_session_pair(&mut pair, client_ch, server_ch, si)?;
1330 
1331 // When we set the reliability value to 0 [times], then it will cause
1332 // the chunk to be abandoned immediately after the first transmission.
1333 pair.client_stream(client_ch, si)?
1334 .set_reliability_params(true, ReliabilityType::Timed, 0)?;
1335 pair.server_stream(server_ch, si)?
1336 .set_reliability_params(true, ReliabilityType::Timed, 0)?; // doesn't matter
1337 
1338 //br.drop_next_nwrites(0, 1).await; // drop the first packet (second one should be sacked)
1339 
1340 sbuf[0..4].copy_from_slice(&0u32.to_be_bytes());
1341 let n = pair.client_stream(client_ch, si)?.write_sctp(
1342 &Bytes::from(sbuf.clone()),
1343 PayloadProtocolIdentifier::Binary,
1344 )?;
1345 assert_eq!(sbuf.len(), n, "unexpected length of received data");
1346 pair.client.drive(pair.time, pair.server.addr);
1347 pair.client.outbound.clear();
1348 
1349 sbuf[0..4].copy_from_slice(&1u32.to_be_bytes());
1350 let n = pair.client_stream(client_ch, si)?.write_sctp(
1351 &Bytes::from(sbuf.clone()),
1352 PayloadProtocolIdentifier::Binary,
1353 )?;
1354 assert_eq!(sbuf.len(), n, "unexpected length of received data");
1355 
1356 //log::debug!("flush_buffers");
1357 pair.drive();
1358 
1359 let mut buf = vec![0u8; 2000];
1360 
1361 //log::debug!("read_sctp");
1362 let chunks = pair.server_stream(server_ch, si)?.read_sctp()?.unwrap();
1363 let (n, ppi) = (chunks.len(), chunks.ppi);
1364 chunks.read(&mut buf)?;
1365 assert_eq!(n, sbuf.len(), "unexpected length of received data");
1366 assert_eq!(ppi, PayloadProtocolIdentifier::Binary, "unexpected ppi");
1367 assert_eq!(
1368 u32::from_be_bytes([buf[0], buf[1], buf[2], buf[3]]),
1369 1,
1370 "unexpected received data"
1371 );
1372 
1373 //log::debug!("process");
1374 pair.drive();
1375 
1376 {
1377 let q = &pair
1378 .client_conn_mut(client_ch)
1379 .streams
1380 .get(&si)
1381 .unwrap()
1382 .reassembly_queue;
1383 assert!(!q.is_readable(), "should no longer be readable");
1384 assert_eq!(
1385 0,
1386 q.unordered.len(),
1387 "should be nothing in the unordered queue"
1388 );
1389 assert_eq!(
1390 0,
1391 q.unordered_chunks.len(),
1392 "should be nothing in the unorderedChunks list"
1393 );
1394 }
1395 
1396 close_association_pair(&mut pair, client_ch, server_ch, si);
1397 
1398 Ok(())
1399}
1400 
1401//TODO: TestAssocT1InitTimer
1402//TODO: TestAssocT1CookieTimer
1403//TODO: TestAssocT3RtxTimer
1404 
1405/*FIXME
1406// 1) Send 4 packets. drop the first one.
1407// 2) Last 3 packet will be received, which triggers fast-retransmission
1408// 3) The first one is retransmitted, which makes s1 readable
1409// Above should be done before RTO occurs (fast recovery)
1410#[test]
1411fn test_assoc_congestion_control_fast_retransmission() -> Result<()> {
1412 let _guard = subscribe();
1413
1414 let si: u16 = 6;
1415 let mut sbuf = vec![0u8; 1000];
1416 for i in 0..sbuf.len() {
1417 sbuf[i] = (i & 0xff) as u8;
1418 }
1419
1420 let (mut pair, client_ch, server_ch) = create_association_pair(AckMode::Normal, 0)?;
1421
1422 establish_session_pair(&mut pair, client_ch, server_ch, si)?;
1423
1424 //br.drop_next_nwrites(0, 1).await; // drop the first packet (second one should be sacked)
1425
1426 for i in 0..4u32 {
1427 sbuf[0..4].copy_from_slice(&i.to_be_bytes());
1428 let n = pair.client_stream(client_ch, si)?.write_sctp(
1429 &Bytes::from(sbuf.clone()),
1430 PayloadProtocolIdentifier::Binary,
1431 )?;
1432 assert_eq!(sbuf.len(), n, "unexpected length of received data");
1433 pair.client.drive(pair.time, pair.server.addr);
1434 if i == 0 {
1435 //drop the first packet
1436 pair.client.outbound.clear();
1437 }
1438 }
1439
1440 // process packets for 500 msec, assuming that the fast retrans/recover
1441 // should complete within 500 msec.
1442 /*for _ in 0..50 {
1443 br.tick().await;
1444 tokio::time::sleep(Duration::from_millis(10)).await;
1445 }*/
1446 debug!("advance 500ms");
1447 pair.time += Duration::from_millis(500);
1448 pair.step();
1449
1450 let mut buf = vec![0u8; 3000];
1451
1452 // Try to read all 4 packets
1453 for i in 0..4 {
1454 {
1455 let q = &pair
1456 .server_conn_mut(server_ch)
1457 .streams
1458 .get(&si)
1459 .unwrap()
1460 .reassembly_queue;
1461 assert!(q.is_readable(), "should be readable at {}", i);
1462 }
1463
1464 let chunks = pair.server_stream(server_ch, si)?.read_sctp()?.unwrap();
1465 let (n, ppi) = (chunks.len(), chunks.ppi);
1466 chunks.read(&mut buf)?;
1467 assert_eq!(n, sbuf.len(), "unexpected length of received data");
1468 assert_eq!(ppi, PayloadProtocolIdentifier::Binary, "unexpected ppi");
1469 assert_eq!(
1470 u32::from_be_bytes([buf[0], buf[1], buf[2], buf[3]]),
1471 i,
1472 "unexpected received data"
1473 );
1474 }
1475
1476 pair.drive();
1477 //br.process().await;
1478
1479 {
1480 let a = pair.client_conn_mut(client_ch);
1481 assert!(!a.in_fast_recovery, "should not be in fast-recovery");
1482 debug!("nSACKs : {}", a.stats.get_num_sacks());
1483 debug!("nFastRetrans: {}", a.stats.get_num_fast_retrans());
1484
1485 assert_eq!(1, a.stats.get_num_fast_retrans(), "should be 1");
1486 }
1487 {
1488 let b = pair.server_conn_mut(server_ch);
1489 debug!("nDATAs : {}", b.stats.get_num_datas());
1490 debug!("nAckTimeouts: {}", b.stats.get_num_ack_timeouts());
1491 }
1492
1493 close_association_pair(&mut pair, client_ch, server_ch, si);
1494
1495 Ok(())
1496}*/
1497 
1498#[test]
1499fn test_assoc_congestion_control_congestion_avoidance() -> Result<()> {
1500 //let _guard = subscribe();
1501 
1502 let max_receive_buffer_size: u32 = 64 * 1024;
1503 let si: u16 = 6;
1504 let n_packets_to_send: u32 = 2000;
1505 
1506 let mut sbuf = vec![0u8; 1000];
1507 for (i, b) in sbuf.iter_mut().enumerate() {
1508 *b = (i & 0xff) as u8;
1509 }
1510 
1511 let (mut pair, client_ch, server_ch) =
1512 create_association_pair(AckMode::Normal, max_receive_buffer_size)?;
1513 
1514 establish_session_pair(&mut pair, client_ch, server_ch, si)?;
1515 
1516 {
1517 pair.client_conn_mut(client_ch).stats.reset();
1518 pair.server_conn_mut(server_ch).stats.reset();
1519 }
1520 
1521 for i in 0..n_packets_to_send {
1522 sbuf[0..4].copy_from_slice(&i.to_be_bytes());
1523 let n = pair.client_stream(client_ch, si)?.write_sctp(
1524 &Bytes::from(sbuf.clone()),
1525 PayloadProtocolIdentifier::Binary,
1526 )?;
1527 assert_eq!(sbuf.len(), n, "unexpected length of received data");
1528 }
1529 pair.drive_client();
1530 //debug!("pair.drive_client() done");
1531 
1532 let mut rbuf = vec![0u8; 3000];
1533 
1534 // Repeat calling br.Tick() until the buffered amount becomes 0
1535 let mut n_packets_received = 0u32;
1536 while pair.client_conn_mut(client_ch).buffered_amount() > 0
1537 && n_packets_received < n_packets_to_send
1538 {
1539 /*println!("timestamp: {:?}", pair.time);
1540 println!(
1541 "buffered_amount {}, pair.server.inbound {}, n_packets_received {}, n_packets_to_send {}",
1542 pair.client_conn_mut(client_ch).buffered_amount(),
1543 pair.server.inbound.len(),
1544 n_packets_received,
1545 n_packets_to_send
1546 );*/
1547 
1548 pair.step();
1549 
1550 while let Some(chunks) = pair.server_stream(server_ch, si)?.read_sctp()? {
1551 let (n, ppi) = (chunks.len(), chunks.ppi);
1552 chunks.read(&mut rbuf)?;
1553 assert_eq!(sbuf.len(), n, "unexpected length of received data");
1554 assert_eq!(
1555 n_packets_received,
1556 u32::from_be_bytes([rbuf[0], rbuf[1], rbuf[2], rbuf[3]]),
1557 "unexpected length of received data"
1558 );
1559 assert_eq!(ppi, PayloadProtocolIdentifier::Binary, "unexpected ppi");
1560 
1561 n_packets_received += 1;
1562 }
1563 
1564 //pair.drive_client();
1565 }
1566 
1567 pair.drive();
1568 //println!("timestamp: {:?}", pair.time);
1569 
1570 assert_eq!(
1571 n_packets_received, n_packets_to_send,
1572 "unexpected num of packets received"
1573 );
1574 
1575 {
1576 let a = pair.client_conn_mut(client_ch);
1577 
1578 assert!(!a.in_fast_recovery, "should not be in fast-recovery");
1579 assert!(
1580 a.cwnd > a.ssthresh,
1581 "should be in congestion avoidance mode"
1582 );
1583 assert!(
1584 a.ssthresh >= max_receive_buffer_size,
1585 "{} should not be less than the initial size of 128KB {}",
1586 a.ssthresh,
1587 max_receive_buffer_size
1588 );
1589 
1590 debug!("nSACKs : {}", a.stats.get_num_sacks());
1591 debug!("nT3Timeouts: {}", a.stats.get_num_t3timeouts());
1592 
1593 assert!(
1594 a.stats.get_num_sacks() <= n_packets_to_send as u64 / 2,
1595 "too many sacks"
1596 );
1597 assert_eq!(0, a.stats.get_num_t3timeouts(), "should be no retransmit");
1598 }
1599 {
1600 assert_eq!(
1601 0,
1602 pair.server_conn_mut(server_ch)
1603 .streams
1604 .get(&si)
1605 .unwrap()
1606 .get_num_bytes_in_reassembly_queue(),
1607 "reassembly queue should be empty"
1608 );
1609 
1610 let b = pair.server_conn_mut(server_ch);
1611 
1612 debug!("nDATAs : {}", b.stats.get_num_datas());
1613 
1614 assert_eq!(
1615 n_packets_to_send as u64,
1616 b.stats.get_num_datas(),
1617 "packet count mismatch"
1618 );
1619 }
1620 
1621 close_association_pair(&mut pair, client_ch, server_ch, si);
1622 
1623 Ok(())
1624}
1625 
1626#[test]
1627fn test_assoc_congestion_control_slow_reader() -> Result<()> {
1628 //let _guard = subscribe();
1629 
1630 let max_receive_buffer_size: u32 = 64 * 1024;
1631 let si: u16 = 6;
1632 let n_packets_to_send: u32 = 130;
1633 
1634 let mut sbuf = vec![0u8; 1000];
1635 for (i, b) in sbuf.iter_mut().enumerate() {
1636 *b = (i & 0xff) as u8;
1637 }
1638 
1639 let (mut pair, client_ch, server_ch) =
1640 create_association_pair(AckMode::Normal, max_receive_buffer_size)?;
1641 
1642 establish_session_pair(&mut pair, client_ch, server_ch, si)?;
1643 
1644 for i in 0..n_packets_to_send {
1645 sbuf[0..4].copy_from_slice(&i.to_be_bytes());
1646 let n = pair.client_stream(client_ch, si)?.write_sctp(
1647 &Bytes::from(sbuf.clone()),
1648 PayloadProtocolIdentifier::Binary,
1649 )?;
1650 assert_eq!(sbuf.len(), n, "unexpected length of received data");
1651 }
1652 pair.drive_client();
1653 
1654 let mut rbuf = vec![0u8; 3000];
1655 
1656 // 1. First forward packets to receiver until rwnd becomes 0
1657 // 2. Wait until the sender's cwnd becomes 1*MTU (RTO occurred)
1658 // 3. Stat reading a1's data
1659 let mut n_packets_received = 0u32;
1660 let mut has_rtoed = false;
1661 while pair.client_conn_mut(client_ch).buffered_amount() > 0
1662 && n_packets_received < n_packets_to_send
1663 {
1664 /*println!(
1665 "buffered_amount {}, pair.server.inbound {}, n_packets_received {}, n_packets_to_send {}",
1666 pair.client_conn_mut(client_ch).buffered_amount(),
1667 pair.server.inbound.len(),
1668 n_packets_received,
1669 n_packets_to_send
1670 );*/
1671 
1672 if !has_rtoed {
1673 let rwnd = pair
1674 .server_conn_mut(server_ch)
1675 .get_my_receiver_window_credit();
1676 let cwnd = pair.client_conn_mut(client_ch).cwnd;
1677 let cmtu = pair.client_conn_mut(client_ch).mtu;
1678 if cwnd > cmtu || rwnd > 0 {
1679 // Do not read until a1.getMyReceiverWindowCredit() becomes zero
1680 pair.step();
1681 continue;
1682 }
1683 
1684 has_rtoed = true;
1685 }
1686 
1687 while let Some(chunks) = pair.server_stream(server_ch, si)?.read_sctp()? {
1688 let (n, ppi) = (chunks.len(), chunks.ppi);
1689 chunks.read(&mut rbuf)?;
1690 assert_eq!(sbuf.len(), n, "unexpected length of received data");
1691 assert_eq!(
1692 n_packets_received,
1693 u32::from_be_bytes([rbuf[0], rbuf[1], rbuf[2], rbuf[3]]),
1694 "unexpected length of received data"
1695 );
1696 assert_eq!(ppi, PayloadProtocolIdentifier::Binary, "unexpected ppi");
1697 
1698 n_packets_received += 1;
1699 }
1700 
1701 pair.step();
1702 }
1703 
1704 pair.drive();
1705 
1706 assert_eq!(
1707 n_packets_received, n_packets_to_send,
1708 "unexpected num of packets received"
1709 );
1710 assert_eq!(
1711 0,
1712 pair.server_conn_mut(server_ch)
1713 .streams
1714 .get(&si)
1715 .unwrap()
1716 .get_num_bytes_in_reassembly_queue(),
1717 "reassembly queue should be empty"
1718 );
1719 
1720 {
1721 let a = pair.client_conn_mut(client_ch);
1722 debug!("nSACKs : {}", a.stats.get_num_sacks());
1723 }
1724 {
1725 let b = pair.server_conn_mut(server_ch);
1726 debug!("nDATAs : {}", b.stats.get_num_datas());
1727 debug!("nAckTimeouts: {}", b.stats.get_num_ack_timeouts());
1728 }
1729 
1730 close_association_pair(&mut pair, client_ch, server_ch, si);
1731 
1732 Ok(())
1733}
1734 
1735#[test]
1736fn test_assoc_delayed_ack() -> Result<()> {
1737 //let _guard = subscribe();
1738 
1739 let si: u16 = 6;
1740 let mut sbuf = vec![0u8; 1000];
1741 let mut rbuf = vec![0u8; 1500];
1742 for (i, b) in sbuf.iter_mut().enumerate() {
1743 *b = (i & 0xff) as u8;
1744 }
1745 
1746 let (mut pair, client_ch, server_ch) = create_association_pair(AckMode::AlwaysDelay, 0)?;
1747 
1748 establish_session_pair(&mut pair, client_ch, server_ch, si)?;
1749 
1750 {
1751 pair.client_conn_mut(client_ch).stats.reset();
1752 pair.server_conn_mut(server_ch).stats.reset();
1753 }
1754 
1755 let n = pair.client_stream(client_ch, si)?.write_sctp(
1756 &Bytes::from(sbuf.clone()),
1757 PayloadProtocolIdentifier::Binary,
1758 )?;
1759 assert_eq!(sbuf.len(), n, "unexpected length of received data");
1760 pair.drive_client();
1761 
1762 // Repeat calling br.Tick() until the buffered amount becomes 0
1763 let since = pair.time;
1764 let mut n_packets_received = 0;
1765 while pair.client_conn_mut(client_ch).buffered_amount() > 0 {
1766 pair.step();
1767 
1768 while let Some(chunks) = pair.server_stream(server_ch, si)?.read_sctp()? {
1769 let (n, ppi) = (chunks.len(), chunks.ppi);
1770 chunks.read(&mut rbuf)?;
1771 assert_eq!(sbuf.len(), n, "unexpected length of received data");
1772 assert_eq!(ppi, PayloadProtocolIdentifier::Binary, "unexpected ppi");
1773 
1774 n_packets_received += 1;
1775 }
1776 }
1777 let delay = (pair.time.duration_since(since).as_millis() as f64) / 1000.0;
1778 debug!("received in {} seconds", delay);
1779 assert!(delay >= 0.2, "should be >= 200msec");
1780 
1781 pair.drive();
1782 
1783 assert_eq!(n_packets_received, 1, "unexpected num of packets received");
1784 assert_eq!(
1785 0,
1786 pair.server_conn_mut(server_ch)
1787 .streams
1788 .get(&si)
1789 .unwrap()
1790 .get_num_bytes_in_reassembly_queue(),
1791 "reassembly queue should be empty"
1792 );
1793 
1794 let a_num_sacks = {
1795 let a = pair.client_conn_mut(client_ch);
1796 debug!("nSACKs : {}", a.stats.get_num_sacks());
1797 assert_eq!(0, a.stats.get_num_t3timeouts(), "should be no retransmit");
1798 a.stats.get_num_sacks()
1799 };
1800 
1801 {
1802 let b = pair.server_conn_mut(server_ch);
1803 
1804 debug!("nDATAs : {}", b.stats.get_num_datas());
1805 debug!("nAckTimeouts: {}", b.stats.get_num_ack_timeouts());
1806 
1807 assert_eq!(1, b.stats.get_num_datas(), "DATA chunk count mismatch");
1808 assert_eq!(
1809 a_num_sacks,
1810 b.stats.get_num_datas(),
1811 "sack count should be equal to the number of data chunks"
1812 );
1813 assert_eq!(
1814 1,
1815 b.stats.get_num_ack_timeouts(),
1816 "ackTimeout count mismatch"
1817 );
1818 }
1819 
1820 close_association_pair(&mut pair, client_ch, server_ch, si);
1821 
1822 Ok(())
1823}
1824 
1825#[test]
1826fn test_assoc_reset_close_one_way() -> Result<()> {
1827 //let _guard = subscribe();
1828 
1829 let si: u16 = 1;
1830 let msg: Bytes = Bytes::from_static(b"ABC");
1831 
1832 let (mut pair, client_ch, server_ch) = create_association_pair(AckMode::NoDelay, 0)?;
1833 
1834 establish_session_pair(&mut pair, client_ch, server_ch, si)?;
1835 
1836 {
1837 let a = pair.client_conn_mut(client_ch);
1838 assert_eq!(0, a.buffered_amount(), "incorrect bufferedAmount");
1839 }
1840 
1841 let n = pair
1842 .client_stream(client_ch, si)?
1843 .write_sctp(&msg, PayloadProtocolIdentifier::Binary)?;
1844 assert_eq!(msg.len(), n, "unexpected length of received data");
1845 {
1846 let a = pair.client_conn_mut(client_ch);
1847 assert_eq!(msg.len(), a.buffered_amount(), "incorrect bufferedAmount");
1848 }
1849 pair.step();
1850 
1851 let mut buf = vec![0u8; 32];
1852 
1853 while pair.server_stream(server_ch, si).is_ok() {
1854 debug!("s1.read_sctp begin");
1855 match pair.server_stream(server_ch, si)?.read_sctp() {
1856 Ok(chunks_opt) => {
1857 if let Some(chunks) = chunks_opt {
1858 let (n, ppi) = (chunks.len(), chunks.ppi);
1859 chunks.read(&mut buf)?;
1860 debug!("s1.read_sctp done with {:?}", &buf[..n]);
1861 assert_eq!(ppi, PayloadProtocolIdentifier::Binary, "unexpected ppi");
1862 assert_eq!(n, msg.len(), "unexpected length of received data");
1863 }
1864 
1865 debug!("s0.close");
1866 pair.client_stream(client_ch, si)?.stop()?; // send reset
1867 
1868 pair.step();
1869 }
1870 Err(err) => {
1871 debug!("s1.read_sctp err {:?}", err);
1872 break;
1873 }
1874 }
1875 }
1876 
1877 pair.drive();
1878 
1879 close_association_pair(&mut pair, client_ch, server_ch, si);
1880 
1881 Ok(())
1882}
1883 
1884#[test]
1885fn test_assoc_reset_close_both_ways() -> Result<()> {
1886 //let _guard = subscribe();
1887 
1888 let si: u16 = 1;
1889 let msg: Bytes = Bytes::from_static(b"ABC");
1890 
1891 let (mut pair, client_ch, server_ch) = create_association_pair(AckMode::NoDelay, 0)?;
1892 
1893 establish_session_pair(&mut pair, client_ch, server_ch, si)?;
1894 
1895 {
1896 let a = pair.client_conn_mut(client_ch);
1897 assert_eq!(0, a.buffered_amount(), "incorrect bufferedAmount");
1898 }
1899 
1900 let n = pair
1901 .client_stream(client_ch, si)?
1902 .write_sctp(&msg, PayloadProtocolIdentifier::Binary)?;
1903 assert_eq!(msg.len(), n, "unexpected length of received data");
1904 {
1905 let a = pair.client_conn_mut(client_ch);
1906 assert_eq!(msg.len(), a.buffered_amount(), "incorrect bufferedAmount");
1907 }
1908 pair.step();
1909 
1910 let mut buf = vec![0u8; 32];
1911 
1912 while pair.server_stream(server_ch, si).is_ok() || pair.client_stream(client_ch, si).is_ok() {
1913 if pair.server_stream(server_ch, si).is_ok() {
1914 debug!("s1.read_sctp begin");
1915 match pair.server_stream(server_ch, si)?.read_sctp() {
1916 Ok(chunks_opt) => {
1917 if let Some(chunks) = chunks_opt {
1918 let (n, ppi) = (chunks.len(), chunks.ppi);
1919 chunks.read(&mut buf)?;
1920 debug!("s1.read_sctp done with {:?}", &buf[..n]);
1921 assert_eq!(ppi, PayloadProtocolIdentifier::Binary, "unexpected ppi");
1922 assert_eq!(n, msg.len(), "unexpected length of received data");
1923 }
1924 }
1925 Err(err) => {
1926 debug!("s1.read_sctp err {:?}", err);
1927 break;
1928 }
1929 }
1930 }
1931 
1932 if pair.client_stream(client_ch, si).is_ok() {
1933 debug!("s0.read_sctp begin");
1934 match pair.client_stream(client_ch, si)?.read_sctp() {
1935 Ok(chunks_opt) => {
1936 if let Some(chunks) = chunks_opt {
1937 let (n, ppi) = (chunks.len(), chunks.ppi);
1938 chunks.read(&mut buf)?;
1939 debug!("s0.read_sctp done with {:?}", &buf[..n]);
1940 assert_eq!(ppi, PayloadProtocolIdentifier::Binary, "unexpected ppi");
1941 assert_eq!(n, msg.len(), "unexpected length of received data");
1942 }
1943 }
1944 Err(err) => {
1945 debug!("s0.read_sctp err {:?}", err);
1946 break;
1947 }
1948 }
1949 }
1950 
1951 if pair.client_stream(client_ch, si).is_ok() {
1952 pair.client_stream(client_ch, si)?.stop()?; // send reset
1953 }
1954 if pair.server_stream(server_ch, si).is_ok() {
1955 pair.server_stream(server_ch, si)?.stop()?; // send reset
1956 }
1957 
1958 pair.step();
1959 }
1960 
1961 pair.drive();
1962 
1963 close_association_pair(&mut pair, client_ch, server_ch, si);
1964 
1965 Ok(())
1966}
1967 
1968#[test]
1969fn test_assoc_abort() -> Result<()> {
1970 //let _guard = subscribe();
1971 
1972 let si: u16 = 1;
1973 
1974 let (mut pair, client_ch, server_ch) = create_association_pair(AckMode::NoDelay, 0)?;
1975 
1976 establish_session_pair(&mut pair, client_ch, server_ch, si)?;
1977 
1978 let transmit = {
1979 let abort = ChunkAbort {
1980 error_causes: vec![ErrorCauseProtocolViolation {
1981 code: PROTOCOL_VIOLATION,
1982 ..Default::default()
1983 }],
1984 };
1985 
1986 let packet = pair
1987 .client_conn_mut(client_ch)
1988 .create_packet(vec![Box::new(abort)])
1989 .marshal()?;
1990 
1991 Transmit {
1992 now: pair.time,
1993 remote: pair.server.addr,
1994 ecn: None,
1995 local_ip: None,
1996 payload: Payload::RawEncode(vec![packet]),
1997 }
1998 };
1999 
2000 // Both associations are established
2001 assert_eq!(
2002 AssociationState::Established,
2003 pair.client_conn_mut(client_ch).state()
2004 );
2005 assert_eq!(
2006 AssociationState::Established,
2007 pair.server_conn_mut(server_ch).state()
2008 );
2009 
2010 debug!("send ChunkAbort");
2011 pair.client.outbound.push_back(transmit);
2012 
2013 pair.drive();
2014 
2015 // The receiving association should be closed because it got an ABORT
2016 assert_eq!(
2017 AssociationState::Established,
2018 pair.client_conn_mut(client_ch).state()
2019 );
2020 assert_eq!(
2021 AssociationState::Closed,
2022 pair.server_conn_mut(server_ch).state()
2023 );
2024 
2025 close_association_pair(&mut pair, client_ch, server_ch, si);
2026 
2027 Ok(())
2028}
2029 
2030#[test]
2031fn test_association_handle_packet_before_init() -> Result<()> {
2032 //let _guard = subscribe();
2033 
2034 let tests = vec![
2035 (
2036 "InitAck",
2037 Packet {
2038 common_header: CommonHeader {
2039 source_port: 1,
2040 destination_port: 1,
2041 verification_tag: 0,
2042 },
2043 chunks: vec![Box::new(ChunkInit {
2044 is_ack: true,
2045 initiate_tag: 1,
2046 num_inbound_streams: 1,
2047 num_outbound_streams: 1,
2048 advertised_receiver_window_credit: 1500,
2049 ..Default::default()
2050 })],
2051 },
2052 ),
2053 (
2054 "Abort",
2055 Packet {
2056 common_header: CommonHeader {
2057 source_port: 1,
2058 destination_port: 1,
2059 verification_tag: 0,
2060 },
2061 chunks: vec![Box::<ChunkAbort>::default()],
2062 },
2063 ),
2064 (
2065 "CoockeEcho",
2066 Packet {
2067 common_header: CommonHeader {
2068 source_port: 1,
2069 destination_port: 1,
2070 verification_tag: 0,
2071 },
2072 chunks: vec![Box::<ChunkCookieEcho>::default()],
2073 },
2074 ),
2075 (
2076 "HeartBeat",
2077 Packet {
2078 common_header: CommonHeader {
2079 source_port: 1,
2080 destination_port: 1,
2081 verification_tag: 0,
2082 },
2083 chunks: vec![Box::<ChunkHeartbeat>::default()],
2084 },
2085 ),
2086 (
2087 "PayloadData",
2088 Packet {
2089 common_header: CommonHeader {
2090 source_port: 1,
2091 destination_port: 1,
2092 verification_tag: 0,
2093 },
2094 chunks: vec![Box::<ChunkPayloadData>::default()],
2095 },
2096 ),
2097 (
2098 "Sack",
2099 Packet {
2100 common_header: CommonHeader {
2101 source_port: 1,
2102 destination_port: 1,
2103 verification_tag: 0,
2104 },
2105 chunks: vec![Box::new(ChunkSelectiveAck {
2106 cumulative_tsn_ack: 1000,
2107 advertised_receiver_window_credit: 1500,
2108 gap_ack_blocks: vec![GapAckBlock {
2109 start: 100,
2110 end: 200,
2111 }],
2112 ..Default::default()
2113 })],
2114 },
2115 ),
2116 (
2117 "Reconfig",
2118 Packet {
2119 common_header: CommonHeader {
2120 source_port: 1,
2121 destination_port: 1,
2122 verification_tag: 0,
2123 },
2124 chunks: vec![Box::new(ChunkReconfig {
2125 param_a: Some(Box::<ParamOutgoingResetRequest>::default()),
2126 param_b: Some(Box::<ParamReconfigResponse>::default()),
2127 })],
2128 },
2129 ),
2130 (
2131 "ForwardTSN",
2132 Packet {
2133 common_header: CommonHeader {
2134 source_port: 1,
2135 destination_port: 1,
2136 verification_tag: 0,
2137 },
2138 chunks: vec![Box::new(ChunkForwardTsn {
2139 new_cumulative_tsn: 100,
2140 ..Default::default()
2141 })],
2142 },
2143 ),
2144 (
2145 "Error",
2146 Packet {
2147 common_header: CommonHeader {
2148 source_port: 1,
2149 destination_port: 1,
2150 verification_tag: 0,
2151 },
2152 chunks: vec![Box::<ChunkError>::default()],
2153 },
2154 ),
2155 (
2156 "Shutdown",
2157 Packet {
2158 common_header: CommonHeader {
2159 source_port: 1,
2160 destination_port: 1,
2161 verification_tag: 0,
2162 },
2163 chunks: vec![Box::<ChunkShutdown>::default()],
2164 },
2165 ),
2166 (
2167 "ShutdownAck",
2168 Packet {
2169 common_header: CommonHeader {
2170 source_port: 1,
2171 destination_port: 1,
2172 verification_tag: 0,
2173 },
2174 chunks: vec![Box::new(ChunkShutdownAck)],
2175 },
2176 ),
2177 (
2178 "ShutdownComplete",
2179 Packet {
2180 common_header: CommonHeader {
2181 source_port: 1,
2182 destination_port: 1,
2183 verification_tag: 0,
2184 },
2185 chunks: vec![Box::new(ChunkShutdownComplete)],
2186 },
2187 ),
2188 ];
2189 
2190 let remote = SocketAddr::from_str("0.0.0.0:0").unwrap();
2191 
2192 for (name, packet) in tests {
2193 debug!("testing {}", name);
2194 
2195 //let (a_conn, charlie_conn) = pipe();
2196 let config = Arc::new(TransportConfig::default());
2197 let mut a = Association::new(None, config, 1400, 0, remote, None, Instant::now());
2198 
2199 let packet = packet.marshal()?;
2200 a.handle_event(AssociationEvent(AssociationEventInner::Datagram(
2201 Transmit {
2202 now: Instant::now(),
2203 remote,
2204 ecn: None,
2205 local_ip: None,
2206 payload: Payload::RawEncode(vec![packet]),
2207 },
2208 )));
2209 
2210 a.close()?;
2211 }
2212 
2213 Ok(())
2214}
2215 
2216// This test reproduces an issue related to having regular messages (regular acks) which keep
2217// rescheduling the T3RTX timer before it can ever fire.
2218#[test]
2219fn test_old_rtx_on_regular_acks() -> Result<()> {
2220 let si: u16 = 6;
2221 let mut sbuf = vec![0u8; 1000];
2222 for (i, b) in sbuf.iter_mut().enumerate() {
2223 *b = (i & 0xff) as u8;
2224 }
2225 
2226 let (mut pair, client_ch, server_ch) = create_association_pair(AckMode::Normal, 0)?;
2227 pair.latency = Duration::from_millis(500);
2228 establish_session_pair(&mut pair, client_ch, server_ch, si)?;
2229 
2230 // Send 20 packet at a regular interval that is < RTO
2231 for i in 0..20u32 {
2232 println!("sending packet {}", i);
2233 sbuf[0..4].copy_from_slice(&i.to_be_bytes());
2234 let n = pair.client_stream(client_ch, si)?.write_sctp(
2235 &Bytes::from(sbuf.clone()),
2236 PayloadProtocolIdentifier::Binary,
2237 )?;
2238 assert_eq!(sbuf.len(), n, "unexpected length of received data");
2239 pair.client.drive(pair.time, pair.server.addr);
2240 
2241 // drop a few transmits
2242 if (5..10).contains(&i) {
2243 pair.client.outbound.clear();
2244 }
2245 
2246 pair.drive_client();
2247 pair.drive_server();
2248 pair.time += Duration::from_millis(500);
2249 }
2250 
2251 pair.drive_client();
2252 pair.drive_server();
2253 
2254 let mut buf = vec![0u8; 3000];
2255 
2256 // All packets must readable correctly
2257 for i in 0..20 {
2258 {
2259 let q = &pair
2260 .server_conn_mut(server_ch)
2261 .streams
2262 .get(&si)
2263 .unwrap()
2264 .reassembly_queue;
2265 println!("q.is_readable()={}", q.is_readable());
2266 assert!(q.is_readable(), "should be readable at {}", i);
2267 }
2268 
2269 let chunks = pair.server_stream(server_ch, si)?.read_sctp()?.unwrap();
2270 let (n, ppi) = (chunks.len(), chunks.ppi);
2271 chunks.read(&mut buf)?;
2272 assert_eq!(n, sbuf.len(), "unexpected length of received data");
2273 assert_eq!(ppi, PayloadProtocolIdentifier::Binary, "unexpected ppi");
2274 assert_eq!(
2275 u32::from_be_bytes([buf[0], buf[1], buf[2], buf[3]]),
2276 i,
2277 "unexpected received data"
2278 );
2279 }
2280 
2281 close_association_pair(&mut pair, client_ch, server_ch, si);
2282 
2283 Ok(())
2284}
2285 
2286/*
2287TODO: The following tests will be moved to sctp-async tests:
2288struct FakeEchoConn {
2289 wr_tx: Mutex<mpsc::Sender<Vec<u8>>>,
2290 rd_rx: Mutex<mpsc::Receiver<Vec<u8>>>,
2291 bytes_sent: AtomicUsize,
2292 bytes_received: AtomicUsize,
2293}
2294
2295impl FakeEchoConn {
2296 fn new() -> impl Conn + AsAny {
2297 let (wr_tx, rd_rx) = mpsc::channel(1);
2298 FakeEchoConn {
2299 wr_tx: Mutex::new(wr_tx),
2300 rd_rx: Mutex::new(rd_rx),
2301 bytes_sent: AtomicUsize::new(0),
2302 bytes_received: AtomicUsize::new(0),
2303 }
2304 }
2305}
2306
2307trait AsAny {
2308 fn as_any(&self) -> &(dyn std::any::Any + Send + Sync);
2309}
2310
2311impl AsAny for FakeEchoConn {
2312 fn as_any(&self) -> &(dyn std::any::Any + Send + Sync) {
2313 self
2314 }
2315}
2316
2317type UResult<T> = std::result::Result<T, util::Error>;
2318
2319#[async_trait]
2320impl Conn for FakeEchoConn {
2321 fn connect(&self, _addr: SocketAddr) -> UResult<()> {
2322 Err(io::Error::new(io::ErrorKind::Other, "Not applicable").into())
2323 }
2324
2325 fn recv(&self, b: &mut [u8]) -> UResult<usize> {
2326 let mut rd_rx = self.rd_rx.lock().await;
2327 let v = match rd_rx.recv().await {
2328 Some(v) => v,
2329 None => {
2330 return Err(io::Error::new(io::ErrorKind::UnexpectedEof, "Unexpected EOF").into())
2331 }
2332 };
2333 let l = std::cmp::min(v.len(), b.len());
2334 b[..l].copy_from_slice(&v[..l]);
2335 self.bytes_received.fetch_add(l, Ordering::SeqCst);
2336 Ok(l)
2337 }
2338
2339 fn recv_from(&self, _buf: &mut [u8]) -> UResult<(usize, SocketAddr)> {
2340 Err(io::Error::new(io::ErrorKind::Other, "Not applicable").into())
2341 }
2342
2343 fn send(&self, b: &[u8]) -> UResult<usize> {
2344 let wr_tx = self.wr_tx.lock().await;
2345 match wr_tx.send(b.to_vec()).await {
2346 Ok(_) => {}
2347 Err(err) => return Err(io::Error::new(io::ErrorKind::Other, err.to_string()).into()),
2348 };
2349 self.bytes_sent.fetch_add(b.len(), Ordering::SeqCst);
2350 Ok(b.len())
2351 }
2352
2353 fn send_to(&self, _buf: &[u8], _target: SocketAddr) -> UResult<usize> {
2354 Err(io::Error::new(io::ErrorKind::Other, "Not applicable").into())
2355 }
2356
2357 fn local_addr(&self) -> UResult<SocketAddr> {
2358 Err(io::Error::new(io::ErrorKind::AddrNotAvailable, "Addr Not Available").into())
2359 }
2360
2361 fn remote_addr(&self) -> Option<SocketAddr> {
2362 None
2363 }
2364
2365 fn close(&self) -> UResult<()> {
2366 Ok(())
2367 }
2368}
2369
2370//use std::io::Write;
2371
2372#[test]
2373fn test_stats() -> Result<()> {
2374 /*env_logger::Builder::new()
2375 .format(|buf, record| {
2376 writeln!(
2377 buf,
2378 "{}:{} [{}] {} - {}",
2379 record.file().unwrap_or("unknown"),
2380 record.line().unwrap_or(0),
2381 record.level(),
2382 chrono::Local::now().format("%H:%M:%S.%6f"),
2383 record.args()
2384 )
2385 })
2386 .filter(None, log::LevelFilter::Trace)
2387 .init();*/
2388
2389 let conn = Arc::new(FakeEchoConn::new());
2390 let a = Association::client(Config {
2391 net_conn: Arc::clone(&conn) as Arc<dyn Conn + Send + Sync>,
2392 max_receive_buffer_size: 0,
2393 max_message_size: 0,
2394 name: "client".to_owned(),
2395 })
2396 .await?;
2397
2398 if let Some(conn) = conn.as_any().downcast_ref::<FakeEchoConn>() {
2399 assert_eq!(
2400 conn.bytes_received.load(Ordering::SeqCst),
2401 a.bytes_received()
2402 );
2403 assert_eq!(conn.bytes_sent.load(Ordering::SeqCst), a.bytes_sent());
2404 } else {
2405 assert!(false, "must be FakeEchoConn");
2406 }
2407
2408 Ok(())
2409}
2410
2411fn create_assocs() -> Result<(Association, Association)> {
2412 let addr1 = SocketAddr::from_str("0.0.0.0:0").unwrap();
2413 let addr2 = SocketAddr::from_str("0.0.0.0:0").unwrap();
2414
2415 let udp1 = UdpSocket::bind(addr1).await.unwrap();
2416 let udp2 = UdpSocket::bind(addr2).await.unwrap();
2417
2418 udp1.connect(udp2.local_addr().unwrap()).await.unwrap();
2419 udp2.connect(udp1.local_addr().unwrap()).await.unwrap();
2420
2421 let (a1chan_tx, mut a1chan_rx) = mpsc::channel(1);
2422 let (a2chan_tx, mut a2chan_rx) = mpsc::channel(1);
2423
2424 tokio::spawn(async move {
2425 let a = Association::client(Config {
2426 net_conn: Arc::new(udp1),
2427 max_receive_buffer_size: 0,
2428 max_message_size: 0,
2429 name: "client".to_owned(),
2430 })
2431 .await?;
2432
2433 let _ = a1chan_tx.send(a).await;
2434
2435 Result::<()>::Ok(())
2436 });
2437
2438 tokio::spawn(async move {
2439 let a = Association::server(Config {
2440 net_conn: Arc::new(udp2),
2441 max_receive_buffer_size: 0,
2442 max_message_size: 0,
2443 name: "server".to_owned(),
2444 })
2445 .await?;
2446
2447 let _ = a2chan_tx.send(a).await;
2448
2449 Result::<()>::Ok(())
2450 });
2451
2452 let timer1 = tokio::time::sleep(Duration::from_secs(1));
2453 tokio::pin!(timer1);
2454 let a1 = tokio::select! {
2455 _ = timer1.as_mut() =>{
2456 assert!(false,"timed out waiting for a1");
2457 return Err(Error::Other("timed out waiting for a1".to_owned()).into());
2458 },
2459 a1 = a1chan_rx.recv() => {
2460 a1.unwrap()
2461 }
2462 };
2463
2464 let timer2 = tokio::time::sleep(Duration::from_secs(1));
2465 tokio::pin!(timer2);
2466 let a2 = tokio::select! {
2467 _ = timer2.as_mut() =>{
2468 assert!(false,"timed out waiting for a2");
2469 return Err(Error::Other("timed out waiting for a2".to_owned()).into());
2470 },
2471 a2 = a2chan_rx.recv() => {
2472 a2.unwrap()
2473 }
2474 };
2475
2476 Ok((a1, a2))
2477}
2478
2479//use std::io::Write;
2480//TODO: remove this conditional test
2481#[cfg(not(target_os = "windows"))]
2482#[test]
2483fn test_association_shutdown() -> Result<()> {
2484 /*env_logger::Builder::new()
2485 .format(|buf, record| {
2486 writeln!(
2487 buf,
2488 "{}:{} [{}] {} - {}",
2489 record.file().unwrap_or("unknown"),
2490 record.line().unwrap_or(0),
2491 record.level(),
2492 chrono::Local::now().format("%H:%M:%S.%6f"),
2493 record.args()
2494 )
2495 })
2496 .filter(None, log::LevelFilter::Trace)
2497 .init();*/
2498
2499 let (a1, a2) = create_assocs().await?;
2500
2501 let s11 = a1.open_stream(1, PayloadProtocolIdentifier::String).await?;
2502 let s21 = a2.open_stream(1, PayloadProtocolIdentifier::String).await?;
2503
2504 let test_data = Bytes::from_static(b"test");
2505
2506 let n = s11.write(&test_data).await?;
2507 assert_eq!(test_data.len(), n);
2508
2509 let mut buf = vec![0u8; test_data.len()];
2510 let n = s21.read(&mut buf).await?;
2511 assert_eq!(test_data.len(), n);
2512 assert_eq!(&test_data, &buf[0..n]);
2513
2514 if let Ok(result) = tokio::time::timeout(Duration::from_secs(1), a1.shutdown()).await {
2515 assert!(result.is_ok(), "shutdown should be ok");
2516 } else {
2517 assert!(false, "shutdown timeout");
2518 }
2519
2520 {
2521 let mut close_loop_ch_rx = a2.close_loop_ch_rx.lock().await;
2522
2523 // Wait for close read loop channels to prevent flaky tests.
2524 let timer2 = tokio::time::sleep(Duration::from_secs(1));
2525 tokio::pin!(timer2);
2526 tokio::select! {
2527 _ = timer2.as_mut() =>{
2528 assert!(false,"timed out waiting for a2 read loop to close");
2529 },
2530 _ = close_loop_ch_rx.recv() => {
2531 log::debug!("recv a2.close_loop_ch_rx");
2532 }
2533 };
2534 }
2535 Ok(())
2536}
2537
2538
2539fn test_association_shutdown_during_write() -> Result<()> {
2540 //let _guard = subscribe();
2541
2542 let (a1, a2) = create_assocs().await?;
2543
2544 let s11 = a1.open_stream(1, PayloadProtocolIdentifier::String).await?;
2545 let s21 = a2.open_stream(1, PayloadProtocolIdentifier::String).await?;
2546
2547 let (writing_done_tx, mut writing_done_rx) = mpsc::channel::<()>(1);
2548 let ss21 = Arc::clone(&s21);
2549 tokio::spawn(async move {
2550 let mut i = 0;
2551 while ss21.write(&Bytes::from(vec![i])).await.is_ok() {
2552 if i == 255 {
2553 i = 0;
2554 } else {
2555 i += 1;
2556 }
2557
2558 if i % 100 == 0 {
2559 tokio::time::sleep(Duration::from_millis(20)).await;
2560 }
2561 }
2562
2563 drop(writing_done_tx);
2564 });
2565
2566 let test_data = Bytes::from_static(b"test");
2567
2568 let n = s11.write(&test_data).await?;
2569 assert_eq!(test_data.len(), n);
2570
2571 let mut buf = vec![0u8; test_data.len()];
2572 let n = s21.read(&mut buf).await?;
2573 assert_eq!(test_data.len(), n);
2574 assert_eq!(&test_data, &buf[0..n]);
2575
2576 {
2577 let mut close_loop_ch_rx = a1.close_loop_ch_rx.lock().await;
2578 tokio::select! {
2579 res = tokio::time::timeout(Duration::from_secs(1), a1.shutdown()) => {
2580 if let Ok(result) = res {
2581 assert!(result.is_ok(), "shutdown should be ok");
2582 } else {
2583 assert!(false, "shutdown timeout");
2584 }
2585 }
2586 _ = writing_done_rx.recv() => {
2587 log::debug!("writing_done_rx");
2588 let result = close_loop_ch_rx.recv().await;
2589 log::debug!("a1.close_loop_ch_rx.recv: {:?}", result);
2590 },
2591 };
2592 }
2593
2594 {
2595 let mut close_loop_ch_rx = a2.close_loop_ch_rx.lock().await;
2596 // Wait for close read loop channels to prevent flaky tests.
2597 let timer2 = tokio::time::sleep(Duration::from_secs(1));
2598 tokio::pin!(timer2);
2599 tokio::select! {
2600 _ = timer2.as_mut() =>{
2601 assert!(false,"timed out waiting for a2 read loop to close");
2602 },
2603 _ = close_loop_ch_rx.recv() => {
2604 log::debug!("recv a2.close_loop_ch_rx");
2605 }
2606 };
2607 }
2608
2609 Ok(())
2610}*/
2611 
2612#[test]
2613fn test_assoc_reset_duplicate_reconfig_request() -> Result<()> {
2614 let si: u16 = 1;
2615 
2616 let (mut pair, client_ch, server_ch) = create_association_pair(AckMode::NoDelay, 0)?;
2617 establish_session_pair(&mut pair, client_ch, server_ch, si)?;
2618 
2619 // Client initiates reset of stream 1
2620 pair.client_stream(client_ch, si)?.stop()?;
2621 
2622 // Drive client to generate the reconfig packet, which lands in server's inbound
2623 pair.drive_client();
2624 
2625 // Capture the raw packet bytes before the server processes them.
2626 // These contain the ChunkReconfig with the ParamOutgoingResetRequest.
2627 let captured_packets: Vec<_> = pair.server.inbound.iter().cloned().collect();
2628 
2629 // Let the entire reset flow complete
2630 pair.drive();
2631 
2632 // Verify stream 1 is gone on the server
2633 assert!(
2634 pair.server_stream(server_ch, si).is_err(),
2635 "stream 1 should be removed after reset"
2636 );
2637 
2638 // Server opens a new stream 1
2639 let _ = pair
2640 .server_conn_mut(server_ch)
2641 .open_stream(si, PayloadProtocolIdentifier::Binary)?;
2642 assert!(
2643 pair.server_stream(server_ch, si).is_ok(),
2644 "new stream 1 should exist"
2645 );
2646 
2647 // Inject the captured reconfig packets again (simulating retransmission)
2648 for packet in captured_packets {
2649 pair.server.inbound.push_back(packet);
2650 }
2651 
2652 // Process the injected packets
2653 pair.drive();
2654 
2655 // The new stream 1 must NOT be destroyed by the duplicate reconfig request
2656 assert!(
2657 pair.server_stream(server_ch, si).is_ok(),
2658 "new stream 1 should NOT be destroyed by duplicate reconfig request"
2659 );
2660 
2661 Ok(())
2662}
2663 
2664/// Verify that a retransmission of a RE-CONFIG request that was initially
2665/// InProgress and then completed via TSN advance (not via handle_reconfig_param)
2666/// does not destroy a reused stream. This exercises the case where
2667/// max_completed_reconfig_rsn must be updated in the TSN-advance loop.
2668#[test]
2669fn test_assoc_reset_inprogress_completed_via_tsn_advance_then_retransmit() -> Result<()> {
2670 let si: u16 = 1;
2671 
2672 let (mut pair, client_ch, server_ch) = create_association_pair(AckMode::NoDelay, 0)?;
2673 establish_session_pair(&mut pair, client_ch, server_ch, si)?;
2674 
2675 // Write data so my_next_tsn advances on the client; the RECONFIG
2676 // sender_last_tsn will include this TSN.
2677 let _ = pair.client_stream(client_ch, si)?.write_sctp(
2678 &Bytes::from_static(b"payload"),
2679 PayloadProtocolIdentifier::Binary,
2680 )?;
2681 
2682 // Drive client to generate the DATA packet(s).
2683 pair.drive_client();
2684 
2685 // Withhold the DATA packets from the server (simulate loss).
2686 let withheld: Vec<_> = pair.server.inbound.drain(..).collect();
2687 
2688 // Client initiates reset of stream 1.
2689 pair.client_stream(client_ch, si)?.stop()?;
2690 
2691 // Drive client to generate EOS data + RECONFIG packets.
2692 pair.drive_client();
2693 
2694 // Separate RECONFIG-bearing packets from DATA-only packets.
2695 let (reconfig_packets, data_packets): (Vec<_>, Vec<_>) =
2696 pair.server.inbound.drain(..).partition(|(_, _, raw)| {
2697 let mut offset = 12usize;
2698 while offset + 4 <= raw.len() {
2699 if raw[offset] == 130 {
2700 return true;
2701 }
2702 let chunk_len = u16::from_be_bytes([raw[offset + 2], raw[offset + 3]]) as usize;
2703 if chunk_len < 4 {
2704 break;
2705 }
2706 offset += (chunk_len + 3) & !3;
2707 }
2708 false
2709 });
2710 
2711 assert!(!reconfig_packets.is_empty(), "expected a RECONFIG packet");
2712 
2713 // Step 1: Deliver only the RECONFIG. Server hasn't seen the withheld
2714 // DATA so peer_last_tsn < sender_last_tsn → InProgress.
2715 for pkt in &reconfig_packets {
2716 pair.server.inbound.push_back(pkt.clone());
2717 }
2718 pair.drive_server();
2719 
2720 assert!(
2721 pair.server_stream(server_ch, si).is_ok(),
2722 "stream should survive InProgress reset"
2723 );
2724 
2725 // Step 2: Now deliver the withheld DATA + the other data packets.
2726 // The TSN-advance loop in handle_data will complete the pending
2727 // reconfig request (removing it from reconfig_requests).
2728 // Crucially, this does NOT go through handle_reconfig_param, so
2729 // max_completed_reconfig_rsn is only updated if the TSN-advance
2730 // path does it.
2731 for pkt in withheld {
2732 pair.server.inbound.push_back(pkt);
2733 }
2734 for pkt in data_packets {
2735 pair.server.inbound.push_back(pkt);
2736 }
2737 pair.drive();
2738 
2739 // The boundary DATA must remain readable even though the reset completed.
2740 // Draining it retires the old stream generation.
2741 let boundary_data = pair
2742 .server_stream(server_ch, si)?
2743 .read()?
2744 .expect("boundary DATA should survive the reset");
2745 assert_eq!(boundary_data.len(), b"payload".len());
2746 
2747 assert!(
2748 pair.server_stream(server_ch, si).is_err(),
2749 "stream should retire after its boundary DATA is drained"
2750 );
2751 
2752 // Step 3: Reopen stream 1 with the same ID.
2753 let _ = pair
2754 .server_conn_mut(server_ch)
2755 .open_stream(si, PayloadProtocolIdentifier::Binary)?;
2756 assert!(
2757 pair.server_stream(server_ch, si).is_ok(),
2758 "new stream 1 should exist"
2759 );
2760 
2761 // Step 4: Replay the original RECONFIG (simulating a late retransmission).
2762 // If the watermark was not updated during the TSN-advance completion,
2763 // this will bypass the dedup guard and destroy the new stream.
2764 for pkt in reconfig_packets {
2765 pair.server.inbound.push_back(pkt);
2766 }
2767 pair.drive();
2768 
2769 // The new stream 1 must survive.
2770 assert!(
2771 pair.server_stream(server_ch, si).is_ok(),
2772 "new stream 1 should NOT be destroyed by retransmitted reconfig \
2773 after InProgress completion via TSN advance"
2774 );
2775 
2776 Ok(())
2777}
2778 
2779/// Verify that a retransmission of an InProgress RE-CONFIG request is
2780/// re-evaluated (not falsely deduplicated) so the reset completes once
2781/// the outstanding TSN is finally received.
2782#[test]
2783fn test_assoc_reset_inprogress_reconfig_retransmission() -> Result<()> {
2784 let si: u16 = 1;
2785 
2786 let (mut pair, client_ch, server_ch) = create_association_pair(AckMode::NoDelay, 0)?;
2787 establish_session_pair(&mut pair, client_ch, server_ch, si)?;
2788 
2789 // Write data so my_next_tsn advances on the client; the RECONFIG
2790 // sender_last_tsn will include this TSN.
2791 let _ = pair.client_stream(client_ch, si)?.write_sctp(
2792 &Bytes::from_static(b"payload"),
2793 PayloadProtocolIdentifier::Binary,
2794 )?;
2795 
2796 // Drive client to generate the DATA packet(s).
2797 pair.drive_client();
2798 
2799 // Withhold the DATA packets from the server (simulate loss).
2800 let withheld: Vec<_> = pair.server.inbound.drain(..).collect();
2801 
2802 // Client initiates reset of stream 1.
2803 pair.client_stream(client_ch, si)?.stop()?;
2804 
2805 // Drive client to generate EOS data + RECONFIG packets.
2806 pair.drive_client();
2807 
2808 // Separate RECONFIG-bearing packets from DATA-only packets.
2809 // Chunk type is the first byte after the 12-byte common header.
2810 let (reconfig_packets, data_packets): (Vec<_>, Vec<_>) =
2811 pair.server.inbound.drain(..).partition(|(_, _, raw)| {
2812 // Scan all chunks in the packet for CT_RECONFIG (130).
2813 let mut offset = 12usize;
2814 while offset + 4 <= raw.len() {
2815 if raw[offset] == 130 {
2816 return true;
2817 }
2818 let chunk_len = u16::from_be_bytes([raw[offset + 2], raw[offset + 3]]) as usize;
2819 if chunk_len < 4 {
2820 break;
2821 }
2822 offset += (chunk_len + 3) & !3; // pad to 4-byte boundary
2823 }
2824 false
2825 });
2826 
2827 assert!(!reconfig_packets.is_empty(), "expected a RECONFIG packet");
2828 
2829 // Deliver only the RECONFIG packets. The server hasn't seen the withheld
2830 // DATA so peer_last_tsn < sender_last_tsn → InProgress.
2831 for pkt in &reconfig_packets {
2832 pair.server.inbound.push_back(pkt.clone());
2833 }
2834 pair.drive_server();
2835 
2836 // Stream should still exist (reset is InProgress, not completed).
2837 assert!(
2838 pair.server_stream(server_ch, si).is_ok(),
2839 "stream should survive InProgress reset"
2840 );
2841 
2842 // Inject the RECONFIG again (simulating retransmission) together with
2843 // the withheld DATA so the TSN can finally advance.
2844 for pkt in reconfig_packets {
2845 pair.server.inbound.push_back(pkt);
2846 }
2847 for pkt in withheld {
2848 pair.server.inbound.push_back(pkt);
2849 }
2850 for pkt in data_packets {
2851 pair.server.inbound.push_back(pkt);
2852 }
2853 
2854 // Let everything settle.
2855 pair.drive();
2856 
2857 // The boundary DATA must remain readable even though the reset completed.
2858 // Draining it retires the old stream generation.
2859 let boundary_data = pair
2860 .server_stream(server_ch, si)?
2861 .read()?
2862 .expect("boundary DATA should survive the reset");
2863 assert_eq!(boundary_data.len(), b"payload".len());
2864 
2865 assert!(
2866 pair.server_stream(server_ch, si).is_err(),
2867 "stream should retire after its boundary DATA is drained"
2868 );
2869 
2870 Ok(())
2871}
2872 
2873#[test]
2874fn test_assoc_open_stream_rejects_pending_reset_id() -> Result<()> {
2875 let si: u16 = 1;
2876 
2877 let (mut pair, client_ch, server_ch) = create_association_pair(AckMode::NoDelay, 0)?;
2878 establish_session_pair(&mut pair, client_ch, server_ch, si)?;
2879 
2880 // SERVER initiates reset of stream 1. When the client processes the server's
2881 // incoming RE-CONFIG, it removes stream 1 from self.streams and generates its
2882 // own outgoing RE-CONFIG (stored in self.reconfigs) awaiting acknowledgment.
2883 pair.server_stream(server_ch, si)?.stop()?;
2884 
2885 // Drive server to generate and deliver the RE-CONFIG to the client's inbound
2886 pair.drive_server();
2887 
2888 // Drive client to process the server's RE-CONFIG. This removes stream 1 from
2889 // the client's stream table and creates a pending outgoing RE-CONFIG.
2890 pair.drive_client();
2891 
2892 // Client's stream 1 is removed, but the outgoing RE-CONFIG is still pending.
2893 // open_stream should reject this stream ID.
2894 assert!(
2895 pair.client_stream(client_ch, si).is_err(),
2896 "stream 1 should be removed from client"
2897 );
2898 match pair
2899 .client_conn_mut(client_ch)
2900 .open_stream(si, PayloadProtocolIdentifier::Binary)
2901 {
2902 Err(Error::ErrStreamResetPending) => {}
2903 other => panic!("expected ErrStreamResetPending, got {:?}", other.err()),
2904 }
2905 
2906 // Let the full RE-CONFIG exchange complete
2907 pair.drive();
2908 
2909 // Now the RE-CONFIG is acknowledged, stream 1 should be available for reuse
2910 let _ = pair
2911 .client_conn_mut(client_ch)
2912 .open_stream(si, PayloadProtocolIdentifier::Binary)?;
2913 assert!(
2914 pair.client_stream(client_ch, si).is_ok(),
2915 "stream 1 should be available after RE-CONFIG is acknowledged"
2916 );
2917 
2918 Ok(())
2919}
2920 
2921#[test]
2922fn test_assoc_reconfig_failure_keeps_stream_quarantined() -> Result<()> {
2923 let si: u16 = 1;
2924 
2925 let (mut pair, client_ch, server_ch) = create_association_pair(AckMode::NoDelay, 0)?;
2926 establish_session_pair(&mut pair, client_ch, server_ch, si)?;
2927 
2928 // SERVER initiates reset of stream 1. When the client processes the incoming
2929 // RE-CONFIG it removes stream 1 from self.streams and generates its own
2930 // outgoing RE-CONFIG (stored in self.reconfigs).
2931 pair.server_stream(server_ch, si)?.stop()?;
2932 
2933 // Drive server to send the RE-CONFIG to client
2934 pair.drive_server();
2935 
2936 // Drive client to process the incoming RE-CONFIG
2937 pair.drive_client();
2938 
2939 // Drop any packets the client sent so the server never ACKs the client's
2940 // outgoing RE-CONFIG.
2941 pair.server.inbound.clear();
2942 
2943 // Verify open_stream is blocked by the pending outgoing RE-CONFIG
2944 match pair
2945 .client_conn_mut(client_ch)
2946 .open_stream(si, PayloadProtocolIdentifier::Binary)
2947 {
2948 Err(Error::ErrStreamResetPending) => {}
2949 Err(e) => panic!("expected ErrStreamResetPending, got Err({:?})", e),
2950 Ok(_) => panic!("expected ErrStreamResetPending, got Ok"),
2951 }
2952 
2953 // Advance time through MAX_INIT_RETRANS (8) retransmissions + 1 to trigger
2954 // failure. Each iteration jumps 61 seconds (> RTO_MAX of 60s) to guarantee
2955 // the Reconfig timer fires every time. We discard all outbound packets so
2956 // the Reconfig is never acknowledged.
2957 for _ in 0..10 {
2958 pair.time += Duration::from_secs(61);
2959 pair.drive_client();
2960 pair.server.inbound.clear();
2961 }
2962 
2963 assert!(
2964 core::iter::from_fn(|| pair.client_conn_mut(client_ch).poll()).any(|event| matches!(
2965 event,
2966 Event::Stream(StreamEvent::ResetFailed {
2967 id,
2968 reason: StreamResetError::Failed,
2969 }) if id == si
2970 )),
2971 "reset retransmission exhaustion should emit a terminal failure event"
2972 );
2973 
2974 // The in-flight request is abandoned, but without a terminal success the
2975 // stream ID remains quarantined and must not be reused.
2976 assert!(matches!(
2977 pair.client_conn_mut(client_ch)
2978 .open_stream(si, PayloadProtocolIdentifier::Binary),
2979 Err(Error::ErrStreamResetPending)
2980 ));
2981 
2982 Ok(())
2983}
2984 
2985#[test]
2986fn test_snap_connect_established_and_transmit_uses_peer_verification_tag() {
2987 let now = Instant::now();
2988 
2989 let local_transport = TransportConfig::default();
2990 let remote_transport = TransportConfig::default();
2991 
2992 let local_init_bytes = generate_snap_token(&local_transport).expect("generate local init");
2993 let remote_init_bytes = generate_snap_token(&remote_transport).expect("generate remote init");
2994 
2995 // Parse remote to check verification tag later
2996 let remote_init =
2997 ChunkInit::unmarshal(&remote_init_bytes).expect("unmarshal remote INIT chunk");
2998 
2999 let mut endpoint = Endpoint::new(Arc::new(EndpointConfig::default()), None);
3000 let remote_addr: SocketAddr = "127.0.0.1:5000".parse().unwrap();
3001 
3002 let client_config = ClientConfig::new().with_snap(local_init_bytes, remote_init_bytes);
3003 let (_ch, mut assoc) = endpoint
3004 .connect(client_config, remote_addr)
3005 .expect("SNAP connect should succeed");
3006 
3007 assert_eq!(assoc.state(), AssociationState::Established);
3008 assert_matches!(assoc.poll(), Some(Event::Connected));
3009 
3010 // Ensure outbound packets use the peer's initiate_tag as the verification tag.
3011 let mut stream = assoc
3012 .open_stream(1, PayloadProtocolIdentifier::Binary)
3013 .expect("open stream");
3014 let msg = Bytes::from_static(b"hello");
3015 stream
3016 .write_sctp(&msg, PayloadProtocolIdentifier::Binary)
3017 .expect("write_sctp");
3018 
3019 let transmit = assoc
3020 .poll_transmit(now)
3021 .expect("expected at least one outbound datagram");
3022 let Payload::RawEncode(datagrams) = transmit.payload else {
3023 panic!("expected RawEncode transmit");
3024 };
3025 assert!(
3026 !datagrams.is_empty(),
3027 "expected at least one outbound packet"
3028 );
3029 
3030 let pkt = Packet::unmarshal(&datagrams[0]).expect("unmarshal outbound packet");
3031 assert_eq!(
3032 pkt.common_header.verification_tag, remote_init.initiate_tag,
3033 "outbound packet should use peer initiate_tag as verification_tag"
3034 );
3035}
3036 
3037#[test]
3038fn test_server_retransmitted_init_routes_to_existing_association() {
3039 let now = Instant::now();
3040 
3041 let mut endpoint = Endpoint::new(
3042 Arc::new(EndpointConfig::default()),
3043 Some(Arc::new(ServerConfig::default())),
3044 );
3045 
3046 let remote: SocketAddr = "127.0.0.1:5000".parse().unwrap();
3047 let init_tag: u32 = 0x1122_3344;
3048 
3049 let init = ChunkInit {
3050 is_ack: false,
3051 initiate_tag: init_tag,
3052 ..Default::default()
3053 };
3054 
3055 let pkt = Packet {
3056 common_header: CommonHeader {
3057 source_port: 5000,
3058 destination_port: 5000,
3059 verification_tag: 0,
3060 },
3061 chunks: vec![Box::new(init)],
3062 };
3063 
3064 let bytes = pkt.marshal().expect("marshal INIT packet");
3065 
3066 let (ch1, ev1) = endpoint
3067 .handle(now, remote, None, None, bytes.clone())
3068 .expect("first INIT should be handled");
3069 assert!(
3070 matches!(ev1, DatagramEvent::NewAssociation(_)),
3071 "first INIT should create a new association"
3072 );
3073 
3074 // Same INIT retransmitted: should route to the same association handle.
3075 let (ch2, ev2) = endpoint
3076 .handle(now, remote, None, None, bytes)
3077 .expect("retransmitted INIT should be handled");
3078 assert_eq!(
3079 ch2, ch1,
3080 "retransmitted INIT should map to same association"
3081 );
3082 assert!(
3083 matches!(ev2, DatagramEvent::AssociationEvent(_)),
3084 "retransmitted INIT should be routed to existing association"
3085 );
3086}
3087 
3088#[test]
3089fn test_snap_rejects_invalid_remote_bytes() {
3090 let mut endpoint = Endpoint::new(Arc::new(EndpointConfig::default()), None);
3091 let remote_addr: SocketAddr = "127.0.0.1:5000".parse().unwrap();
3092 
3093 let local_init = generate_snap_token(&TransportConfig::default()).expect("local init");
3094 let client_config = ClientConfig::new().with_snap(local_init, Bytes::from_static(b"nope"));
3095 let res = endpoint.connect(client_config, remote_addr);
3096 assert!(
3097 matches!(
3098 res,
3099 Err(ConnectError::Snap(SnapError::ParseFailed {
3100 side: SnapSide::Remote,
3101 ..
3102 }))
3103 ),
3104 "invalid remote bytes should fail with ParseFailed"
3105 );
3106}
3107 
3108#[test]
3109fn test_snap_rejects_invalid_local_bytes() {
3110 let mut endpoint = Endpoint::new(Arc::new(EndpointConfig::default()), None);
3111 let remote_addr: SocketAddr = "127.0.0.1:5000".parse().unwrap();
3112 
3113 let remote_init = generate_snap_token(&TransportConfig::default()).expect("remote init");
3114 let client_config = ClientConfig::new().with_snap(Bytes::from_static(b"nope"), remote_init);
3115 let res = endpoint.connect(client_config, remote_addr);
3116 assert!(
3117 matches!(
3118 res,
3119 Err(ConnectError::Snap(SnapError::ParseFailed {
3120 side: SnapSide::Local,
3121 ..
3122 }))
3123 ),
3124 "invalid local bytes should fail with ParseFailed"
3125 );
3126}
3127 
3128#[test]
3129fn test_snap_rejects_remote_init_ack() {
3130 let mut endpoint = Endpoint::new(Arc::new(EndpointConfig::default()), None);
3131 let remote_addr: SocketAddr = "127.0.0.1:5000".parse().unwrap();
3132 
3133 let local_init = generate_snap_token(&TransportConfig::default()).expect("local init");
3134 let remote_init_ack = ChunkInit {
3135 is_ack: true,
3136 initiate_tag: 1,
3137 ..Default::default()
3138 };
3139 let remote_bytes = remote_init_ack.marshal().expect("marshal remote INIT-ACK");
3140 
3141 let client_config = ClientConfig::new().with_snap(local_init, remote_bytes);
3142 let res = endpoint.connect(client_config, remote_addr);
3143 assert!(
3144 matches!(
3145 res,
3146 Err(ConnectError::Snap(SnapError::InvalidInitAck {
3147 side: SnapSide::Remote
3148 }))
3149 ),
3150 "remote INIT-ACK should be rejected"
3151 );
3152}
3153 
3154#[test]
3155fn test_snap_rejects_remote_init_with_zero_initiate_tag() {
3156 let mut endpoint = Endpoint::new(Arc::new(EndpointConfig::default()), None);
3157 let remote_addr: SocketAddr = "127.0.0.1:5000".parse().unwrap();
3158 
3159 let local_init = generate_snap_token(&TransportConfig::default()).expect("local init");
3160 let remote_init = ChunkInit {
3161 is_ack: false,
3162 initiate_tag: 0,
3163 ..Default::default()
3164 };
3165 let remote_bytes = remote_init.marshal().expect("marshal remote INIT");
3166 
3167 let client_config = ClientConfig::new().with_snap(local_init, remote_bytes);
3168 let res = endpoint.connect(client_config, remote_addr);
3169 assert!(
3170 matches!(
3171 res,
3172 Err(ConnectError::Snap(SnapError::ZeroInitiateTag {
3173 side: SnapSide::Remote
3174 }))
3175 ),
3176 "remote INIT with initiate_tag=0 should be rejected"
3177 );
3178}
3179 
3180#[test]
3181fn test_snap_partial_config_falls_back_to_handshake() {
3182 let mut endpoint = Endpoint::new(Arc::new(EndpointConfig::default()), None);
3183 let remote_addr: SocketAddr = "127.0.0.1:5000".parse().unwrap();
3184 
3185 // Only local_sctp_init set, remote missing — should fall back to normal handshake
3186 let local_init = generate_snap_token(&TransportConfig::default()).expect("local init");
3187 let client_config = ClientConfig {
3188 local_sctp_init: Some(local_init),
3189 ..Default::default()
3190 };
3191 let (_ch, assoc) = endpoint
3192 .connect(client_config, remote_addr)
3193 .expect("partial SNAP config should fall back to handshake");
3194 assert_eq!(
3195 assoc.state(),
3196 AssociationState::CookieWait,
3197 "should use normal handshake when only local init provided"
3198 );
3199 
3200 // Only remote_sctp_init set, local missing — should fall back to normal handshake
3201 let remote_init = generate_snap_token(&TransportConfig::default()).expect("remote init");
3202 let client_config = ClientConfig {
3203 remote_sctp_init: Some(remote_init),
3204 ..Default::default()
3205 };
3206 let (_ch, assoc) = endpoint
3207 .connect(client_config, remote_addr)
3208 .expect("partial SNAP config should fall back to handshake");
3209 assert_eq!(
3210 assoc.state(),
3211 AssociationState::CookieWait,
3212 "should use normal handshake when only remote init provided"
3213 );
3214}
3215 
3216#[test]
3217fn test_snap_rejects_reused_local_init() {
3218 let mut endpoint = Endpoint::new(Arc::new(EndpointConfig::default()), None);
3219 let remote_addr: SocketAddr = "127.0.0.1:5000".parse().unwrap();
3220 
3221 // Generate a single local init — reusing it causes a collision
3222 let local_init = generate_snap_token(&TransportConfig::default()).expect("local init");
3223 
3224 let remote1 = generate_snap_token(&TransportConfig::default()).expect("remote init 1");
3225 let remote2 = generate_snap_token(&TransportConfig::default()).expect("remote init 2");
3226 
3227 // First connect succeeds
3228 let cfg1 = ClientConfig::new().with_snap(local_init.clone(), remote1);
3229 let res1 = endpoint.connect(cfg1, remote_addr);
3230 assert!(res1.is_ok(), "first SNAP connect should succeed");
3231 
3232 // Second connect with same local_init collides on local_aid
3233 let cfg2 = ClientConfig::new().with_snap(local_init, remote2);
3234 let res2 = endpoint.connect(cfg2, remote_addr);
3235 assert!(
3236 matches!(
3237 res2,
3238 Err(ConnectError::Snap(SnapError::AidCollision { .. }))
3239 ),
3240 "second SNAP connect should fail due to AID collision"
3241 );
3242}
3243 
3244#[test]
3245fn test_snap_rejects_local_init_ack() {
3246 let mut endpoint = Endpoint::new(Arc::new(EndpointConfig::default()), None);
3247 let remote_addr: SocketAddr = "127.0.0.1:5000".parse().unwrap();
3248 
3249 let remote_init = generate_snap_token(&TransportConfig::default()).expect("remote init");
3250 let local_init_ack = ChunkInit {
3251 is_ack: true,
3252 initiate_tag: 1,
3253 ..Default::default()
3254 };
3255 let local_bytes = local_init_ack.marshal().expect("marshal local INIT-ACK");
3256 
3257 let client_config = ClientConfig::new().with_snap(local_bytes, remote_init);
3258 let res = endpoint.connect(client_config, remote_addr);
3259 assert!(
3260 matches!(
3261 res,
3262 Err(ConnectError::Snap(SnapError::InvalidInitAck {
3263 side: SnapSide::Local
3264 }))
3265 ),
3266 "local INIT-ACK should be rejected"
3267 );
3268}
3269 
3270#[test]
3271fn test_snap_rejects_local_init_with_zero_initiate_tag() {
3272 let mut endpoint = Endpoint::new(Arc::new(EndpointConfig::default()), None);
3273 let remote_addr: SocketAddr = "127.0.0.1:5000".parse().unwrap();
3274 
3275 let remote_init = generate_snap_token(&TransportConfig::default()).expect("remote init");
3276 let local_init = ChunkInit {
3277 is_ack: false,
3278 initiate_tag: 0,
3279 ..Default::default()
3280 };
3281 let local_bytes = local_init.marshal().expect("marshal local INIT");
3282 
3283 let client_config = ClientConfig::new().with_snap(local_bytes, remote_init);
3284 let res = endpoint.connect(client_config, remote_addr);
3285 assert!(
3286 matches!(
3287 res,
3288 Err(ConnectError::Snap(SnapError::ZeroInitiateTag {
3289 side: SnapSide::Local
3290 }))
3291 ),
3292 "local INIT with initiate_tag=0 should be rejected"
3293 );
3294}
3295 
3296// ---------------------------------------------------------------------------
3297// SNAP hardening tests: size limits, truncated / empty / oversized input,
3298// wrong chunk types, RFC field validation, round-trip token consistency.
3299// ---------------------------------------------------------------------------
3300 
3301#[test]
3302fn test_snap_rejects_empty_bytes() {
3303 let mut endpoint = Endpoint::new(Arc::new(EndpointConfig::default()), None);
3304 let remote_addr: SocketAddr = "127.0.0.1:5000".parse().unwrap();
3305 
3306 let valid = generate_snap_token(&TransportConfig::default()).expect("valid init");
3307 
3308 // Empty remote
3309 let cfg = ClientConfig::new().with_snap(valid.clone(), Bytes::new());
3310 let res = endpoint.connect(cfg, remote_addr);
3311 assert!(
3312 matches!(
3313 res,
3314 Err(ConnectError::Snap(SnapError::ParseFailed {
3315 side: SnapSide::Remote,
3316 ..
3317 }))
3318 ),
3319 "empty remote bytes should fail parsing"
3320 );
3321 
3322 // Empty local
3323 let cfg = ClientConfig::new().with_snap(Bytes::new(), valid);
3324 let res = endpoint.connect(cfg, remote_addr);
3325 assert!(
3326 matches!(
3327 res,
3328 Err(ConnectError::Snap(SnapError::ParseFailed {
3329 side: SnapSide::Local,
3330 ..
3331 }))
3332 ),
3333 "empty local bytes should fail parsing"
3334 );
3335}
3336 
3337#[test]
3338fn test_snap_rejects_truncated_bytes() {
3339 let mut endpoint = Endpoint::new(Arc::new(EndpointConfig::default()), None);
3340 let remote_addr: SocketAddr = "127.0.0.1:5000".parse().unwrap();
3341 
3342 let valid = generate_snap_token(&TransportConfig::default()).expect("valid init");
3343 
3344 // Truncate to just 10 bytes — not enough for INIT header + fixed fields
3345 let truncated = valid.slice(..10.min(valid.len()));
3346 
3347 let cfg = ClientConfig::new().with_snap(valid.clone(), truncated.clone());
3348 let res = endpoint.connect(cfg, remote_addr);
3349 assert!(
3350 matches!(
3351 res,
3352 Err(ConnectError::Snap(SnapError::ParseFailed {
3353 side: SnapSide::Remote,
3354 ..
3355 }))
3356 ),
3357 "truncated remote bytes should fail parsing"
3358 );
3359 
3360 let cfg = ClientConfig::new().with_snap(truncated, valid);
3361 let res = endpoint.connect(cfg, remote_addr);
3362 assert!(
3363 matches!(
3364 res,
3365 Err(ConnectError::Snap(SnapError::ParseFailed {
3366 side: SnapSide::Local,
3367 ..
3368 }))
3369 ),
3370 "truncated local bytes should fail parsing"
3371 );
3372}
3373 
3374#[test]
3375fn test_snap_rejects_oversized_bytes() {
3376 use crate::config::MAX_SNAP_INIT_BYTES;
3377 
3378 let mut endpoint = Endpoint::new(Arc::new(EndpointConfig::default()), None);
3379 let remote_addr: SocketAddr = "127.0.0.1:5000".parse().unwrap();
3380 
3381 let valid = generate_snap_token(&TransportConfig::default()).expect("valid init");
3382 let oversized = Bytes::from(vec![0u8; MAX_SNAP_INIT_BYTES + 1]);
3383 
3384 // Oversized remote
3385 let cfg = ClientConfig::new().with_snap(valid.clone(), oversized.clone());
3386 let res = endpoint.connect(cfg, remote_addr);
3387 assert!(
3388 matches!(
3389 res,
3390 Err(ConnectError::Snap(SnapError::OversizedInit {
3391 side: SnapSide::Remote,
3392 ..
3393 }))
3394 ),
3395 "oversized remote bytes should be rejected"
3396 );
3397 
3398 // Oversized local
3399 let cfg = ClientConfig::new().with_snap(oversized, valid);
3400 let res = endpoint.connect(cfg, remote_addr);
3401 assert!(
3402 matches!(
3403 res,
3404 Err(ConnectError::Snap(SnapError::OversizedInit {
3405 side: SnapSide::Local,
3406 ..
3407 }))
3408 ),
3409 "oversized local bytes should be rejected"
3410 );
3411}
3412 
3413#[test]
3414fn test_snap_rejects_remote_init_zero_inbound_streams() {
3415 let mut endpoint = Endpoint::new(Arc::new(EndpointConfig::default()), None);
3416 let remote_addr: SocketAddr = "127.0.0.1:5000".parse().unwrap();
3417 
3418 let local_init = generate_snap_token(&TransportConfig::default()).expect("local init");
3419 let bad_remote = ChunkInit {
3420 is_ack: false,
3421 initiate_tag: 0xDEAD_BEEF,
3422 initial_tsn: 1,
3423 num_outbound_streams: 1,
3424 num_inbound_streams: 0, // RFC violation
3425 advertised_receiver_window_credit: 65535,
3426 ..Default::default()
3427 };
3428 let remote_bytes = bad_remote.marshal().expect("marshal bad remote INIT");
3429 
3430 let cfg = ClientConfig::new().with_snap(local_init, remote_bytes);
3431 let res = endpoint.connect(cfg, remote_addr);
3432 assert!(
3433 matches!(
3434 res,
3435 Err(ConnectError::Snap(SnapError::InvalidInit {
3436 side: SnapSide::Remote,
3437 ..
3438 }))
3439 ),
3440 "remote INIT with num_inbound_streams=0 should be rejected: got {:?}",
3441 res
3442 );
3443}
3444 
3445#[test]
3446fn test_snap_rejects_remote_init_zero_outbound_streams() {
3447 let mut endpoint = Endpoint::new(Arc::new(EndpointConfig::default()), None);
3448 let remote_addr: SocketAddr = "127.0.0.1:5000".parse().unwrap();
3449 
3450 let local_init = generate_snap_token(&TransportConfig::default()).expect("local init");
3451 let bad_remote = ChunkInit {
3452 is_ack: false,
3453 initiate_tag: 0xDEAD_BEEF,
3454 initial_tsn: 1,
3455 num_outbound_streams: 0, // RFC violation
3456 num_inbound_streams: 1,
3457 advertised_receiver_window_credit: 65535,
3458 ..Default::default()
3459 };
3460 let remote_bytes = bad_remote.marshal().expect("marshal bad remote INIT");
3461 
3462 let cfg = ClientConfig::new().with_snap(local_init, remote_bytes);
3463 let res = endpoint.connect(cfg, remote_addr);
3464 assert!(
3465 matches!(
3466 res,
3467 Err(ConnectError::Snap(SnapError::InvalidInit {
3468 side: SnapSide::Remote,
3469 ..
3470 }))
3471 ),
3472 "remote INIT with num_outbound_streams=0 should be rejected: got {:?}",
3473 res
3474 );
3475}
3476 
3477#[test]
3478fn test_snap_rejects_remote_init_small_arwnd() {
3479 let mut endpoint = Endpoint::new(Arc::new(EndpointConfig::default()), None);
3480 let remote_addr: SocketAddr = "127.0.0.1:5000".parse().unwrap();
3481 
3482 let local_init = generate_snap_token(&TransportConfig::default()).expect("local init");
3483 let bad_remote = ChunkInit {
3484 is_ack: false,
3485 initiate_tag: 0xDEAD_BEEF,
3486 initial_tsn: 1,
3487 num_outbound_streams: 1,
3488 num_inbound_streams: 1,
3489 advertised_receiver_window_credit: 100, // RFC says >= 1500
3490 ..Default::default()
3491 };
3492 let remote_bytes = bad_remote.marshal().expect("marshal bad remote INIT");
3493 
3494 let cfg = ClientConfig::new().with_snap(local_init, remote_bytes);
3495 let res = endpoint.connect(cfg, remote_addr);
3496 assert!(
3497 matches!(
3498 res,
3499 Err(ConnectError::Snap(SnapError::InvalidInit {
3500 side: SnapSide::Remote,
3501 ..
3502 }))
3503 ),
3504 "remote INIT with a_rwnd < 1500 should be rejected: got {:?}",
3505 res
3506 );
3507}
3508 
3509#[test]
3510fn test_snap_rejects_local_init_zero_inbound_streams() {
3511 let mut endpoint = Endpoint::new(Arc::new(EndpointConfig::default()), None);
3512 let remote_addr: SocketAddr = "127.0.0.1:5000".parse().unwrap();
3513 
3514 let remote_init = generate_snap_token(&TransportConfig::default()).expect("remote init");
3515 let bad_local = ChunkInit {
3516 is_ack: false,
3517 initiate_tag: 0xCAFE_BABE,
3518 initial_tsn: 1,
3519 num_outbound_streams: 1,
3520 num_inbound_streams: 0,
3521 advertised_receiver_window_credit: 65535,
3522 ..Default::default()
3523 };
3524 let local_bytes = bad_local.marshal().expect("marshal bad local INIT");
3525 
3526 let cfg = ClientConfig::new().with_snap(local_bytes, remote_init);
3527 let res = endpoint.connect(cfg, remote_addr);
3528 assert!(
3529 matches!(
3530 res,
3531 Err(ConnectError::Snap(SnapError::InvalidInit {
3532 side: SnapSide::Local,
3533 ..
3534 }))
3535 ),
3536 "local INIT with num_inbound_streams=0 should be rejected: got {:?}",
3537 res
3538 );
3539}
3540 
3541#[test]
3542fn test_snap_rejects_local_init_small_arwnd() {
3543 let mut endpoint = Endpoint::new(Arc::new(EndpointConfig::default()), None);
3544 let remote_addr: SocketAddr = "127.0.0.1:5000".parse().unwrap();
3545 
3546 let remote_init = generate_snap_token(&TransportConfig::default()).expect("remote init");
3547 let bad_local = ChunkInit {
3548 is_ack: false,
3549 initiate_tag: 0xCAFE_BABE,
3550 initial_tsn: 1,
3551 num_outbound_streams: 1,
3552 num_inbound_streams: 1,
3553 advertised_receiver_window_credit: 0,
3554 ..Default::default()
3555 };
3556 let local_bytes = bad_local.marshal().expect("marshal bad local INIT");
3557 
3558 let cfg = ClientConfig::new().with_snap(local_bytes, remote_init);
3559 let res = endpoint.connect(cfg, remote_addr);
3560 assert!(
3561 matches!(
3562 res,
3563 Err(ConnectError::Snap(SnapError::InvalidInit {
3564 side: SnapSide::Local,
3565 ..
3566 }))
3567 ),
3568 "local INIT with a_rwnd=0 should be rejected: got {:?}",
3569 res
3570 );
3571}
3572 
3573#[test]
3574fn test_snap_generate_token_roundtrip_is_valid() {
3575 let config = TransportConfig::default();
3576 let bytes = generate_snap_token(&config).expect("generate token");
3577 
3578 assert!(
3579 bytes.len() <= crate::config::MAX_SNAP_INIT_BYTES,
3580 "generated token should be within MAX_SNAP_INIT_BYTES"
3581 );
3582 
3583 let init = ChunkInit::unmarshal(&bytes).expect("unmarshal generated token");
3584 assert!(!init.is_ack, "generated token should be INIT, not INIT-ACK");
3585 assert_ne!(init.initiate_tag, 0, "initiate_tag must not be zero");
3586 assert_ne!(init.initial_tsn, 0, "initial_tsn must not be zero");
3587 assert_eq!(init.num_outbound_streams, u16::MAX);
3588 assert_eq!(init.num_inbound_streams, u16::MAX);
3589 assert_eq!(
3590 init.advertised_receiver_window_credit,
3591 config.max_receive_buffer_size()
3592 );
3593 
3594 // check() should pass for a token we generated ourselves
3595 init.check()
3596 .expect("check() should pass on generated token");
3597}
3598 
3599#[test]
3600fn test_snap_generate_token_unique_per_call() {
3601 let config = TransportConfig::default();
3602 let t1 = generate_snap_token(&config).expect("token 1");
3603 let t2 = generate_snap_token(&config).expect("token 2");
3604 assert_ne!(t1, t2, "two generated tokens should differ (random values)");
3605 
3606 let i1 = ChunkInit::unmarshal(&t1).expect("parse t1");
3607 let i2 = ChunkInit::unmarshal(&t2).expect("parse t2");
3608 assert_ne!(
3609 i1.initiate_tag, i2.initiate_tag,
3610 "initiate_tags should differ"
3611 );
3612}
3613 
3614#[test]
3615fn test_snap_rejects_wrong_chunk_type_bytes() {
3616 use crate::chunk::Chunk;
3617 use crate::chunk::chunk_selective_ack::ChunkSelectiveAck;
3618 
3619 let mut endpoint = Endpoint::new(Arc::new(EndpointConfig::default()), None);
3620 let remote_addr: SocketAddr = "127.0.0.1:5000".parse().unwrap();
3621 
3622 let valid = generate_snap_token(&TransportConfig::default()).expect("valid init");
3623 
3624 // Build a SACK chunk and try to use it as a SNAP INIT
3625 let sack = ChunkSelectiveAck {
3626 cumulative_tsn_ack: 1,
3627 advertised_receiver_window_credit: 65535,
3628 gap_ack_blocks: vec![],
3629 duplicate_tsn: vec![],
3630 };
3631 let sack_bytes = sack.marshal().expect("marshal SACK");
3632 
3633 let cfg = ClientConfig::new().with_snap(valid, sack_bytes);
3634 let res = endpoint.connect(cfg, remote_addr);
3635 assert!(
3636 matches!(
3637 res,
3638 Err(ConnectError::Snap(SnapError::ParseFailed {
3639 side: SnapSide::Remote,
3640 ..
3641 }))
3642 ),
3643 "SACK bytes should fail to parse as INIT: got {:?}",
3644 res
3645 );
3646}
3647 
3648#[test]
3649fn test_snap_rejects_both_sides_empty() {
3650 let mut endpoint = Endpoint::new(Arc::new(EndpointConfig::default()), None);
3651 let remote_addr: SocketAddr = "127.0.0.1:5000".parse().unwrap();
3652 
3653 let cfg = ClientConfig::new().with_snap(Bytes::new(), Bytes::new());
3654 let res = endpoint.connect(cfg, remote_addr);
3655 assert!(
3656 matches!(res, Err(ConnectError::Snap(SnapError::ParseFailed { .. }))),
3657 "two empty byte slices should fail"
3658 );
3659}
3660 
3661#[test]
3662fn test_snap_rejects_random_garbage_bytes() {
3663 let mut endpoint = Endpoint::new(Arc::new(EndpointConfig::default()), None);
3664 let remote_addr: SocketAddr = "127.0.0.1:5000".parse().unwrap();
3665 
3666 let valid = generate_snap_token(&TransportConfig::default()).expect("valid init");
3667 
3668 // 50 bytes of pseudo-random garbage
3669 let garbage = Bytes::from(vec![0xAB; 50]);
3670 
3671 let cfg = ClientConfig::new().with_snap(valid.clone(), garbage.clone());
3672 let res = endpoint.connect(cfg, remote_addr);
3673 assert!(
3674 matches!(
3675 res,
3676 Err(ConnectError::Snap(SnapError::ParseFailed {
3677 side: SnapSide::Remote,
3678 ..
3679 }))
3680 ),
3681 "garbage remote bytes should fail: got {:?}",
3682 res
3683 );
3684 
3685 let cfg = ClientConfig::new().with_snap(garbage, valid);
3686 let res = endpoint.connect(cfg, remote_addr);
3687 assert!(
3688 matches!(
3689 res,
3690 Err(ConnectError::Snap(SnapError::ParseFailed {
3691 side: SnapSide::Local,
3692 ..
3693 }))
3694 ),
3695 "garbage local bytes should fail: got {:?}",
3696 res
3697 );
3698}
3699 
3700#[test]
3701fn test_snap_identical_local_and_remote_init() {
3702 let mut endpoint = Endpoint::new(Arc::new(EndpointConfig::default()), None);
3703 let remote_addr: SocketAddr = "127.0.0.1:5000".parse().unwrap();
3704 
3705 let token = generate_snap_token(&TransportConfig::default()).expect("token");
3706 
3707 let cfg = ClientConfig::new().with_snap(token.clone(), token);
3708 let res = endpoint.connect(cfg, remote_addr);
3709 assert!(
3710 matches!(res, Err(ConnectError::Snap(SnapError::AidCollision { .. }))),
3711 "same token for both sides should collide: got {:?}",
3712 res
3713 );
3714}
3715 
3716#[test]
3717fn test_snap_end_to_end_bidirectional_data() {
3718 // Two endpoints, both using Endpoint::connect with SNAP, exchanging data
3719 // through the full Endpoint::handle pipeline.
3720 let endpoint_config = Arc::new(EndpointConfig::default());
3721 let transport = TransportConfig::default();
3722 
3723 let init_a = generate_snap_token(&transport).expect("init A");
3724 let init_b = generate_snap_token(&transport).expect("init B");
3725 
3726 let addr_a: SocketAddr = SocketAddr::new(
3727 Ipv6Addr::LOCALHOST.into(),
3728 CLIENT_PORTS.lock().unwrap().next().unwrap(),
3729 );
3730 let addr_b: SocketAddr = SocketAddr::new(
3731 Ipv6Addr::LOCALHOST.into(),
3732 CLIENT_PORTS.lock().unwrap().next().unwrap(),
3733 );
3734 
3735 // Both sides are clients (no server config) — this is the SNAP model.
3736 let mut ep_a = Endpoint::new(endpoint_config.clone(), None);
3737 let mut ep_b = Endpoint::new(endpoint_config, None);
3738 
3739 let cfg_a = ClientConfig::new().with_snap(init_a.clone(), init_b.clone());
3740 let (ch_a, mut assoc_a) = ep_a.connect(cfg_a, addr_b).expect("SNAP connect A");
3741 
3742 let cfg_b = ClientConfig::new().with_snap(init_b, init_a);
3743 let (ch_b, mut assoc_b) = ep_b.connect(cfg_b, addr_a).expect("SNAP connect B");
3744 
3745 // Both should be established immediately.
3746 assert_eq!(assoc_a.state(), AssociationState::Established);
3747 assert_eq!(assoc_b.state(), AssociationState::Established);
3748 assert_matches!(assoc_a.poll(), Some(Event::Connected));
3749 assert_matches!(assoc_b.poll(), Some(Event::Connected));
3750 
3751 let now = Instant::now();
3752 
3753 // A writes to B.
3754 let si = 1u16;
3755 let msg_a = Bytes::from_static(b"hello from A");
3756 let mut stream_a = assoc_a
3757 .open_stream(si, PayloadProtocolIdentifier::Binary)
3758 .expect("open stream A");
3759 stream_a
3760 .write_sctp(&msg_a, PayloadProtocolIdentifier::Binary)
3761 .expect("write A");
3762 
3763 // Pump transmits from A into B via Endpoint::handle.
3764 while let Some(transmit) = assoc_a.poll_transmit(now) {
3765 if let Payload::RawEncode(datagrams) = transmit.payload {
3766 for dgram in datagrams {
3767 if let Some((handle, event)) = ep_b.handle(now, addr_a, None, None, dgram) {
3768 assert_eq!(handle, ch_b);
3769 if let DatagramEvent::AssociationEvent(ev) = event {
3770 assoc_b.handle_event(ev);
3771 }
3772 }
3773 }
3774 }
3775 }
3776 
3777 // B should now have data.
3778 let accepted = assoc_b.accept_stream().expect("accept stream on B");
3779 assert_eq!(accepted.stream_identifier, si);
3780 
3781 let mut buf = vec![0u8; 64];
3782 let chunks = assoc_b
3783 .stream(si)
3784 .expect("get stream B")
3785 .read_sctp()
3786 .expect("read B")
3787 .expect("expected data");
3788 let n = chunks.read(&mut buf).expect("read payload");
3789 assert_eq!(&buf[..n], b"hello from A");
3790 
3791 // B writes back to A.
3792 let msg_b = Bytes::from_static(b"hello from B");
3793 assoc_b
3794 .stream(si)
3795 .expect("get stream B for write")
3796 .write_sctp(&msg_b, PayloadProtocolIdentifier::Binary)
3797 .expect("write B");
3798 
3799 // Pump transmits from B into A.
3800 while let Some(transmit) = assoc_b.poll_transmit(now) {
3801 if let Payload::RawEncode(datagrams) = transmit.payload {
3802 for dgram in datagrams {
3803 if let Some((handle, event)) = ep_a.handle(now, addr_b, None, None, dgram) {
3804 assert_eq!(handle, ch_a);
3805 if let DatagramEvent::AssociationEvent(ev) = event {
3806 assoc_a.handle_event(ev);
3807 }
3808 }
3809 }
3810 }
3811 }
3812 
3813 // A should now have data back from B.
3814 let chunks = assoc_a
3815 .stream(si)
3816 .expect("get stream A")
3817 .read_sctp()
3818 .expect("read A")
3819 .expect("expected data from B");
3820 let n = chunks.read(&mut buf).expect("read payload from B");
3821 assert_eq!(&buf[..n], b"hello from B");
3822}
3823 
3824#[test]
3825fn test_snap_association_drains_cleanly() {
3826 // Verify that a SNAP association, once drained, properly cleans up both
3827 // association_ids and association_ids_init in the endpoint.
3828 let mut endpoint = Endpoint::new(Arc::new(EndpointConfig::default()), None);
3829 let remote_addr: SocketAddr = "127.0.0.1:5000".parse().unwrap();
3830 
3831 let local_init = generate_snap_token(&TransportConfig::default()).expect("local init");
3832 let remote_init = generate_snap_token(&TransportConfig::default()).expect("remote init");
3833 
3834 let local_tag = ChunkInit::unmarshal(&local_init)
3835 .expect("parse local")
3836 .initiate_tag;
3837 let remote_tag = ChunkInit::unmarshal(&remote_init)
3838 .expect("parse remote")
3839 .initiate_tag;
3840 
3841 let cfg = ClientConfig::new().with_snap(local_init, remote_init);
3842 let (ch, _assoc) = endpoint.connect(cfg, remote_addr).expect("SNAP connect");
3843 
3844 // Endpoint should have entries in both routing tables.
3845 assert!(
3846 endpoint.association_ids.contains_key(&local_tag),
3847 "local_tag should be in association_ids"
3848 );
3849 assert!(
3850 endpoint.association_ids_init.contains_key(&remote_tag),
3851 "remote_tag should be in association_ids_init"
3852 );
3853 
3854 // Simulate draining: send the Drained endpoint event.
3855 let drain_event = EndpointEvent(EndpointEventInner::Drained);
3856 endpoint.handle_event(ch, drain_event);
3857 
3858 // Both routing tables should be cleaned up.
3859 assert!(
3860 !endpoint.association_ids.contains_key(&local_tag),
3861 "local_tag should be removed from association_ids after drain"
3862 );
3863 assert!(
3864 !endpoint.association_ids_init.contains_key(&remote_tag),
3865 "remote_tag should be removed from association_ids_init after drain"
3866 );
3867}