1use std::sync::Arc;
2use std::sync::atomic::{AtomicU32, Ordering};
3use tracing::debug;
4
5use crate::constants;
6use crate::selector::server_stat_man::ServerStatMan;
7use crate::selector::uri_selector::UriSelector;
8
9#[derive(Debug, Clone, PartialEq)]
10pub enum SegmentStatus {
11 Pending,
12 Downloading,
13 Done,
14 Failed,
15}
16
17#[derive(Debug)]
18pub struct Segment {
19 pub index: u32,
20 pub offset: u64,
21 pub length: u64,
22 pub status: SegmentStatus,
23 pub data: Option<bytes::Bytes>, pub assigned_mirror: Option<usize>,
25 pub retry_count: u32,
26}
27
28impl Segment {
29 fn new(index: u32, offset: u64, length: u64) -> Self {
30 Self {
31 index,
32 offset,
33 length,
34 status: SegmentStatus::Pending,
35 data: None,
36 assigned_mirror: None,
37 retry_count: 0,
38 }
39 }
40}
41
42#[derive(Debug)]
43pub struct MirrorState {
44 pub url: String,
45 pub speed: f64,
46 pub active_segments: usize,
47 pub max_connections: usize,
48 pub consecutive_failures: usize,
49 pub disabled: bool,
50}
51
52impl MirrorState {
53 fn new(url: String) -> Self {
54 Self {
55 url,
56 speed: 0.0,
57 active_segments: 0,
58 max_connections: constants::DEFAULT_MAX_CONNECTIONS_PER_MIRROR,
59 consecutive_failures: 0,
60 disabled: false,
61 }
62 }
63
64 pub fn is_available(&self) -> bool {
65 !self.disabled && self.active_segments < self.max_connections
66 }
67
68 pub fn can_accept_more(&self) -> bool {
69 !self.disabled && self.active_segments < self.max_connections
70 }
71}
72
73pub struct ConcurrentSegmentManager {
74 total_size: u64,
75 segments: Vec<Segment>,
76 mirrors: Vec<MirrorState>,
77 mirror_urls: Vec<String>,
79 completed_bytes: u64,
80 max_retries_per_segment: u32,
81 max_mirror_failures: usize,
82 stat_man: Option<Arc<ServerStatMan>>,
84 uri_selector: Option<Box<dyn UriSelector>>,
86 next_segment_idx: AtomicU32,
94}
95
96impl ConcurrentSegmentManager {
97 pub fn new(total_size: u64, urls: Vec<String>, segment_size: Option<u64>) -> Self {
103 let seg_size = segment_size.unwrap_or(constants::DEFAULT_SEGMENT_SIZE as u64);
104 let num_segments = if total_size == 0 {
105 0
106 } else {
107 total_size.div_ceil(seg_size) as usize
108 };
109
110 let mut segments = Vec::with_capacity(num_segments);
111 for i in 0..num_segments {
112 let offset = (i as u64) * seg_size;
113 let remaining = total_size.saturating_sub(offset);
114 let length = seg_size.min(remaining);
115 segments.push(Segment::new(i as u32, offset, length));
116 }
117
118 let mirrors = urls.iter().cloned().map(MirrorState::new).collect();
119
120 Self {
121 total_size,
122 segments,
123 mirrors,
124 mirror_urls: urls,
125 completed_bytes: 0,
126 max_retries_per_segment: constants::MAX_RETRIES_PER_SEGMENT,
127 max_mirror_failures: constants::MAX_MIRROR_FAILURES as usize,
128 stat_man: None,
129 uri_selector: None,
130 next_segment_idx: AtomicU32::new(0),
131 }
132 }
133
134 pub fn new_with_selector(
168 total_size: u64,
169 urls: Vec<String>,
170 segment_size: Option<u64>,
171 stat_man: Arc<ServerStatMan>,
172 uri_selector: Box<dyn UriSelector>,
173 ) -> Self {
174 let seg_size = segment_size.unwrap_or(constants::DEFAULT_SEGMENT_SIZE as u64);
175 let num_segments = if total_size == 0 {
176 0
177 } else {
178 total_size.div_ceil(seg_size) as usize
179 };
180
181 let mut segments = Vec::with_capacity(num_segments);
182 for i in 0..num_segments {
183 let offset = (i as u64) * seg_size;
184 let remaining = total_size.saturating_sub(offset);
185 let length = seg_size.min(remaining);
186 segments.push(Segment::new(i as u32, offset, length));
187 }
188
189 let mirrors = urls.iter().cloned().map(MirrorState::new).collect();
190
191 Self {
192 total_size,
193 segments,
194 mirrors,
195 mirror_urls: urls,
196 completed_bytes: 0,
197 max_retries_per_segment: constants::MAX_RETRIES_PER_SEGMENT,
198 max_mirror_failures: constants::MAX_MIRROR_FAILURES as usize,
199 stat_man: Some(stat_man),
200 uri_selector: Some(uri_selector),
201 next_segment_idx: AtomicU32::new(0),
202 }
203 }
204
205 pub fn allocate_segments(&mut self) {
206 for mirror_idx in 0..self.mirrors.len() {
207 while self.mirrors[mirror_idx].can_accept_more() {
208 if let Some(seg) = self.find_pending_segment() {
209 seg.status = SegmentStatus::Downloading;
210 seg.assigned_mirror = Some(mirror_idx);
211 self.mirrors[mirror_idx].active_segments += 1;
212 } else {
213 break;
214 }
215 }
216 }
217 }
218
219 fn find_pending_segment(&mut self) -> Option<&mut Segment> {
220 let len = self.segments.len();
221 if len == 0 {
222 return None;
223 }
224 let start = (self.next_segment_idx.load(Ordering::Relaxed) as usize) % len;
227 for i in 0..len {
228 let idx = (start + i) % len;
229 if self.segments[idx].status == SegmentStatus::Pending {
230 self.next_segment_idx
233 .store((idx + 1) as u32, Ordering::Relaxed);
234 return Some(&mut self.segments[idx]);
235 }
236 }
237 None
238 }
239
240 pub fn next_pending_segment_for_mirror(
241 &mut self,
242 mirror_idx: usize,
243 ) -> Option<(u32, u64, u64)> {
244 if !self
245 .mirrors
246 .get(mirror_idx)
247 .is_some_and(|m| m.can_accept_more())
248 {
249 return None;
250 }
251
252 let len = self.segments.len();
253 if len == 0 {
254 return None;
255 }
256 let start = (self.next_segment_idx.load(Ordering::Relaxed) as usize) % len;
260 for i in 0..len {
261 let idx = (start + i) % len;
262 let seg = &mut self.segments[idx];
263 if seg.status == SegmentStatus::Pending {
264 seg.status = SegmentStatus::Downloading;
265 seg.assigned_mirror = Some(mirror_idx);
266 self.next_segment_idx
267 .store((idx + 1) as u32, Ordering::Relaxed);
268 if let Some(m) = self.mirrors.get_mut(mirror_idx) {
269 m.active_segments += 1;
270 }
271 return Some((seg.index, seg.offset, seg.length));
272 }
273 }
274 None
275 }
276
277 pub fn next_pending_segment(&mut self) -> Option<(u32, u64, u64)> {
278 self.next_pending_segment_for_mirror(0)
279 }
280
281 pub fn allocate_next_index(&self) -> Option<u32> {
297 let idx = self.next_segment_idx.fetch_add(1, Ordering::Relaxed);
298 if (idx as usize) < self.segments.len() {
299 Some(idx)
300 } else {
301 None
302 }
303 }
304
305 pub fn reset_allocation_index(&self) {
316 self.next_segment_idx.store(0, Ordering::Relaxed);
317 }
318
319 pub fn complete_segment(&mut self, index: u32, data: bytes::Bytes) -> bool {
320 if let Some(seg) = self.segments.get_mut(index as usize) {
321 seg.status = SegmentStatus::Done;
322 seg.data = Some(data);
323
324 if let Some(mi) = seg.assigned_mirror
325 && let Some(m) = self.mirrors.get_mut(mi)
326 {
327 m.active_segments = m.active_segments.saturating_sub(1);
328 m.consecutive_failures = 0;
329 }
330
331 self.completed_bytes += seg.length;
332 true
333 } else {
334 false
335 }
336 }
337
338 pub fn fail_segment(&mut self, index: u32) -> Option<usize> {
339 let (prev_mirror, new_retry) = {
340 let seg = self.segments.get(index as usize)?;
341 (seg.assigned_mirror, seg.retry_count + 1)
342 };
343
344 if let Some(mi) = prev_mirror
345 && let Some(m) = self.mirrors.get_mut(mi)
346 {
347 m.active_segments = m.active_segments.saturating_sub(1);
348 m.consecutive_failures += 1;
349 if m.consecutive_failures >= self.max_mirror_failures {
350 m.disabled = true;
351 }
352 }
353
354 if new_retry >= self.max_retries_per_segment {
355 if let Some(seg) = self.segments.get_mut(index as usize) {
356 seg.status = SegmentStatus::Failed;
357 seg.retry_count = new_retry;
358 }
359 None
360 } else {
361 let reassign = self.find_available_mirror_for_reassignment(prev_mirror.unwrap_or(0));
362 if let Some(seg) = self.segments.get_mut(index as usize) {
363 seg.status = SegmentStatus::Pending;
364 seg.assigned_mirror = reassign;
365 seg.retry_count = new_retry;
366 }
367 reassign
368 }
369 }
370
371 fn find_available_mirror_for_reassignment(&self, exclude: usize) -> Option<usize> {
372 self.mirrors
373 .iter()
374 .enumerate()
375 .filter(|(i, m)| *i != exclude && m.is_available())
376 .map(|(i, _)| i)
377 .next()
378 }
379
380 pub fn is_complete(&self) -> bool {
381 self.segments
382 .iter()
383 .all(|s| s.status == SegmentStatus::Done)
384 }
385
386 pub fn has_failed_segments(&self) -> bool {
387 self.segments
388 .iter()
389 .any(|s| s.status == SegmentStatus::Failed)
390 }
391
392 pub fn has_pending_segments(&self) -> bool {
393 self.segments
394 .iter()
395 .any(|s| s.status == SegmentStatus::Pending)
396 }
397
398 pub fn completed_ranges(&self) -> Vec<(u64, u64)> {
399 let mut ranges = Vec::new();
400 for seg in &self.segments {
401 if seg.status == SegmentStatus::Done {
402 ranges.push((seg.offset, seg.length));
403 }
404 }
405 ranges.sort_by_key(|r| r.0);
406 ranges
407 }
408
409 pub fn assemble(&self) -> Option<Vec<u8>> {
410 if !self.is_complete() || self.total_size == 0 {
411 return None;
412 }
413
414 let mut result = Vec::with_capacity(self.total_size as usize);
415 for seg in &self.segments {
416 let data = seg.data.as_ref()?;
417 result.extend_from_slice(data);
418 }
419 Some(result)
420 }
421
422 pub fn progress(&self) -> f64 {
423 if self.total_size == 0 {
424 return 100.0;
425 }
426 let done = self
427 .segments
428 .iter()
429 .filter(|s| s.status == SegmentStatus::Done)
430 .count();
431 done as f64 / self.segments.len() as f64 * 100.0
432 }
433
434 pub fn num_segments(&self) -> usize {
435 self.segments.len()
436 }
437 pub fn segment_status(&self, index: usize) -> Option<SegmentStatus> {
438 self.segments.get(index).map(|s| s.status.clone())
439 }
440 pub fn num_mirrors(&self) -> usize {
441 self.mirrors.len()
442 }
443 pub fn total_size(&self) -> u64 {
444 self.total_size
445 }
446 pub fn completed_bytes(&self) -> u64 {
447 self.completed_bytes
448 }
449
450 pub fn mirror_url(&self, index: usize) -> Option<&str> {
451 self.mirrors.get(index).map(|m| m.url.as_str())
452 }
453
454 pub fn available_mirrors(&self) -> Vec<usize> {
455 self.mirrors
456 .iter()
457 .enumerate()
458 .filter(|(_, m)| m.is_available())
459 .map(|(i, _)| i)
460 .collect()
461 }
462
463 pub fn any_mirror_available(&self) -> bool {
464 self.mirrors.iter().any(|m| m.is_available())
465 }
466
467 pub fn set_max_connections_per_mirror(&mut self, max: usize) {
468 for m in &mut self.mirrors {
469 m.max_connections = max;
470 }
471 }
472
473 pub fn set_mirror_max_connections(&mut self, mirror_idx: usize, max: usize) {
482 if let Some(mirror) = self.mirrors.get_mut(mirror_idx) {
483 mirror.max_connections = max;
484 }
485 }
486
487 pub fn get_mirror_max_connections(&self, mirror_idx: usize) -> Option<usize> {
493 self.mirrors.get(mirror_idx).map(|m| m.max_connections)
494 }
495
496 pub fn set_max_retries(&mut self, retries: u32) {
497 self.max_retries_per_segment = retries;
498 }
499
500 pub fn segment_retry_count(&self, seg_idx: u32) -> u32 {
501 self.segments
502 .iter()
503 .find(|s| s.index == seg_idx)
504 .map(|s| s.retry_count)
505 .unwrap_or(0)
506 }
507
508 pub fn has_permanently_failed_segments(&self) -> bool {
509 self.segments.iter().any(|s| {
510 s.status == SegmentStatus::Failed && s.retry_count >= self.max_retries_per_segment
511 })
512 }
513
514 pub fn mark_completed_up_to(&mut self, offset: u64, length: u64) {
515 let end_offset = offset + length;
516 for segment in &mut self.segments {
517 if segment.offset + segment.length <= offset {
518 if segment.status != SegmentStatus::Done {
519 segment.status = SegmentStatus::Done;
520 self.completed_bytes += segment.length;
521 }
522 } else if segment.offset < end_offset {
523 let overlap_start = std::cmp::max(segment.offset, offset);
524 let overlap_end = std::cmp::min(segment.offset + segment.length, end_offset);
525 if overlap_end > overlap_start {
526 debug!(
527 "段 {} 部分已完成: {}/{} bytes",
528 segment.index,
529 overlap_end - segment.offset,
530 segment.length
531 );
532 }
533 }
534 }
535 }
536
537 pub fn segment_info(&self, index: usize) -> Option<(u64, u64, &SegmentStatus)> {
538 self.segments
539 .get(index)
540 .map(|s| (s.offset, s.length, &s.status))
541 }
542
543 pub fn select_mirror_for_next_segment(&mut self) -> Option<(usize, (u32, u64, u64))> {
556 let pending_seg = self
558 .segments
559 .iter()
560 .find(|s| s.status == SegmentStatus::Pending)?;
561
562 let seg_index = pending_seg.index;
563
564 if let Some(ref selector) = self.uri_selector {
566 let used_hosts: Vec<(usize, String)> = self
568 .segments
569 .iter()
570 .filter(|s| s.status == SegmentStatus::Downloading)
571 .filter_map(|s| {
572 s.assigned_mirror.and_then(|idx| {
573 self.mirror_urls.get(idx).map(|url| {
574 let host = extract_host_from_url(url);
576 (idx, host)
577 })
578 })
579 })
580 .collect();
581
582 if let Some(mirror_idx) = selector.select(&self.mirror_urls, &used_hosts) {
584 if self
586 .mirrors
587 .get(mirror_idx)
588 .is_some_and(|m| m.can_accept_more())
589 {
590 if let Some(seg) = self.segments.get_mut(seg_index as usize) {
592 seg.status = SegmentStatus::Downloading;
593 seg.assigned_mirror = Some(mirror_idx);
594 if let Some(m) = self.mirrors.get_mut(mirror_idx) {
595 m.active_segments += 1;
596 }
597 return Some((mirror_idx, (seg.index, seg.offset, seg.length)));
598 }
599 }
600 }
601 }
602
603 for mirror_idx in 0..self.mirrors.len() {
605 if self.mirrors[mirror_idx].can_accept_more()
606 && let Some(seg) = self.segments.get_mut(seg_index as usize)
607 {
608 seg.status = SegmentStatus::Downloading;
609 seg.assigned_mirror = Some(mirror_idx);
610 self.mirrors[mirror_idx].active_segments += 1;
611 return Some((mirror_idx, (seg.index, seg.offset, seg.length)));
612 }
613 }
614
615 None
616 }
617
618 pub fn report_segment_complete(
635 &mut self,
636 seg_idx: u32,
637 data: bytes::Bytes,
638 bytes_per_sec: u64,
639 is_multi_connection: bool,
640 ) -> bool {
641 let mirror_idx = self
643 .segments
644 .get(seg_idx as usize)
645 .and_then(|s| s.assigned_mirror);
646
647 let success = self.complete_segment(seg_idx, data);
649
650 if success {
652 if let (Some(idx), Some(stat_man)) = (mirror_idx, &self.stat_man)
653 && let Some(url) = self.mirror_urls.get(idx)
654 {
655 let host = extract_host_from_url(url);
656 stat_man.update(&host, bytes_per_sec, is_multi_connection);
657
658 if let Some(stat) = stat_man.find_stat(&host) {
660 stat.reset_status();
661 }
662 }
663
664 if let Some(ref selector) = self.uri_selector {
666 selector.tune_command(&self.mirror_urls, bytes_per_sec);
667 }
668 }
669
670 success
671 }
672
673 pub fn report_segment_failed(&mut self, seg_idx: u32, error_code: u16) -> Option<usize> {
689 let mirror_idx = self
691 .segments
692 .get(seg_idx as usize)
693 .and_then(|s| s.assigned_mirror);
694
695 let reassign = self.fail_segment(seg_idx);
697
698 if let (Some(idx), Some(stat_man)) = (mirror_idx, &self.stat_man)
700 && let Some(url) = self.mirror_urls.get(idx)
701 {
702 let host = extract_host_from_url(url);
703 stat_man.get_or_create(&host);
705 stat_man.mark_failure(&host, error_code);
706
707 if let Some(stat) = stat_man.find_stat(&host)
709 && !stat.is_available()
710 {
711 if let Some(m) = self.mirrors.get_mut(idx) {
713 m.disabled = true;
714 }
715 }
716 }
717
718 reassign
719 }
720
721 pub fn get_mirror_url(&self, mirror_idx: usize) -> Option<&str> {
723 self.mirror_urls.get(mirror_idx).map(|s| s.as_str())
724 }
725
726 pub fn mirror_active_segments(&self, mirror_idx: usize) -> usize {
728 self.mirrors
729 .get(mirror_idx)
730 .map(|m| m.active_segments)
731 .unwrap_or(0)
732 }
733
734 pub fn has_intelligent_selection(&self) -> bool {
736 self.uri_selector.is_some() && self.stat_man.is_some()
737 }
738}
739
740fn extract_host_from_url(url: &str) -> String {
742 let url = url.trim();
743 if !url.contains("://") {
744 return url.to_string();
745 }
746 let after_scheme = &url[url.find("://").unwrap() + 3..];
747 let host_part = if let Some(slash_idx) = after_scheme.find('/') {
748 &after_scheme[..slash_idx]
749 } else {
750 after_scheme
751 };
752 host_part.to_string()
753}
754
755#[cfg(test)]
756mod tests {
757 use super::*;
758
759 #[test]
760 fn test_manager_creation_small_file() {
761 let mgr = ConcurrentSegmentManager::new(1024, vec!["http://a.com/f".to_string()], None);
762 assert_eq!(mgr.num_segments(), 1);
763 assert_eq!(mgr.num_mirrors(), 1);
764 assert_eq!(mgr.total_size(), 1024);
765 assert!(!mgr.is_complete());
766 assert!(mgr.has_pending_segments());
767 }
768
769 #[test]
770 fn test_manager_large_file_multi_segment() {
771 let mgr = ConcurrentSegmentManager::new(
772 3_000_000,
773 vec!["http://a.com/f".to_string(), "http://b.com/f".to_string()],
774 Some(1_000_000),
775 );
776 assert_eq!(mgr.num_segments(), 3);
777 assert_eq!(mgr.num_mirrors(), 2);
778 }
779
780 #[test]
781 fn test_allocate_segments_round_robin() {
782 let mut mgr = ConcurrentSegmentManager::new(
783 3_000_000,
784 vec!["http://a.com/f".to_string(), "http://b.com/f".to_string()],
785 Some(1_000_000),
786 );
787
788 mgr.allocate_segments();
789
790 let assigned_a: Vec<_> = mgr
791 .segments
792 .iter()
793 .filter(|s| s.assigned_mirror == Some(0))
794 .map(|s| s.index)
795 .collect();
796 let assigned_b: Vec<_> = mgr
797 .segments
798 .iter()
799 .filter(|s| s.assigned_mirror == Some(1))
800 .map(|s| s.index)
801 .collect();
802
803 assert!(!assigned_a.is_empty());
804 assert!(!assigned_b.is_empty());
805 assert_eq!(assigned_a.len() + assigned_b.len(), 3);
806 }
807
808 #[test]
809 fn test_complete_and_assemble() {
810 let mut mgr =
811 ConcurrentSegmentManager::new(200, vec!["http://x.com/f".to_string()], Some(100));
812
813 mgr.allocate_segments();
814 assert_eq!(mgr.progress(), 0.0);
815
816 mgr.complete_segment(0, bytes::Bytes::from(vec![0xAB; 100]));
817 assert!(!mgr.is_complete());
818 assert!((mgr.progress() - 50.0).abs() < 0.01);
819
820 mgr.complete_segment(1, bytes::Bytes::from(vec![0xCD; 100]));
821 assert!(mgr.is_complete());
822 assert!((mgr.progress() - 100.0).abs() < 0.01);
823
824 let assembled = mgr.assemble().unwrap();
825 assert_eq!(assembled.len(), 200);
826 assert_eq!(&assembled[..100], &[0xAB; 100][..]);
827 assert_eq!(&assembled[100..], &[0xCD; 100][..]);
828 }
829
830 #[test]
831 fn test_fail_and_reassign() {
832 let mut mgr = ConcurrentSegmentManager::new(
833 200,
834 vec!["http://a.com/f".to_string(), "http://b.com/f".to_string()],
835 Some(100),
836 );
837
838 mgr.allocate_segments();
839
840 let reassign = mgr.fail_segment(0);
841 assert!(reassign.is_some());
842
843 let seg = &mgr.segments[0];
844 assert_eq!(seg.status, SegmentStatus::Pending);
845 assert_eq!(seg.assigned_mirror, reassign);
846 assert_eq!(seg.retry_count, 1);
847 }
848
849 #[test]
850 fn test_max_retries_exhausted() {
851 let mut mgr =
852 ConcurrentSegmentManager::new(100, vec!["http://a.com/f".to_string()], Some(100));
853 mgr.set_max_retries(2);
854
855 mgr.fail_segment(0);
856 assert!(mgr.has_pending_segments());
857
858 mgr.fail_segment(0);
859 assert!(mgr.has_failed_segments());
860 assert!(!mgr.has_pending_segments());
861 }
862
863 #[test]
864 fn test_empty_file() {
865 let mgr = ConcurrentSegmentManager::new(0, vec!["http://x.com/f".to_string()], None);
866 assert_eq!(mgr.num_segments(), 0);
867 assert!(mgr.is_complete());
868 assert!(mgr.assemble().is_none());
869 }
870
871 #[test]
872 fn test_next_pending_for_specific_mirror() {
873 let mut mgr = ConcurrentSegmentManager::new(
874 300,
875 vec!["http://a.com/f".to_string(), "http://b.com/f".to_string()],
876 Some(100),
877 );
878
879 let r = mgr.next_pending_segment_for_mirror(0);
880 assert!(r.is_some());
881 let (idx, off, len) = r.unwrap();
882 assert_eq!(idx, 0);
883 assert_eq!(off, 0);
884 assert_eq!(len, 100);
885
886 let r2 = mgr.next_pending_segment_for_mirror(1);
887 assert!(r2.is_some());
888 let (idx2, _, _) = r2.unwrap();
889 assert_eq!(idx2, 1);
890 }
891
892 #[test]
897 fn test_new_with_selector() {
898 use crate::selector::adaptive_uri_selector::AdaptiveUriSelector;
899
900 let stat_man = Arc::new(ServerStatMan::new());
901 let urls = vec![
902 "http://mirror1.com/file".to_string(),
903 "http://mirror2.com/file".to_string(),
904 ];
905 let selector = Box::new(AdaptiveUriSelector::new_with_uris(
906 Arc::clone(&stat_man),
907 urls.clone(),
908 ));
909
910 let mgr = ConcurrentSegmentManager::new_with_selector(
911 1_000_000,
912 urls,
913 Some(500_000),
914 stat_man,
915 selector,
916 );
917
918 assert_eq!(mgr.num_segments(), 2);
919 assert_eq!(mgr.num_mirrors(), 2);
920 assert!(mgr.has_intelligent_selection());
921 }
922
923 #[test]
924 fn test_select_mirror_for_next_segment_without_selector() {
925 let mut mgr = ConcurrentSegmentManager::new(
926 300,
927 vec!["http://a.com/f".to_string(), "http://b.com/f".to_string()],
928 Some(100),
929 );
930
931 let result = mgr.select_mirror_for_next_segment();
933 assert!(result.is_some());
934
935 let (mirror_idx, (seg_idx, offset, len)) = result.unwrap();
936 assert_eq!(seg_idx, 0);
937 assert_eq!(offset, 0);
938 assert_eq!(len, 100);
939 assert!(mirror_idx < 2);
940 }
941
942 #[test]
943 fn test_select_mirror_for_next_segment_with_selector() {
944 use crate::selector::adaptive_uri_selector::AdaptiveUriSelector;
945
946 let stat_man = Arc::new(ServerStatMan::new());
947 let urls = vec![
948 "http://fast.com/f".to_string(),
949 "http://slow.com/f".to_string(),
950 ];
951
952 stat_man.update("fast.com", 1_000_000, false);
954 stat_man.update("slow.com", 1000, false);
955 let fast_stat = stat_man.find_stat("fast.com").unwrap();
956 fast_stat.increment_counter();
957 let slow_stat = stat_man.find_stat("slow.com").unwrap();
958 slow_stat.increment_counter();
959
960 let selector = Box::new(AdaptiveUriSelector::new_with_uris(
961 Arc::clone(&stat_man),
962 urls.clone(),
963 ));
964
965 let mut mgr =
966 ConcurrentSegmentManager::new_with_selector(300, urls, Some(100), stat_man, selector);
967
968 let result = mgr.select_mirror_for_next_segment();
969 assert!(result.is_some());
970
971 let (mirror_idx, _) = result.unwrap();
972 assert_eq!(mirror_idx, 0);
974 }
975
976 #[test]
977 fn test_report_segment_complete_updates_stats() {
978 use crate::selector::adaptive_uri_selector::AdaptiveUriSelector;
979
980 let stat_man = Arc::new(ServerStatMan::new());
981 let urls = vec!["http://test.mirror.com/f".to_string()];
982
983 let selector = Box::new(AdaptiveUriSelector::new_with_uris(
984 Arc::clone(&stat_man),
985 urls.clone(),
986 ));
987
988 let mut mgr = ConcurrentSegmentManager::new_with_selector(
989 100,
990 urls,
991 Some(100),
992 stat_man.clone(),
993 selector,
994 );
995
996 mgr.allocate_segments();
997
998 let success =
1000 mgr.report_segment_complete(0, bytes::Bytes::from(vec![0xAB; 100]), 1_000_000, false);
1001 assert!(success);
1002
1003 let stat = stat_man.find_stat("test.mirror.com").unwrap();
1005 assert!(stat.get_download_speed() > 0);
1006 }
1007
1008 #[test]
1009 fn test_report_segment_failed_updates_stats() {
1010 use crate::selector::adaptive_uri_selector::AdaptiveUriSelector;
1011
1012 let stat_man = Arc::new(ServerStatMan::new());
1013 let urls = vec![
1014 "http://failing.mirror.com/f".to_string(),
1015 "http://backup.mirror.com/f".to_string(),
1016 ];
1017
1018 let selector = Box::new(AdaptiveUriSelector::new_with_uris(
1019 Arc::clone(&stat_man),
1020 urls.clone(),
1021 ));
1022
1023 let mut mgr = ConcurrentSegmentManager::new_with_selector(
1024 100,
1025 urls,
1026 Some(100),
1027 stat_man.clone(),
1028 selector,
1029 );
1030
1031 mgr.allocate_segments();
1032
1033 let reassign = mgr.report_segment_failed(0, 503);
1035 assert!(reassign.is_some());
1036
1037 let stat = stat_man.find_stat("failing.mirror.com").unwrap();
1039 assert_eq!(stat.get_consecutive_failures(), 1);
1040 assert_eq!(stat.get_last_error_code(), 503);
1041 }
1042
1043 #[test]
1044 fn test_extract_host_from_url() {
1045 assert_eq!(
1046 extract_host_from_url("http://example.com/path"),
1047 "example.com"
1048 );
1049 assert_eq!(
1050 extract_host_from_url("https://host:8080/file?q=1"),
1051 "host:8080"
1052 );
1053 assert_eq!(extract_host_from_url("ftp://server.com"), "server.com");
1054 assert_eq!(extract_host_from_url("not-a-url"), "not-a-url");
1055 }
1056
1057 #[test]
1058 fn test_get_mirror_url() {
1059 let mgr = ConcurrentSegmentManager::new(
1060 100,
1061 vec!["http://a.com/f".to_string(), "http://b.com/f".to_string()],
1062 Some(100),
1063 );
1064
1065 assert_eq!(mgr.get_mirror_url(0), Some("http://a.com/f"));
1066 assert_eq!(mgr.get_mirror_url(1), Some("http://b.com/f"));
1067 assert_eq!(mgr.get_mirror_url(999), None);
1068 }
1069
1070 #[test]
1071 fn test_mirror_active_segments() {
1072 let mut mgr =
1073 ConcurrentSegmentManager::new(300, vec!["http://a.com/f".to_string()], Some(100));
1074
1075 assert_eq!(mgr.mirror_active_segments(0), 0);
1076 assert_eq!(mgr.num_segments(), 3);
1077
1078 mgr.set_max_connections_per_mirror(3);
1080
1081 mgr.allocate_segments();
1082 assert_eq!(mgr.mirror_active_segments(0), 3);
1084 }
1085
1086 #[test]
1087 fn test_no_intelligent_selection_by_default() {
1088 let mgr = ConcurrentSegmentManager::new(100, vec!["http://a.com/f".to_string()], Some(100));
1089
1090 assert!(!mgr.has_intelligent_selection());
1091 }
1092
1093 #[tokio::test(flavor = "multi_thread", worker_threads = 16)]
1105 async fn test_segment_allocation_is_lock_free() {
1106 use std::collections::HashSet;
1107
1108 let manager = Arc::new(ConcurrentSegmentManager::new(
1110 16000,
1111 vec!["http://test".into()],
1112 Some(1),
1113 ));
1114 assert_eq!(manager.num_segments(), 16000);
1115
1116 let collected: Arc<tokio::sync::Mutex<Vec<u32>>> =
1117 Arc::new(tokio::sync::Mutex::new(Vec::new()));
1118
1119 let mut handles = Vec::with_capacity(16);
1120 for _ in 0..16 {
1121 let m = manager.clone();
1122 let c = collected.clone();
1123 handles.push(tokio::spawn(async move {
1124 let mut local = Vec::with_capacity(1000);
1126 for _ in 0..1000 {
1127 if let Some(idx) = m.allocate_next_index() {
1128 local.push(idx);
1129 }
1130 }
1131 c.lock().await.extend(local);
1132 }));
1133 }
1134
1135 for h in handles {
1136 h.await.unwrap();
1137 }
1138
1139 let indices = collected.lock().await.clone();
1140 assert_eq!(
1141 indices.len(),
1142 16000,
1143 "should have allocated all 16000 segments"
1144 );
1145
1146 let set: HashSet<u32> = indices.iter().copied().collect();
1148 assert_eq!(
1149 set.len(),
1150 16000,
1151 "all indices must be unique (no duplicates)"
1152 );
1153
1154 for i in 0..16000u32 {
1156 assert!(set.contains(&i), "missing index {}", i);
1157 }
1158 }
1159
1160 #[test]
1168 fn test_allocation_hint_advances_and_resets() {
1169 let mut mgr =
1170 ConcurrentSegmentManager::new(500, vec!["http://a.com/f".to_string()], Some(100));
1171 mgr.set_max_connections_per_mirror(10);
1173 assert_eq!(mgr.num_segments(), 5);
1174
1175 let mut claimed = Vec::new();
1179 while let Some((idx, _, _)) = mgr.next_pending_segment_for_mirror(0) {
1180 claimed.push(idx);
1181 }
1182 assert_eq!(claimed, vec![0, 1, 2, 3, 4]);
1183
1184 assert!(mgr.next_pending_segment_for_mirror(0).is_none());
1186
1187 mgr.segments[1].status = SegmentStatus::Pending;
1191 let next = mgr.next_pending_segment_for_mirror(0);
1192 assert!(next.is_some());
1193 assert_eq!(next.unwrap().0, 1);
1194
1195 mgr.reset_allocation_index();
1197 mgr.segments[3].status = SegmentStatus::Pending;
1200 let next = mgr.next_pending_segment_for_mirror(0);
1201 assert!(next.is_some());
1202 assert_eq!(next.unwrap().0, 3);
1203 }
1204}