File
Blob: firmware/vendor/sctp-proto/src/endpoint/endpoint_test.rs
| 1 | use alloc::borrow::ToOwned; |
| 2 | use alloc::vec; |
| 3 | use alloc::vec::Vec; |
| 4 | use std::println; |
| 5 | |
| 6 | use super::*; |
| 7 | use crate::association::Event; |
| 8 | use crate::config::generate_snap_token; |
| 9 | use crate::error::{Error, Result}; |
| 10 | |
| 11 | use crate::association::state::{AckMode, AssociationState}; |
| 12 | use crate::association::stream::{ReliabilityType, Stream, StreamEvent, StreamResetError}; |
| 13 | use crate::chunk::chunk_abort::ChunkAbort; |
| 14 | use crate::chunk::chunk_cookie_echo::ChunkCookieEcho; |
| 15 | use crate::chunk::chunk_error::ChunkError; |
| 16 | use crate::chunk::chunk_forward_tsn::ChunkForwardTsn; |
| 17 | use crate::chunk::chunk_heartbeat::ChunkHeartbeat; |
| 18 | use crate::chunk::chunk_init::ChunkInit; |
| 19 | use crate::chunk::chunk_payload_data::{ChunkPayloadData, PayloadProtocolIdentifier}; |
| 20 | use crate::chunk::chunk_reconfig::ChunkReconfig; |
| 21 | use crate::chunk::chunk_selective_ack::{ChunkSelectiveAck, GapAckBlock}; |
| 22 | use crate::chunk::chunk_shutdown::ChunkShutdown; |
| 23 | use crate::chunk::chunk_shutdown_ack::ChunkShutdownAck; |
| 24 | use crate::chunk::chunk_shutdown_complete::ChunkShutdownComplete; |
| 25 | use crate::chunk::{ErrorCauseProtocolViolation, PROTOCOL_VIOLATION}; |
| 26 | use crate::packet::{CommonHeader, Packet}; |
| 27 | use crate::param::param_outgoing_reset_request::ParamOutgoingResetRequest; |
| 28 | use crate::param::param_reconfig_response::ParamReconfigResponse; |
| 29 | use assert_matches::assert_matches; |
| 30 | use core::net::Ipv6Addr; |
| 31 | use core::ops::RangeFrom; |
| 32 | use core::str::FromStr; |
| 33 | use core::time::Duration; |
| 34 | use core::{cmp, mem}; |
| 35 | use log::{info, trace}; |
| 36 | use std::net::UdpSocket; |
| 37 | use std::sync::{LazyLock, Mutex}; |
| 38 | use std::time::Instant; |
| 39 | |
| 40 | pub static SERVER_PORTS: LazyLock<Mutex<RangeFrom<u16>>> = LazyLock::new(|| Mutex::new(4433..)); |
| 41 | pub static CLIENT_PORTS: LazyLock<Mutex<RangeFrom<u16>>> = LazyLock::new(|| Mutex::new(44433..)); |
| 42 | |
| 43 | fn 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` |
| 53 | const MAX_DATAGRAMS: usize = 10; |
| 54 | |
| 55 | fn 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 | |
| 72 | pub fn client_config() -> ClientConfig { |
| 73 | ClientConfig::new() |
| 74 | } |
| 75 | |
| 76 | pub fn server_config() -> ServerConfig { |
| 77 | ServerConfig::new() |
| 78 | } |
| 79 | |
| 80 | struct 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 | |
| 93 | impl 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 | |
| 198 | impl ::core::ops::Deref for TestEndpoint { |
| 199 | type Target = Endpoint; |
| 200 | fn deref(&self) -> &Endpoint { |
| 201 | &self.endpoint |
| 202 | } |
| 203 | } |
| 204 | |
| 205 | impl ::core::ops::DerefMut for TestEndpoint { |
| 206 | fn deref_mut(&mut self) -> &mut Endpoint { |
| 207 | &mut self.endpoint |
| 208 | } |
| 209 | } |
| 210 | |
| 211 | struct Pair { |
| 212 | server: TestEndpoint, |
| 213 | client: TestEndpoint, |
| 214 | time: Instant, |
| 215 | latency: Duration, // One-way |
| 216 | } |
| 217 | |
| 218 | impl 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 | |
| 363 | impl Default for Pair { |
| 364 | fn default() -> Self { |
| 365 | Pair::new(Default::default(), server_config()) |
| 366 | } |
| 367 | } |
| 368 | |
| 369 | fn 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 | |
| 397 | fn 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 | |
| 440 | fn 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] |
| 468 | fn 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] |
| 520 | fn 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] |
| 600 | fn 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] |
| 662 | fn 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] |
| 724 | fn 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] |
| 810 | fn 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] |
| 875 | fn 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] |
| 934 | fn 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] |
| 1009 | fn 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] |
| 1088 | fn 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] |
| 1162 | fn 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] |
| 1244 | fn 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] |
| 1318 | fn 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] |
| 1411 | fn 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] |
| 1499 | fn 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] |
| 1627 | fn 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] |
| 1736 | fn 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] |
| 1826 | fn 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] |
| 1885 | fn 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] |
| 1969 | fn 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] |
| 2031 | fn 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] |
| 2219 | fn 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 | /* |
| 2287 | TODO: The following tests will be moved to sctp-async tests: |
| 2288 | struct 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 | |
| 2295 | impl 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 | |
| 2307 | trait AsAny { |
| 2308 | fn as_any(&self) -> &(dyn std::any::Any + Send + Sync); |
| 2309 | } |
| 2310 | |
| 2311 | impl AsAny for FakeEchoConn { |
| 2312 | fn as_any(&self) -> &(dyn std::any::Any + Send + Sync) { |
| 2313 | self |
| 2314 | } |
| 2315 | } |
| 2316 | |
| 2317 | type UResult<T> = std::result::Result<T, util::Error>; |
| 2318 | |
| 2319 | #[async_trait] |
| 2320 | impl 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] |
| 2373 | fn 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 | |
| 2411 | fn 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] |
| 2483 | fn 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 | |
| 2539 | fn 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] |
| 2613 | fn 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] |
| 2669 | fn 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] |
| 2783 | fn 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] |
| 2874 | fn 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] |
| 2922 | fn 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] |
| 2986 | fn 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] |
| 3038 | fn 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] |
| 3089 | fn 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] |
| 3109 | fn 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] |
| 3129 | fn 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] |
| 3155 | fn 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] |
| 3181 | fn 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] |
| 3217 | fn 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] |
| 3245 | fn 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] |
| 3271 | fn 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] |
| 3302 | fn 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] |
| 3338 | fn 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] |
| 3375 | fn 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] |
| 3414 | fn 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] |
| 3446 | fn 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] |
| 3478 | fn 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] |
| 3510 | fn 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] |
| 3542 | fn 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] |
| 3574 | fn 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] |
| 3600 | fn 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] |
| 3615 | fn 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] |
| 3649 | fn 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] |
| 3662 | fn 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] |
| 3701 | fn 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] |
| 3717 | fn 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] |
| 3825 | fn 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 | } |