File
Blob: firmware/vendor/sctp-proto/src/association/receive_limits_test.rs
| 1 | use super::*; |
| 2 | use crate::ReceiveLimits; |
| 3 | use crate::packet::PartialDecode; |
| 4 | |
| 5 | fn association(message: u32, bytes: u32, chunks: usize, streams: usize) -> Association { |
| 6 | Association { |
| 7 | state: AssociationState::Established, |
| 8 | peer_last_tsn: 0, |
| 9 | my_next_tsn: 1, |
| 10 | source_port: 5000, |
| 11 | destination_port: 5000, |
| 12 | max_receive_buffer_size: bytes, |
| 13 | max_receive_message_size: message, |
| 14 | receive_limits: Some(ReceiveLimits::new(message, bytes, chunks, streams)), |
| 15 | max_payload_size: 1200, |
| 16 | ..Default::default() |
| 17 | } |
| 18 | } |
| 19 | |
| 20 | fn data(tsn: u32, stream: u16, ssn: u16, len: usize) -> ChunkPayloadData { |
| 21 | ChunkPayloadData { |
| 22 | tsn, |
| 23 | stream_identifier: stream, |
| 24 | stream_sequence_number: ssn, |
| 25 | beginning_fragment: true, |
| 26 | ending_fragment: true, |
| 27 | payload_type: PayloadProtocolIdentifier::Binary, |
| 28 | user_data: Bytes::from(vec![42; len]), |
| 29 | ..Default::default() |
| 30 | } |
| 31 | } |
| 32 | |
| 33 | fn receive(a: &mut Association, chunk: ChunkPayloadData) { |
| 34 | let packet = a.create_packet(vec![Box::new(chunk)]).marshal().unwrap(); |
| 35 | let partial = PartialDecode::unmarshal(&packet).unwrap(); |
| 36 | a.handle_event(AssociationEvent(AssociationEventInner::Datagram( |
| 37 | Transmit { |
| 38 | now: Instant::now(), |
| 39 | remote: a.remote_addr, |
| 40 | ecn: None, |
| 41 | local_ip: None, |
| 42 | payload: Payload::PartialDecode(partial), |
| 43 | }, |
| 44 | ))); |
| 45 | } |
| 46 | |
| 47 | #[test] |
| 48 | fn policy_is_opt_in_and_preserves_high_stream_identifiers() { |
| 49 | assert!(TransportConfig::default().receive_limits().is_none()); |
| 50 | let limits = ReceiveLimits::new(8192, 32768, 64, 8); |
| 51 | let config = TransportConfig::default().with_receive_limits(limits); |
| 52 | assert_eq!(config.max_receive_buffer_size(), 32768); |
| 53 | assert_eq!(config.max_receive_message_size(), 8192); |
| 54 | assert_eq!(config.max_num_inbound_streams(), u16::MAX); |
| 55 | let overridden = config |
| 56 | .with_max_receive_message_size(u32::MAX) |
| 57 | .with_max_receive_buffer_size(u32::MAX); |
| 58 | assert_eq!(overridden.max_receive_message_size(), 8192); |
| 59 | assert_eq!(overridden.max_receive_buffer_size(), 32768); |
| 60 | assert_eq!( |
| 61 | overridden |
| 62 | .with_max_receive_message_size(1024) |
| 63 | .max_receive_message_size(), |
| 64 | 1024 |
| 65 | ); |
| 66 | let mut a = association(8192, 32768, 64, 8); |
| 67 | for (n, stream) in [0, 64, 257, 32768, 65530, 65534, 4000, 5000] |
| 68 | .into_iter() |
| 69 | .enumerate() |
| 70 | { |
| 71 | a.handle_data(&data(n as u32 + 1, stream, 0, 1)).unwrap(); |
| 72 | assert!(a.stream(stream).unwrap().read().unwrap().is_some()); |
| 73 | } |
| 74 | assert_eq!(a.streams.len(), 8); |
| 75 | assert_eq!( |
| 76 | a.handle_data(&data(9, 6000, 0, 1)).unwrap_err(), |
| 77 | Error::ErrReceiveLimitExceeded |
| 78 | ); |
| 79 | assert_eq!(a.streams.len(), 8); |
| 80 | assert!( |
| 81 | a.open_stream(65000, PayloadProtocolIdentifier::Binary) |
| 82 | .is_err() |
| 83 | ); |
| 84 | } |
| 85 | |
| 86 | #[test] |
| 87 | fn fragmented_message_limit_is_checked_before_reassembly_grows() { |
| 88 | for unordered in [false, true] { |
| 89 | let mut a = association(8192, 32768, 64, 8); |
| 90 | for n in 0..8 { |
| 91 | let mut chunk = data(n + 1, 65000, 0, 1024); |
| 92 | chunk.unordered = unordered; |
| 93 | chunk.beginning_fragment = n == 0; |
| 94 | chunk.ending_fragment = false; |
| 95 | a.handle_data(&chunk).unwrap(); |
| 96 | } |
| 97 | assert_eq!(a.retained_receive_data(), (8192, 8)); |
| 98 | let mut tail = data(9, 65000, 0, 1); |
| 99 | tail.unordered = unordered; |
| 100 | tail.beginning_fragment = false; |
| 101 | assert_eq!( |
| 102 | a.handle_data(&tail).unwrap_err(), |
| 103 | Error::ErrInboundPacketTooLarge |
| 104 | ); |
| 105 | assert_eq!(a.retained_receive_data(), (8192, 8)); |
| 106 | } |
| 107 | } |
| 108 | |
| 109 | #[test] |
| 110 | fn all_receive_queues_share_the_byte_and_fragment_budget() { |
| 111 | let mut a = association(16, 32, 64, 8); |
| 112 | // Complete data on an unordered stream can be consumed while a missing |
| 113 | // earlier TSN keeps the same payload alive in the association queue. |
| 114 | let mut first = data(2, 0, 0, 16); |
| 115 | first.unordered = true; |
| 116 | a.handle_data(&first).unwrap(); |
| 117 | assert_eq!(a.retained_receive_data(), (16, 1)); |
| 118 | drop(a.stream(0).unwrap().read().unwrap().unwrap()); |
| 119 | assert_eq!(a.get_my_receiver_window_credit(), 16); |
| 120 | assert_eq!(a.retained_receive_data(), (16, 1)); |
| 121 | let mut second = data(3, 65000, 0, 16); |
| 122 | second.ending_fragment = false; |
| 123 | a.handle_data(&second).unwrap(); |
| 124 | assert_eq!(a.retained_receive_data(), (32, 2)); |
| 125 | // Missing TSNs are not allowed to bypass the hard resource budget. |
| 126 | assert_eq!( |
| 127 | a.handle_data(&data(1, 0, 1, 1)).unwrap_err(), |
| 128 | Error::ErrReceiveLimitExceeded |
| 129 | ); |
| 130 | assert_eq!(a.retained_receive_data(), (32, 2)); |
| 131 | } |
| 132 | |
| 133 | #[test] |
| 134 | fn missing_tsn_can_fill_a_gap_and_restore_receive_credit_within_the_budget() { |
| 135 | let mut a = association(16, 32, 64, 8); |
| 136 | let mut first = data(2, 0, 0, 8); |
| 137 | first.unordered = true; |
| 138 | a.handle_data(&first).unwrap(); |
| 139 | drop(a.stream(0).unwrap().read().unwrap().unwrap()); |
| 140 | assert_eq!(a.get_my_receiver_window_credit(), 24); |
| 141 | let mut partial = data(3, 65000, 0, 8); |
| 142 | partial.ending_fragment = false; |
| 143 | a.handle_data(&partial).unwrap(); |
| 144 | assert_eq!(a.get_my_receiver_window_credit(), 16); |
| 145 | a.handle_data(&data(1, 257, 0, 8)).unwrap(); |
| 146 | assert_eq!(a.peer_last_tsn, 3); |
| 147 | assert!(a.payload_queue.is_empty()); |
| 148 | drop(a.stream(257).unwrap().read().unwrap().unwrap()); |
| 149 | assert_eq!(a.get_my_receiver_window_credit(), 24); |
| 150 | let mut tail = data(4, 65000, 0, 8); |
| 151 | tail.beginning_fragment = false; |
| 152 | a.handle_data(&tail).unwrap(); |
| 153 | assert_eq!(a.stream(65000).unwrap().read().unwrap().unwrap().len(), 16); |
| 154 | assert_eq!(a.retained_receive_data(), (0, 0)); |
| 155 | assert_eq!(a.get_my_receiver_window_credit(), 32); |
| 156 | assert!(!a.is_closed()); |
| 157 | } |
| 158 | |
| 159 | #[test] |
| 160 | fn tiny_incomplete_messages_cannot_exhaust_chunk_metadata() { |
| 161 | let mut a = association(8192, 32768, 64, 8); |
| 162 | for n in 0..64 { |
| 163 | let mut chunk = data(n + 1, 0, n as u16, 1); |
| 164 | chunk.ending_fragment = false; |
| 165 | a.handle_data(&chunk).unwrap(); |
| 166 | } |
| 167 | assert_eq!(a.retained_receive_data(), (64, 64)); |
| 168 | assert_eq!(a.get_my_receiver_window_credit(), 0); |
| 169 | assert_eq!( |
| 170 | a.handle_data(&data(65, 0, 64, 1)).unwrap_err(), |
| 171 | Error::ErrReceiveLimitExceeded |
| 172 | ); |
| 173 | assert_eq!(a.retained_receive_data(), (64, 64)); |
| 174 | } |
| 175 | |
| 176 | #[test] |
| 177 | fn reset_deferred_data_is_counted_before_new_generation_is_readable() { |
| 178 | let mut a = association(8, 16, 64, 8); |
| 179 | a.handle_data(&data(1, 0, 0, 8)).unwrap(); |
| 180 | let reset: Box<dyn Param + Send + Sync> = Box::new(ParamOutgoingResetRequest { |
| 181 | reconfig_request_sequence_number: 7, |
| 182 | reconfig_response_sequence_number: u32::MAX, |
| 183 | sender_last_tsn: 1, |
| 184 | stream_identifiers: vec![0], |
| 185 | }); |
| 186 | a.handle_reconfig_param(&reset, &mut vec![]).unwrap(); |
| 187 | assert!(a.retiring_streams.contains_key(&0)); |
| 188 | a.handle_data(&data(2, 0, 0, 8)).unwrap(); |
| 189 | assert!(a.deferred_reset_data.contains_key(&2)); |
| 190 | assert_eq!(a.retained_receive_data(), (16, 2)); |
| 191 | assert_eq!( |
| 192 | a.handle_data(&data(3, 0, 1, 1)).unwrap_err(), |
| 193 | Error::ErrReceiveLimitExceeded |
| 194 | ); |
| 195 | assert_eq!(a.retained_receive_data(), (16, 2)); |
| 196 | drop(a.stream(0).unwrap().read().unwrap().unwrap()); |
| 197 | assert!(a.retained_receive_data().0 <= 8); |
| 198 | } |
| 199 | |
| 200 | #[test] |
| 201 | fn compact_payloads_do_not_keep_unrelated_packet_bytes_alive() { |
| 202 | let mut a = association(8, 16, 64, 8); |
| 203 | let packet = Bytes::from(vec![42; 8192]); |
| 204 | let mut chunk = data(2, 0, 0, 1); |
| 205 | chunk.user_data = packet.slice(4000..4001); |
| 206 | a.handle_data(&chunk).unwrap(); |
| 207 | let retained = &a.payload_queue.get(2).unwrap().user_data; |
| 208 | assert_eq!(retained.as_ref(), chunk.user_data.as_ref()); |
| 209 | assert_ne!(retained.as_ptr(), chunk.user_data.as_ptr()); |
| 210 | let queued = &a.streams.get(&0).unwrap().reassembly_queue.ordered[0].chunks[0]; |
| 211 | assert_eq!(retained.as_ptr(), queued.user_data.as_ptr()); |
| 212 | } |
| 213 | |
| 214 | #[test] |
| 215 | fn overflow_on_the_wire_closes_the_association_and_fresh_peer_can_receive() { |
| 216 | let mut a = association(8, 8, 2, 8); |
| 217 | receive(&mut a, data(1, 0, 0, 8)); |
| 218 | assert!(!a.is_closed()); |
| 219 | receive(&mut a, data(2, 0, 1, 1)); |
| 220 | assert!(a.is_closed()); |
| 221 | assert!(core::iter::from_fn(|| a.poll()).any(|e| matches!(e, Event::AssociationLost { .. }))); |
| 222 | assert!(a.streams.is_empty()); |
| 223 | assert!(a.deferred_reset_data.is_empty()); |
| 224 | let mut fresh = association(8, 8, 2, 8); |
| 225 | receive(&mut fresh, data(1, 65000, 0, 8)); |
| 226 | assert!(!fresh.is_closed()); |
| 227 | assert_eq!( |
| 228 | fresh.stream(65000).unwrap().read().unwrap().unwrap().len(), |
| 229 | 8 |
| 230 | ); |
| 231 | } |