1use std::collections::BTreeMap;
28use std::collections::VecDeque;
29
30use std::net::SocketAddr;
31
32use std::time::Duration;
33use std::time::Instant;
34
35use smallvec::SmallVec;
36
37use slab::Slab;
38
39use crate::Config;
40use crate::Error;
41use crate::Result;
42use crate::StartupExit;
43
44use crate::pmtud;
45use crate::recovery;
46use crate::recovery::Bandwidth;
47use crate::recovery::HandshakeStatus;
48use crate::recovery::OnLossDetectionTimeoutOutcome;
49use crate::recovery::RecoveryOps;
50
51#[derive(Debug, Copy, Clone, PartialEq, Eq, PartialOrd, Ord)]
53pub enum PathState {
54 Failed,
56
57 Unknown,
59
60 Validating,
62
63 ValidatingMTU,
65
66 Validated,
68}
69
70impl PathState {
71 #[cfg(feature = "ffi")]
72 pub fn to_c(self) -> libc::ssize_t {
73 match self {
74 PathState::Failed => -1,
75 PathState::Unknown => 0,
76 PathState::Validating => 1,
77 PathState::ValidatingMTU => 2,
78 PathState::Validated => 3,
79 }
80 }
81}
82
83#[derive(Clone, Debug, PartialEq, Eq)]
85pub enum PathEvent {
86 New(SocketAddr, SocketAddr),
91
92 Validated(SocketAddr, SocketAddr),
95
96 FailedValidation(SocketAddr, SocketAddr),
100
101 Closed(SocketAddr, SocketAddr),
104
105 ReusedSourceConnectionId(
109 u64,
110 (SocketAddr, SocketAddr),
111 (SocketAddr, SocketAddr),
112 ),
113
114 PeerMigrated(SocketAddr, SocketAddr),
120}
121
122#[derive(Debug)]
124pub struct Path {
125 local_addr: SocketAddr,
127
128 peer_addr: SocketAddr,
130
131 pub active_scid_seq: Option<u64>,
133
134 pub active_dcid_seq: Option<u64>,
136
137 state: PathState,
139
140 active: bool,
142
143 pub recovery: recovery::Recovery,
145
146 pub pmtud: Option<pmtud::Pmtud>,
148
149 in_flight_challenges: VecDeque<([u8; 8], usize, Instant)>,
152
153 max_challenge_size: usize,
155
156 probing_lost: usize,
158
159 last_probe_lost_time: Option<Instant>,
161
162 received_challenges: VecDeque<[u8; 8]>,
164
165 received_challenges_max_len: usize,
167
168 pub sent_count: usize,
170
171 pub recv_count: usize,
173
174 pub retrans_count: usize,
176
177 pub total_pto_count: usize,
184
185 pub dgram_sent_count: usize,
187
188 pub dgram_lost_count: usize,
190
191 pub dgram_recv_count: usize,
193
194 pub sent_bytes: u64,
196
197 pub recv_bytes: u64,
199
200 pub stream_retrans_bytes: u64,
203
204 pub max_send_bytes: usize,
207
208 pub verified_peer_address: bool,
210
211 pub peer_verified_local_address: bool,
213
214 challenge_requested: bool,
216
217 failure_notified: bool,
219
220 migrating: bool,
223
224 pub needs_ack_eliciting: bool,
226}
227
228impl Path {
229 pub fn new(
232 local_addr: SocketAddr, peer_addr: SocketAddr,
233 recovery_config: &recovery::RecoveryConfig,
234 path_challenge_recv_max_queue_len: usize, is_initial: bool,
235 config: Option<&Config>,
236 ) -> Self {
237 let (state, active_scid_seq, active_dcid_seq) = if is_initial {
238 (PathState::Validated, Some(0), Some(0))
239 } else {
240 (PathState::Unknown, None, None)
241 };
242
243 let pmtud = config.and_then(|c| {
244 if c.pmtud {
245 let maximum_supported_mtu: usize = std::cmp::min(
246 c.local_transport_params
249 .max_udp_payload_size
250 .try_into()
251 .unwrap_or(c.max_send_udp_payload_size),
252 c.max_send_udp_payload_size,
253 );
254 Some(pmtud::Pmtud::new(maximum_supported_mtu, c.pmtud_max_probes))
255 } else {
256 None
257 }
258 });
259
260 Self {
261 local_addr,
262 peer_addr,
263 active_scid_seq,
264 active_dcid_seq,
265 state,
266 active: false,
267 recovery: recovery::Recovery::new_with_config(recovery_config),
268 pmtud,
269 in_flight_challenges: VecDeque::new(),
270 max_challenge_size: 0,
271 probing_lost: 0,
272 last_probe_lost_time: None,
273 received_challenges: VecDeque::with_capacity(
274 path_challenge_recv_max_queue_len,
275 ),
276 received_challenges_max_len: path_challenge_recv_max_queue_len,
277 sent_count: 0,
278 recv_count: 0,
279 retrans_count: 0,
280 total_pto_count: 0,
281 dgram_sent_count: 0,
282 dgram_lost_count: 0,
283 dgram_recv_count: 0,
284 sent_bytes: 0,
285 recv_bytes: 0,
286 stream_retrans_bytes: 0,
287 max_send_bytes: 0,
288 verified_peer_address: false,
289 peer_verified_local_address: false,
290 challenge_requested: false,
291 failure_notified: false,
292 migrating: false,
293 needs_ack_eliciting: false,
294 }
295 }
296
297 #[inline]
299 pub fn local_addr(&self) -> SocketAddr {
300 self.local_addr
301 }
302
303 #[inline]
305 pub fn peer_addr(&self) -> SocketAddr {
306 self.peer_addr
307 }
308
309 #[inline]
311 fn working(&self) -> bool {
312 self.state > PathState::Failed
313 }
314
315 #[inline]
317 pub fn active(&self) -> bool {
318 self.active && self.working() && self.active_dcid_seq.is_some()
319 }
320
321 #[inline]
323 pub fn usable(&self) -> bool {
324 self.active() ||
325 (self.state == PathState::Validated &&
326 self.active_dcid_seq.is_some())
327 }
328
329 #[inline]
331 fn unused(&self) -> bool {
332 !self.active() && self.active_dcid_seq.is_none()
334 }
335
336 #[inline]
338 pub fn probing_required(&self) -> bool {
339 !self.received_challenges.is_empty() || self.validation_requested()
340 }
341
342 fn promote_to(&mut self, state: PathState) {
345 if self.state < state {
346 self.state = state;
347 }
348 }
349
350 #[inline]
352 pub fn validated(&self) -> bool {
353 self.state == PathState::Validated
354 }
355
356 #[inline]
358 fn validation_failed(&self) -> bool {
359 self.state == PathState::Failed
360 }
361
362 #[inline]
364 pub fn under_validation(&self) -> bool {
365 matches!(self.state, PathState::Validating | PathState::ValidatingMTU)
366 }
367
368 #[inline]
370 pub fn request_validation(&mut self) {
371 self.challenge_requested = true;
372 }
373
374 #[inline]
376 pub fn validation_requested(&self) -> bool {
377 self.challenge_requested
378 }
379
380 pub fn should_send_pmtu_probe(
381 &mut self, hs_confirmed: bool, hs_done: bool, out_len: usize,
382 is_closing: bool, frames_empty: bool,
383 ) -> bool {
384 let Some(pmtud) = self.pmtud.as_mut() else {
385 return false;
386 };
387
388 (hs_confirmed && hs_done) &&
389 self.recovery.cwnd_available() > pmtud.get_probe_size() &&
390 out_len >= pmtud.get_probe_size() &&
391 pmtud.should_probe() &&
392 !is_closing &&
393 frames_empty
394 }
395
396 pub fn on_challenge_sent(&mut self) {
397 self.promote_to(PathState::Validating);
398 self.challenge_requested = false;
399 }
400
401 pub fn add_challenge_sent(
403 &mut self, data: [u8; 8], pkt_size: usize, sent_time: Instant,
404 ) {
405 self.on_challenge_sent();
406 self.in_flight_challenges
407 .push_back((data, pkt_size, sent_time));
408 }
409
410 pub fn on_challenge_received(&mut self, data: [u8; 8]) {
411 if self.received_challenges.len() == self.received_challenges_max_len {
413 return;
414 }
415
416 self.received_challenges.push_back(data);
417 self.peer_verified_local_address = true;
418 }
419
420 pub fn has_pending_challenge(&self, data: [u8; 8]) -> bool {
421 self.in_flight_challenges.iter().any(|(d, ..)| *d == data)
422 }
423
424 pub fn on_response_received(&mut self, data: [u8; 8]) -> bool {
426 self.verified_peer_address = true;
427 self.probing_lost = 0;
428
429 let mut challenge_size = 0;
430 self.in_flight_challenges.retain(|(d, s, _)| {
431 if *d == data {
432 challenge_size = *s;
433 false
434 } else {
435 true
436 }
437 });
438
439 self.promote_to(PathState::ValidatingMTU);
441
442 self.max_challenge_size =
443 std::cmp::max(self.max_challenge_size, challenge_size);
444
445 if self.state == PathState::ValidatingMTU {
446 if self.max_challenge_size >= crate::MIN_CLIENT_INITIAL_LEN {
447 self.promote_to(PathState::Validated);
449 return true;
450 }
451
452 self.request_validation();
454 }
455
456 false
457 }
458
459 fn on_failed_validation(&mut self) {
460 self.state = PathState::Failed;
461 self.active = false;
462 }
463
464 #[inline]
465 pub fn pop_received_challenge(&mut self) -> Option<[u8; 8]> {
466 self.received_challenges.pop_front()
467 }
468
469 pub fn on_loss_detection_timeout(
470 &mut self, handshake_status: HandshakeStatus, now: Instant,
471 is_server: bool, trace_id: &str,
472 ) -> OnLossDetectionTimeoutOutcome {
473 let outcome = self.recovery.on_loss_detection_timeout(
474 handshake_status,
475 now,
476 trace_id,
477 );
478
479 let mut lost_probe_time = None;
480 self.in_flight_challenges.retain(|(_, _, sent_time)| {
481 if *sent_time <= now {
482 if lost_probe_time.is_none() {
483 lost_probe_time = Some(*sent_time);
484 }
485 false
486 } else {
487 true
488 }
489 });
490
491 if let Some(lost_probe_time) = lost_probe_time {
494 self.last_probe_lost_time = match self.last_probe_lost_time {
495 Some(last) => {
496 if lost_probe_time - last >= self.recovery.rtt() {
498 self.probing_lost += 1;
499 Some(lost_probe_time)
500 } else {
501 Some(last)
502 }
503 },
504 None => {
505 self.probing_lost += 1;
506 Some(lost_probe_time)
507 },
508 };
509 if self.probing_lost >= crate::MAX_PROBING_TIMEOUTS ||
513 (is_server && self.max_send_bytes < crate::MIN_PROBING_SIZE)
514 {
515 self.on_failed_validation();
516 } else {
517 self.request_validation();
518 }
519 }
520
521 self.total_pto_count += 1;
523
524 outcome
525 }
526
527 pub fn can_reinit_recovery(&self) -> bool {
531 self.recovery.bytes_in_flight() == 0 &&
538 self.recovery.bytes_in_flight_duration() == Duration::ZERO
539 }
540
541 pub fn reinit_recovery(
542 &mut self, recovery_config: &recovery::RecoveryConfig,
543 ) {
544 self.recovery = recovery::Recovery::new_with_config(recovery_config)
545 }
546
547 pub fn stats(&self) -> PathStats {
548 let pmtu = match self.pmtud.as_ref().map(|p| p.get_current_mtu()) {
549 Some(v) => v,
550
551 None => self.recovery.max_datagram_size(),
552 };
553
554 PathStats {
555 local_addr: self.local_addr,
556 peer_addr: self.peer_addr,
557 validation_state: self.state,
558 active: self.active,
559 recv: self.recv_count,
560 sent: self.sent_count,
561 lost: self.recovery.lost_count(),
562 retrans: self.retrans_count,
563 total_pto_count: self.total_pto_count,
564 dgram_recv: self.dgram_recv_count,
565 dgram_sent: self.dgram_sent_count,
566 dgram_lost: self.dgram_lost_count,
567 rtt: self.recovery.rtt(),
568 min_rtt: self.recovery.min_rtt(),
569 max_rtt: self.recovery.max_rtt(),
570 rttvar: self.recovery.rttvar(),
571 cwnd: self.recovery.cwnd(),
572 sent_bytes: self.sent_bytes,
573 recv_bytes: self.recv_bytes,
574 lost_bytes: self.recovery.bytes_lost(),
575 stream_retrans_bytes: self.stream_retrans_bytes,
576 pmtu,
577 delivery_rate: self.recovery.delivery_rate().to_bytes_per_second(),
578 max_bandwidth: self
579 .recovery
580 .max_bandwidth()
581 .map(Bandwidth::to_bytes_per_second),
582 rtt_persistent_jump_count: self.recovery.rtt_persistent_jump_count(),
583 startup_exit: self.recovery.startup_exit(),
584 }
585 }
586
587 pub fn bytes_in_flight_duration(&self) -> Duration {
588 self.recovery.bytes_in_flight_duration()
589 }
590}
591
592#[derive(Default)]
594pub struct SocketAddrIter {
595 pub(crate) sockaddrs: SmallVec<[SocketAddr; 8]>,
596 pub(crate) index: usize,
597}
598
599impl Iterator for SocketAddrIter {
600 type Item = SocketAddr;
601
602 #[inline]
603 fn next(&mut self) -> Option<Self::Item> {
604 let v = self.sockaddrs.get(self.index)?;
605 self.index += 1;
606 Some(*v)
607 }
608}
609
610impl ExactSizeIterator for SocketAddrIter {
611 #[inline]
612 fn len(&self) -> usize {
613 self.sockaddrs.len() - self.index
614 }
615}
616
617pub struct PathMap {
619 paths: Slab<Path>,
622
623 max_concurrent_paths: usize,
625
626 addrs_to_paths: BTreeMap<(SocketAddr, SocketAddr), usize>,
629
630 events: VecDeque<PathEvent>,
632
633 is_server: bool,
635}
636
637impl PathMap {
638 pub fn new(
641 mut initial_path: Path, max_concurrent_paths: usize, is_server: bool,
642 ) -> Self {
643 let mut paths = Slab::with_capacity(1); let mut addrs_to_paths = BTreeMap::new();
645
646 let local_addr = initial_path.local_addr;
647 let peer_addr = initial_path.peer_addr;
648
649 initial_path.active = true;
651
652 let active_path_id = paths.insert(initial_path);
653 addrs_to_paths.insert((local_addr, peer_addr), active_path_id);
654
655 Self {
656 paths,
657 max_concurrent_paths,
658 addrs_to_paths,
659 events: VecDeque::new(),
660 is_server,
661 }
662 }
663
664 #[inline]
670 pub fn get(&self, path_id: usize) -> Result<&Path> {
671 self.paths.get(path_id).ok_or(Error::InvalidState)
672 }
673
674 #[inline]
680 pub fn get_mut(&mut self, path_id: usize) -> Result<&mut Path> {
681 self.paths.get_mut(path_id).ok_or(Error::InvalidState)
682 }
683
684 #[inline]
685 pub fn get_active_with_pid(&self) -> Option<(usize, &Path)> {
688 self.paths.iter().find(|(_, p)| p.active())
689 }
690
691 #[inline]
696 pub fn get_active(&self) -> Result<&Path> {
697 self.get_active_with_pid()
698 .map(|(_, p)| p)
699 .ok_or(Error::InvalidState)
700 }
701
702 #[inline]
707 pub fn get_active_path_id(&self) -> Result<usize> {
708 self.get_active_with_pid()
709 .map(|(pid, _)| pid)
710 .ok_or(Error::InvalidState)
711 }
712
713 #[inline]
718 pub fn get_active_mut(&mut self) -> Result<&mut Path> {
719 self.paths
720 .iter_mut()
721 .map(|(_, p)| p)
722 .find(|p| p.active())
723 .ok_or(Error::InvalidState)
724 }
725
726 #[inline]
728 pub fn iter(&self) -> slab::Iter<'_, Path> {
729 self.paths.iter()
730 }
731
732 #[inline]
734 pub fn iter_mut(&mut self) -> slab::IterMut<'_, Path> {
735 self.paths.iter_mut()
736 }
737
738 #[inline]
740 pub fn len(&self) -> usize {
741 self.paths.len()
742 }
743
744 #[inline]
746 pub fn path_id_from_addrs(
747 &self, addrs: &(SocketAddr, SocketAddr),
748 ) -> Option<usize> {
749 self.addrs_to_paths.get(addrs).copied()
750 }
751
752 fn make_room_for_new_path(&mut self) -> Result<()> {
758 if self.paths.len() < self.max_concurrent_paths {
759 return Ok(());
760 }
761
762 let (pid_to_remove, _) = self
763 .paths
764 .iter()
765 .find(|(_, p)| p.unused())
766 .ok_or(Error::Done)?;
767
768 let path = self.paths.remove(pid_to_remove);
769 self.addrs_to_paths
770 .remove(&(path.local_addr, path.peer_addr));
771
772 self.notify_event(PathEvent::Closed(path.local_addr, path.peer_addr));
773
774 Ok(())
775 }
776
777 pub fn insert_path(&mut self, path: Path, is_server: bool) -> Result<usize> {
788 self.make_room_for_new_path()?;
789
790 let local_addr = path.local_addr;
791 let peer_addr = path.peer_addr;
792
793 let pid = self.paths.insert(path);
794 self.addrs_to_paths.insert((local_addr, peer_addr), pid);
795
796 if is_server {
798 self.notify_event(PathEvent::New(local_addr, peer_addr));
799 }
800
801 Ok(pid)
802 }
803
804 pub fn notify_event(&mut self, ev: PathEvent) {
806 self.events.push_back(ev);
807 }
808
809 pub fn pop_event(&mut self) -> Option<PathEvent> {
811 self.events.pop_front()
812 }
813
814 pub fn notify_failed_validations(&mut self) {
816 let validation_failed = self
817 .paths
818 .iter_mut()
819 .filter(|(_, p)| p.validation_failed() && !p.failure_notified);
820
821 for (_, p) in validation_failed {
822 self.events.push_back(PathEvent::FailedValidation(
823 p.local_addr,
824 p.peer_addr,
825 ));
826
827 p.failure_notified = true;
828 }
829 }
830
831 pub fn find_candidate_path(&self) -> Option<usize> {
833 self.paths
835 .iter()
836 .find(|(_, p)| p.usable())
837 .map(|(pid, _)| pid)
838 }
839
840 pub fn on_response_received(&mut self, data: [u8; 8]) -> Result<()> {
842 let active_pid = self.get_active_path_id()?;
843
844 let challenge_pending =
845 self.iter_mut().find(|(_, p)| p.has_pending_challenge(data));
846
847 if let Some((pid, p)) = challenge_pending {
848 if p.on_response_received(data) {
849 let local_addr = p.local_addr;
850 let peer_addr = p.peer_addr;
851 let was_migrating = p.migrating;
852
853 p.migrating = false;
854
855 self.notify_event(PathEvent::Validated(local_addr, peer_addr));
857
858 if pid == active_pid && was_migrating {
861 self.notify_event(PathEvent::PeerMigrated(
862 local_addr, peer_addr,
863 ));
864 }
865 }
866 }
867 Ok(())
868 }
869
870 pub fn set_active_path(&mut self, path_id: usize) -> Result<()> {
881 let is_server = self.is_server;
882
883 if let Ok(old_active_path) = self.get_active_mut() {
884 old_active_path.active = false;
885 }
886
887 let new_active_path = self.get_mut(path_id)?;
888 new_active_path.active = true;
889
890 if is_server {
891 if new_active_path.validated() {
892 let local_addr = new_active_path.local_addr();
893 let peer_addr = new_active_path.peer_addr();
894
895 self.notify_event(PathEvent::PeerMigrated(local_addr, peer_addr));
896 } else {
897 new_active_path.migrating = true;
898
899 if !new_active_path.under_validation() {
901 new_active_path.request_validation();
902 }
903 }
904 }
905
906 Ok(())
907 }
908
909 pub fn set_discover_pmtu_on_existing_paths(
911 &mut self, discover: bool, max_send_udp_payload_size: usize,
912 pmtud_max_probes: u8,
913 ) {
914 for (_, path) in self.paths.iter_mut() {
915 path.pmtud = if discover {
916 Some(pmtud::Pmtud::new(
917 max_send_udp_payload_size,
918 pmtud_max_probes,
919 ))
920 } else {
921 None
922 };
923 }
924 }
925}
926
927#[derive(Clone)]
934#[non_exhaustive]
935pub struct PathStats {
936 pub local_addr: SocketAddr,
938
939 pub peer_addr: SocketAddr,
941
942 pub validation_state: PathState,
944
945 pub active: bool,
947
948 pub recv: usize,
950
951 pub sent: usize,
953
954 pub lost: usize,
956
957 pub retrans: usize,
959
960 pub total_pto_count: usize,
967
968 pub dgram_recv: usize,
970
971 pub dgram_sent: usize,
973
974 pub dgram_lost: usize,
976
977 pub rtt: Duration,
979
980 pub min_rtt: Option<Duration>,
982
983 pub max_rtt: Option<Duration>,
985
986 pub rttvar: Duration,
989
990 pub cwnd: usize,
992
993 pub sent_bytes: u64,
995
996 pub recv_bytes: u64,
998
999 pub lost_bytes: u64,
1001
1002 pub stream_retrans_bytes: u64,
1004
1005 pub pmtu: usize,
1007
1008 pub delivery_rate: u64,
1017
1018 pub max_bandwidth: Option<u64>,
1023
1024 pub rtt_persistent_jump_count: u64,
1026
1027 pub startup_exit: Option<StartupExit>,
1029}
1030
1031impl std::fmt::Debug for PathStats {
1032 #[inline]
1033 fn fmt(&self, f: &mut std::fmt::Formatter) -> std::fmt::Result {
1034 write!(
1035 f,
1036 "local_addr={:?} peer_addr={:?} ",
1037 self.local_addr, self.peer_addr,
1038 )?;
1039 write!(
1040 f,
1041 "validation_state={:?} active={} ",
1042 self.validation_state, self.active,
1043 )?;
1044 write!(
1045 f,
1046 "recv={} sent={} lost={} retrans={} rtt={:?} min_rtt={:?} rttvar={:?} cwnd={}",
1047 self.recv, self.sent, self.lost, self.retrans, self.rtt, self.min_rtt, self.rttvar, self.cwnd,
1048 )?;
1049
1050 write!(
1051 f,
1052 " sent_bytes={} recv_bytes={} lost_bytes={}",
1053 self.sent_bytes, self.recv_bytes, self.lost_bytes,
1054 )?;
1055
1056 write!(
1057 f,
1058 " stream_retrans_bytes={} pmtu={} delivery_rate={} rtt_persistent_jump_count={}",
1059 self.stream_retrans_bytes,
1060 self.pmtu,
1061 self.delivery_rate,
1062 self.rtt_persistent_jump_count,
1063 )
1064 }
1065}
1066
1067#[cfg(test)]
1068mod tests {
1069 use crate::rand;
1070 use crate::MIN_CLIENT_INITIAL_LEN;
1071
1072 use crate::recovery::RecoveryConfig;
1073 use crate::Config;
1074
1075 use super::*;
1076
1077 #[test]
1078 fn path_validation_limited_mtu() {
1079 let client_addr = "127.0.0.1:1234".parse().unwrap();
1080 let client_addr_2 = "127.0.0.1:5678".parse().unwrap();
1081 let server_addr = "127.0.0.1:4321".parse().unwrap();
1082
1083 let config = Config::new(crate::PROTOCOL_VERSION).unwrap();
1084 let recovery_config = RecoveryConfig::from_config(&config);
1085
1086 let path = Path::new(
1087 client_addr,
1088 server_addr,
1089 &recovery_config,
1090 config.path_challenge_recv_max_queue_len,
1091 true,
1092 None,
1093 );
1094 let mut path_mgr = PathMap::new(path, 2, false);
1095
1096 let probed_path = Path::new(
1097 client_addr_2,
1098 server_addr,
1099 &recovery_config,
1100 config.path_challenge_recv_max_queue_len,
1101 false,
1102 None,
1103 );
1104 path_mgr.insert_path(probed_path, false).unwrap();
1105
1106 let pid = path_mgr
1107 .path_id_from_addrs(&(client_addr_2, server_addr))
1108 .unwrap();
1109 path_mgr.get_mut(pid).unwrap().request_validation();
1110 assert!(path_mgr.get_mut(pid).unwrap().validation_requested());
1111 assert!(path_mgr.get_mut(pid).unwrap().probing_required());
1112
1113 let data = rand::rand_u64().to_be_bytes();
1116 path_mgr.get_mut(pid).unwrap().add_challenge_sent(
1117 data,
1118 MIN_CLIENT_INITIAL_LEN - 1,
1119 Instant::now(),
1120 );
1121
1122 assert!(!path_mgr.get_mut(pid).unwrap().validation_requested());
1123 assert!(!path_mgr.get_mut(pid).unwrap().probing_required());
1124 assert!(path_mgr.get_mut(pid).unwrap().under_validation());
1125 assert!(!path_mgr.get_mut(pid).unwrap().validated());
1126 assert_eq!(path_mgr.get_mut(pid).unwrap().state, PathState::Validating);
1127 assert_eq!(path_mgr.pop_event(), None);
1128
1129 path_mgr.on_response_received(data).unwrap();
1132
1133 assert!(path_mgr.get_mut(pid).unwrap().validation_requested());
1134 assert!(path_mgr.get_mut(pid).unwrap().probing_required());
1135 assert!(path_mgr.get_mut(pid).unwrap().under_validation());
1136 assert!(!path_mgr.get_mut(pid).unwrap().validated());
1137 assert_eq!(
1138 path_mgr.get_mut(pid).unwrap().state,
1139 PathState::ValidatingMTU
1140 );
1141 assert_eq!(path_mgr.pop_event(), None);
1142
1143 let data = rand::rand_u64().to_be_bytes();
1146 path_mgr.get_mut(pid).unwrap().add_challenge_sent(
1147 data,
1148 MIN_CLIENT_INITIAL_LEN,
1149 Instant::now(),
1150 );
1151
1152 path_mgr.on_response_received(data).unwrap();
1153
1154 assert!(!path_mgr.get_mut(pid).unwrap().validation_requested());
1155 assert!(!path_mgr.get_mut(pid).unwrap().probing_required());
1156 assert!(!path_mgr.get_mut(pid).unwrap().under_validation());
1157 assert!(path_mgr.get_mut(pid).unwrap().validated());
1158 assert_eq!(path_mgr.get_mut(pid).unwrap().state, PathState::Validated);
1159 assert_eq!(
1160 path_mgr.pop_event(),
1161 Some(PathEvent::Validated(client_addr_2, server_addr))
1162 );
1163 }
1164
1165 #[test]
1166 fn multiple_probes() {
1167 let client_addr = "127.0.0.1:1234".parse().unwrap();
1168 let server_addr = "127.0.0.1:4321".parse().unwrap();
1169
1170 let config = Config::new(crate::PROTOCOL_VERSION).unwrap();
1171 let recovery_config = RecoveryConfig::from_config(&config);
1172
1173 let path = Path::new(
1174 client_addr,
1175 server_addr,
1176 &recovery_config,
1177 config.path_challenge_recv_max_queue_len,
1178 true,
1179 None,
1180 );
1181 let mut client_path_mgr = PathMap::new(path, 2, false);
1182 let mut server_path = Path::new(
1183 server_addr,
1184 client_addr,
1185 &recovery_config,
1186 config.path_challenge_recv_max_queue_len,
1187 false,
1188 None,
1189 );
1190
1191 let client_pid = client_path_mgr
1192 .path_id_from_addrs(&(client_addr, server_addr))
1193 .unwrap();
1194
1195 let data = rand::rand_u64().to_be_bytes();
1197
1198 client_path_mgr
1199 .get_mut(client_pid)
1200 .unwrap()
1201 .add_challenge_sent(data, MIN_CLIENT_INITIAL_LEN, Instant::now());
1202
1203 let data_2 = rand::rand_u64().to_be_bytes();
1205
1206 client_path_mgr
1207 .get_mut(client_pid)
1208 .unwrap()
1209 .add_challenge_sent(data_2, MIN_CLIENT_INITIAL_LEN, Instant::now());
1210 assert_eq!(
1211 client_path_mgr
1212 .get(client_pid)
1213 .unwrap()
1214 .in_flight_challenges
1215 .len(),
1216 2
1217 );
1218
1219 server_path.on_challenge_received(data);
1221 assert_eq!(server_path.received_challenges.len(), 1);
1222 server_path.on_challenge_received(data_2);
1223 assert_eq!(server_path.received_challenges.len(), 2);
1224
1225 client_path_mgr.on_response_received(data).unwrap();
1227 assert_eq!(
1228 client_path_mgr
1229 .get(client_pid)
1230 .unwrap()
1231 .in_flight_challenges
1232 .len(),
1233 1
1234 );
1235
1236 client_path_mgr.on_response_received(data_2).unwrap();
1238 assert_eq!(
1239 client_path_mgr
1240 .get(client_pid)
1241 .unwrap()
1242 .in_flight_challenges
1243 .len(),
1244 0
1245 );
1246 }
1247
1248 #[test]
1249 fn too_many_probes() {
1250 let client_addr = "127.0.0.1:1234".parse().unwrap();
1251 let server_addr = "127.0.0.1:4321".parse().unwrap();
1252
1253 let config = Config::new(crate::PROTOCOL_VERSION).unwrap();
1255 let recovery_config = RecoveryConfig::from_config(&config);
1256
1257 let path = Path::new(
1258 client_addr,
1259 server_addr,
1260 &recovery_config,
1261 config.path_challenge_recv_max_queue_len,
1262 true,
1263 None,
1264 );
1265 let mut client_path_mgr = PathMap::new(path, 2, false);
1266 let mut server_path = Path::new(
1267 server_addr,
1268 client_addr,
1269 &recovery_config,
1270 config.path_challenge_recv_max_queue_len,
1271 false,
1272 None,
1273 );
1274
1275 let client_pid = client_path_mgr
1276 .path_id_from_addrs(&(client_addr, server_addr))
1277 .unwrap();
1278
1279 let data = rand::rand_u64().to_be_bytes();
1281
1282 client_path_mgr
1283 .get_mut(client_pid)
1284 .unwrap()
1285 .add_challenge_sent(data, MIN_CLIENT_INITIAL_LEN, Instant::now());
1286
1287 let data_2 = rand::rand_u64().to_be_bytes();
1289
1290 client_path_mgr
1291 .get_mut(client_pid)
1292 .unwrap()
1293 .add_challenge_sent(data_2, MIN_CLIENT_INITIAL_LEN, Instant::now());
1294 assert_eq!(
1295 client_path_mgr
1296 .get(client_pid)
1297 .unwrap()
1298 .in_flight_challenges
1299 .len(),
1300 2
1301 );
1302
1303 let data_3 = rand::rand_u64().to_be_bytes();
1305
1306 client_path_mgr
1307 .get_mut(client_pid)
1308 .unwrap()
1309 .add_challenge_sent(data_3, MIN_CLIENT_INITIAL_LEN, Instant::now());
1310 assert_eq!(
1311 client_path_mgr
1312 .get(client_pid)
1313 .unwrap()
1314 .in_flight_challenges
1315 .len(),
1316 3
1317 );
1318
1319 let data_4 = rand::rand_u64().to_be_bytes();
1321
1322 client_path_mgr
1323 .get_mut(client_pid)
1324 .unwrap()
1325 .add_challenge_sent(data_4, MIN_CLIENT_INITIAL_LEN, Instant::now());
1326 assert_eq!(
1327 client_path_mgr
1328 .get(client_pid)
1329 .unwrap()
1330 .in_flight_challenges
1331 .len(),
1332 4
1333 );
1334
1335 server_path.on_challenge_received(data);
1338 assert_eq!(server_path.received_challenges.len(), 1);
1339 server_path.on_challenge_received(data_2);
1340 assert_eq!(server_path.received_challenges.len(), 2);
1341 server_path.on_challenge_received(data_3);
1342 assert_eq!(server_path.received_challenges.len(), 3);
1343 server_path.on_challenge_received(data_4);
1344 assert_eq!(server_path.received_challenges.len(), 3);
1345
1346 client_path_mgr.on_response_received(data).unwrap();
1348 assert_eq!(
1349 client_path_mgr
1350 .get(client_pid)
1351 .unwrap()
1352 .in_flight_challenges
1353 .len(),
1354 3
1355 );
1356
1357 client_path_mgr.on_response_received(data_2).unwrap();
1359 assert_eq!(
1360 client_path_mgr
1361 .get(client_pid)
1362 .unwrap()
1363 .in_flight_challenges
1364 .len(),
1365 2
1366 );
1367
1368 client_path_mgr.on_response_received(data_3).unwrap();
1370 assert_eq!(
1371 client_path_mgr
1372 .get(client_pid)
1373 .unwrap()
1374 .in_flight_challenges
1375 .len(),
1376 1
1377 );
1378
1379 }
1381}