1use std::path::{Path, PathBuf};
2use std::sync::{Arc, Mutex as StdMutex, Weak};
3
4use fs2::FileExt;
5use reqwest::header::HeaderMap;
6use reqwest::Method;
7use tokio::io::{AsyncSeekExt, AsyncWriteExt};
8use tokio::sync::Mutex;
9
10use crate::direction::Direction;
11use crate::http_breakpoint::{BreakpointDownload, BreakpointDownloadHttpConfig, BreakpointUpload};
12use crate::inner::inner_task::InnerTask;
13use crate::upload_file::UploadFileSnapshot;
14use crate::upload_source::UploadSource;
15
16fn map_download_validation_error(error: std::io::Error) -> crate::MeowError {
17 match error.kind() {
18 std::io::ErrorKind::NotFound | std::io::ErrorKind::UnexpectedEof => {
19 crate::MeowError::from_source(
20 crate::InnerErrorCode::LocalFileRemoved,
21 "completed download disappeared or was truncated".to_owned(),
22 error,
23 )
24 }
25 std::io::ErrorKind::InvalidData => crate::MeowError::from_source(
26 crate::InnerErrorCode::ChecksumMismatch,
27 "validate completed download content failed".to_owned(),
28 error,
29 ),
30 _ => crate::MeowError::from_io(
31 "read completed download for validation failed".to_owned(),
32 error,
33 ),
34 }
35}
36
37fn map_download_target_lock_error(
45 message: impl Into<String>,
46 error: std::io::Error,
47) -> crate::MeowError {
48 let message = message.into();
49 if error.kind() == std::io::ErrorKind::WouldBlock {
50 crate::MeowError::from_source(crate::InnerErrorCode::InvalidTaskState, message, error)
51 } else {
52 crate::MeowError::from_io(message, error)
53 }
54}
55
56fn arm_download_checkpoint_timer(
57 progress: Weak<StdMutex<Option<crate::dflt::download_progress::DownloadProgress>>>,
58 file_slot: Weak<Mutex<Option<tokio::fs::File>>>,
59 barrier: Weak<Mutex<()>>,
60) -> tokio::task::JoinHandle<()> {
61 tokio::spawn(async move {
62 tokio::time::sleep(crate::dflt::download_progress::DEFAULT_CHECKPOINT_INTERVAL).await;
63 let (Some(progress), Some(file_slot), Some(barrier)) =
64 (progress.upgrade(), file_slot.upgrade(), barrier.upgrade())
65 else {
66 return;
69 };
70 let result = async {
71 let _barrier = barrier.lock().await;
72 let begin_progress = Arc::clone(&progress);
73 let checkpoint_due = tokio::task::spawn_blocking(move || {
74 let mut guard = begin_progress.lock().map_err(|_| {
75 std::io::Error::other("download checkpoint timer lock poisoned")
76 })?;
77 match guard.as_mut() {
78 Some(state) => state.begin_timer_checkpoint(),
79 None => Ok(false),
80 }
81 })
82 .await
83 .map_err(|error| std::io::Error::other(format!("timer worker failed: {error}")))??;
84 if !checkpoint_due {
85 return Ok::<(), std::io::Error>(());
86 }
87
88 {
89 let mut slot = file_slot.lock().await;
90 let file = slot.as_mut().ok_or_else(|| {
91 std::io::Error::other("locked download target missing during timed checkpoint")
92 })?;
93 file.sync_data().await?;
94 }
95 let commit_progress = Arc::clone(&progress);
96 tokio::task::spawn_blocking(move || {
97 let mut guard = commit_progress.lock().map_err(|_| {
98 std::io::Error::other("download checkpoint timer lock poisoned")
99 })?;
100 if let Some(state) = guard.as_mut() {
101 state.commit_checkpoint_after_data_sync()?;
102 }
103 Ok::<(), std::io::Error>(())
104 })
105 .await
106 .map_err(|error| std::io::Error::other(format!("timer worker failed: {error}")))??;
107 Ok(())
108 }
109 .await;
110 if let Err(error) = result {
111 crate::meow_warn_log!(
112 "download_checkpoint",
113 "timed .rcdl checkpoint failed: {}",
114 error
115 );
116 }
117 })
118}
119
120#[derive(Clone)]
125pub struct TransferTask {
126 file_sign: Arc<str>,
128 file_name: Arc<str>,
130 file_path: PathBuf,
132 upload_source: Option<UploadSource>,
134 upload_file_snapshot: Option<UploadFileSnapshot>,
135 direction: Direction,
137 total_size: u64,
139 chunk_size: u64,
141 url: String,
143 method: Method,
145 headers: HeaderMap,
147 breakpoint_download_http: BreakpointDownloadHttpConfig,
149 breakpoint_upload: Arc<dyn BreakpointUpload + Send + Sync>,
151 breakpoint_download: Arc<dyn BreakpointDownload + Send + Sync>,
153 http_client: Option<reqwest::Client>,
155 download_file_slot: Arc<Mutex<Option<tokio::fs::File>>>,
157 download_checkpoint_barrier: Arc<Mutex<()>>,
159 transfer_lifecycle: Arc<crate::inner::inner_task::TransferLifecycle>,
162 target_lease: Arc<StdMutex<Option<crate::target_lease::TargetLease>>>,
164 max_parts_in_flight: usize,
167 download_progress: Arc<StdMutex<Option<crate::dflt::download_progress::DownloadProgress>>>,
171 max_upload_prepare_retries: u32,
173}
174
175impl std::fmt::Debug for TransferTask {
176 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
177 f.debug_struct("TransferTask")
178 .field("file_sign", &self.file_sign)
179 .field("file_name", &self.file_name)
180 .field("file_path", &self.file_path)
181 .field("upload_source", &self.upload_source)
182 .field("direction", &self.direction)
183 .field("total_size", &self.total_size)
184 .field("chunk_size", &self.chunk_size)
185 .field("url", &self.url)
186 .field("method", &self.method)
187 .field("headers", &self.headers)
188 .field("breakpoint_upload", &"<dyn BreakpointUpload>")
189 .field("breakpoint_download", &"<dyn BreakpointDownload>")
190 .field("breakpoint_download_http", &self.breakpoint_download_http)
191 .field(
192 "max_upload_prepare_retries",
193 &self.max_upload_prepare_retries,
194 )
195 .finish()
196 }
197}
198
199impl TransferTask {
200 pub(crate) fn from_inner(inner: &InnerTask) -> Self {
202 Self {
203 file_sign: inner.file_sign_arc(),
204 file_name: inner.file_name_arc(),
205 file_path: inner.file_path().to_path_buf(),
206 upload_source: inner.upload_source().cloned(),
207 upload_file_snapshot: inner.upload_file_snapshot().cloned(),
208 direction: inner.direction(),
209 total_size: inner.total_size(),
210 chunk_size: inner.chunk_size(),
211 url: inner.url().to_string(),
212 method: inner.method(),
213 headers: inner.headers().clone(),
214 breakpoint_download_http: inner.breakpoint_download_http().clone(),
215 breakpoint_upload: inner.breakpoint_upload().clone(),
216 breakpoint_download: inner.breakpoint_download().clone(),
217 http_client: inner.http_client_ref().cloned(),
218 download_file_slot: Arc::new(Mutex::new(None)),
219 download_checkpoint_barrier: Arc::new(Mutex::new(())),
220 transfer_lifecycle: inner.transfer_lifecycle(),
221 target_lease: Arc::new(StdMutex::new(None)),
222 max_parts_in_flight: inner.max_parts_in_flight(),
223 download_progress: Arc::new(StdMutex::new(None)),
224 max_upload_prepare_retries: inner.max_upload_prepare_retries(),
225 }
226 }
227
228 pub fn direction(&self) -> Direction {
240 self.direction
241 }
242
243 pub fn total_size(&self) -> u64 {
255 self.total_size
256 }
257
258 pub fn chunk_size(&self) -> u64 {
270 self.chunk_size
271 }
272
273 pub fn file_sign(&self) -> &str {
285 &self.file_sign
286 }
287
288 pub fn file_name(&self) -> &str {
300 &self.file_name
301 }
302
303 pub fn file_path(&self) -> &Path {
315 &self.file_path
316 }
317
318 pub(crate) fn upload_source(&self) -> Option<&UploadSource> {
320 self.upload_source.as_ref()
321 }
322
323 pub fn url(&self) -> &str {
335 &self.url
336 }
337
338 pub fn method(&self) -> Method {
350 self.method.clone()
351 }
352
353 pub fn headers(&self) -> &HeaderMap {
365 &self.headers
366 }
367
368 pub fn breakpoint_download_http(&self) -> Option<&BreakpointDownloadHttpConfig> {
383 Some(&self.breakpoint_download_http)
384 }
385
386 pub(crate) fn breakpoint_upload(&self) -> Option<&Arc<dyn BreakpointUpload + Send + Sync>> {
388 Some(&self.breakpoint_upload)
389 }
390
391 pub(crate) fn breakpoint_download(&self) -> Option<&Arc<dyn BreakpointDownload + Send + Sync>> {
393 Some(&self.breakpoint_download)
394 }
395
396 pub(crate) fn max_upload_prepare_retries(&self) -> u32 {
398 self.max_upload_prepare_retries
399 }
400
401 pub(crate) fn http_client_ref(&self) -> Option<&reqwest::Client> {
403 self.http_client.as_ref()
404 }
405
406 pub(crate) fn upload_file_snapshot(&self) -> Option<&UploadFileSnapshot> {
407 self.upload_file_snapshot.as_ref()
408 }
409
410 pub(crate) fn begin_upload_abort(&self) -> bool {
411 self.transfer_lifecycle.begin_abort()
412 }
413
414 pub(crate) fn require_upload_abort(&self) {
415 if self.direction == Direction::Upload {
416 self.transfer_lifecycle.require_abort();
417 }
418 }
419
420 pub(crate) fn upload_abort_required(&self) -> bool {
421 self.direction == Direction::Upload && self.transfer_lifecycle.abort_required()
422 }
423
424 pub(crate) fn begin_terminal_completion(&self) -> bool {
425 self.transfer_lifecycle.begin_completion()
426 }
427
428 pub(crate) fn acknowledge_terminal_completion(&self) {
429 self.transfer_lifecycle.acknowledge_completion();
430 }
431
432 pub(crate) fn finish_incomplete_completion(
433 &self,
434 ) -> Option<crate::inner::inner_task::DeferredStop> {
435 self.transfer_lifecycle.finish_incomplete_completion()
436 }
437
438 pub(crate) async fn ensure_download_target_lease(&self) -> Result<(), crate::MeowError> {
439 let slot = Arc::clone(&self.target_lease);
440 let path = self.file_path.clone();
441 tokio::task::spawn_blocking(move || {
442 let mut guard = slot.lock().map_err(|_| {
443 crate::MeowError::from_code_str(
444 crate::InnerErrorCode::LockPoisoned,
445 "download target lease lock poisoned",
446 )
447 })?;
448 if guard.is_none() {
449 *guard = Some(
450 crate::target_lease::TargetLease::acquire(&path).map_err(|e| {
451 map_download_target_lock_error(
452 format!("acquire download target lease failed: {}", path.display()),
453 e,
454 )
455 })?,
456 );
457 }
458 Ok(())
459 })
460 .await
461 .map_err(|e| {
462 crate::MeowError::from_code(
463 crate::InnerErrorCode::IoError,
464 format!("download target lease worker failed: {e}"),
465 )
466 })?
467 }
468
469 pub(crate) async fn ensure_download_target_file_locked(&self) -> Result<u64, crate::MeowError> {
473 let mut slot = self.download_file_slot.lock().await;
474 if let Some(file) = slot.as_ref() {
475 return file
476 .metadata()
477 .await
478 .map(|metadata| metadata.len())
479 .map_err(|e| {
480 crate::MeowError::from_io(
481 format!(
482 "stat locked download target failed: {}",
483 self.file_path.display()
484 ),
485 e,
486 )
487 });
488 }
489
490 let path = self.file_path.clone();
491 let display = path.display().to_string();
492 let file = tokio::task::spawn_blocking(move || {
493 crate::target_lease::open_locked_target(&path, true)
494 })
495 .await
496 .map_err(|e| {
497 crate::MeowError::from_code(
498 crate::InnerErrorCode::IoError,
499 format!("download target lock worker failed: {e}"),
500 )
501 })?
502 .map_err(|e| {
503 map_download_target_lock_error(
504 format!("lock actual download target failed: {display}"),
505 e,
506 )
507 })?;
508 let len = file
509 .metadata()
510 .map(|metadata| metadata.len())
511 .map_err(|e| {
512 crate::MeowError::from_io(
513 format!("stat newly locked download target failed: {display}"),
514 e,
515 )
516 })?;
517 *slot = Some(tokio::fs::File::from_std(file));
518 Ok(len)
519 }
520
521 pub(crate) async fn release_download_target_file_lock(&self) -> Result<(), crate::MeowError> {
525 let file = {
526 let mut slot = self.download_file_slot.lock().await;
527 slot.take()
528 };
529 let Some(file) = file else {
530 return Ok(());
531 };
532 file.sync_all().await.map_err(|e| {
533 crate::MeowError::from_io(
534 format!(
535 "sync locked download target failed: {}",
536 self.file_path.display()
537 ),
538 e,
539 )
540 })?;
541 let file = file.into_std().await;
542 tokio::task::spawn_blocking(move || FileExt::unlock(&file))
543 .await
544 .map_err(|e| {
545 crate::MeowError::from_code(
546 crate::InnerErrorCode::IoError,
547 format!("download target unlock worker failed: {e}"),
548 )
549 })?
550 .map_err(|e| crate::MeowError::from_io("unlock download target failed".to_owned(), e))
551 }
552
553 pub(crate) fn release_download_target_lease(&self) -> Result<(), crate::MeowError> {
558 let lease = self
559 .target_lease
560 .lock()
561 .map_err(|_| {
562 crate::MeowError::from_code_str(
563 crate::InnerErrorCode::LockPoisoned,
564 "download target lease lock poisoned",
565 )
566 })?
567 .take();
568 drop(lease);
569 Ok(())
570 }
571
572 pub(crate) fn download_file_slot(&self) -> &Arc<Mutex<Option<tokio::fs::File>>> {
574 &self.download_file_slot
575 }
576
577 pub(crate) fn max_parts_in_flight(&self) -> usize {
579 self.max_parts_in_flight
580 }
581
582 pub(crate) fn download_progress(
584 &self,
585 ) -> &Arc<StdMutex<Option<crate::dflt::download_progress::DownloadProgress>>> {
586 &self.download_progress
587 }
588
589 pub(crate) async fn write_and_stage_download_part(
593 &self,
594 offset: u64,
595 body: &[u8],
596 digest: [u8; 32],
597 ) -> Result<(), crate::MeowError> {
598 let _barrier = self.download_checkpoint_barrier.lock().await;
599 let body_len = u64::try_from(body.len()).map_err(|_| {
600 crate::MeowError::from_code_str(
601 crate::InnerErrorCode::InvalidRange,
602 "parallel download part length does not fit u64",
603 )
604 })?;
605 let part_end = offset.checked_add(body_len).ok_or_else(|| {
606 crate::MeowError::from_code_str(
607 crate::InnerErrorCode::InvalidRange,
608 "parallel download part end overflow",
609 )
610 })?;
611 {
612 let mut slot = self.download_file_slot.lock().await;
613 let file = slot.as_mut().ok_or_else(|| {
614 crate::MeowError::from_code_str(
615 crate::InnerErrorCode::InvalidTaskState,
616 "locked download target missing during positioned write",
617 )
618 })?;
619 let file_len = file
620 .metadata()
621 .await
622 .map(|metadata| metadata.len())
623 .map_err(|e| {
624 crate::MeowError::from_io(
625 format!(
626 "stat locked download target failed: {}",
627 self.file_path.display()
628 ),
629 e,
630 )
631 })?;
632 if file_len < part_end {
633 return Err(crate::MeowError::from_code(
634 crate::InnerErrorCode::LocalFileRemoved,
635 format!(
636 "download target was truncated before positioned write: len={file_len} need>={part_end}"
637 ),
638 ));
639 }
640 file.seek(std::io::SeekFrom::Start(offset))
641 .await
642 .map_err(|e| {
643 crate::MeowError::from_io(
644 format!("seek locked download target failed: offset={offset}"),
645 e,
646 )
647 })?;
648 file.write_all(body).await.map_err(|e| {
649 crate::MeowError::from_io(
650 format!("write locked download target failed: offset={offset}"),
651 e,
652 )
653 })?;
654 file.flush().await.map_err(|e| {
655 crate::MeowError::from_io(
656 format!("flush locked download target failed: offset={offset}"),
657 e,
658 )
659 })?;
660 }
661
662 self.stage_download_part_digest_locked(offset, digest, true)
663 .await
664 }
665
666 pub(crate) async fn stage_serial_download_part(
670 &self,
671 offset: u64,
672 digest: [u8; 32],
673 ) -> Result<(), crate::MeowError> {
674 let _barrier = self.download_checkpoint_barrier.lock().await;
675 self.stage_download_part_digest_locked(offset, digest, true)
676 .await
677 }
678
679 pub(crate) async fn retain_serial_download_contiguous_progress(
684 &self,
685 ) -> Result<u64, crate::MeowError> {
686 let _barrier = self.download_checkpoint_barrier.lock().await;
687 let watermark = self
688 .download_progress
689 .lock()
690 .map_err(|_| {
691 crate::MeowError::from_code_str(
692 crate::InnerErrorCode::LockPoisoned,
693 "download checkpoint lock poisoned",
694 )
695 })?
696 .as_ref()
697 .ok_or_else(|| {
698 crate::MeowError::from_code_str(
699 crate::InnerErrorCode::InvalidTaskState,
700 "download checkpoint state missing during serial compaction",
701 )
702 })?
703 .contiguous_watermark();
704
705 {
706 let mut slot = self.download_file_slot.lock().await;
707 let file = slot.as_mut().ok_or_else(|| {
708 crate::MeowError::from_code_str(
709 crate::InnerErrorCode::InvalidTaskState,
710 "locked download target missing during serial compaction",
711 )
712 })?;
713 file.set_len(watermark).await.map_err(|e| {
714 crate::MeowError::from_io(
715 format!(
716 "truncate serial download to verified prefix failed: {}",
717 self.file_path.display()
718 ),
719 e,
720 )
721 })?;
722 file.sync_all().await.map_err(|e| {
723 crate::MeowError::from_io("sync serial verified prefix failed".to_owned(), e)
724 })?;
725 }
726
727 let progress = Arc::clone(&self.download_progress);
728 let persisted = tokio::task::spawn_blocking(move || {
729 let mut guard = progress.lock().map_err(|_| {
730 crate::MeowError::from_code_str(
731 crate::InnerErrorCode::LockPoisoned,
732 "download checkpoint lock poisoned",
733 )
734 })?;
735 let state = guard.as_mut().ok_or_else(|| {
736 crate::MeowError::from_code_str(
737 crate::InnerErrorCode::InvalidTaskState,
738 "download checkpoint state missing during serial compaction",
739 )
740 })?;
741 state
742 .retain_contiguous_prefix_after_data_sync()
743 .map_err(|e| {
744 crate::MeowError::from_io(
745 "persist compacted serial checkpoint failed".to_owned(),
746 e,
747 )
748 })
749 })
750 .await
751 .map_err(|e| {
752 crate::MeowError::from_code(
753 crate::InnerErrorCode::IoError,
754 format!("serial compaction checkpoint worker failed: {e}"),
755 )
756 })??;
757 if persisted != watermark {
758 return Err(crate::MeowError::from_code(
759 crate::InnerErrorCode::InvalidTaskState,
760 format!(
761 "serial checkpoint watermark changed during compaction: expected={watermark} actual={persisted}"
762 ),
763 ));
764 }
765 Ok(watermark)
766 }
767
768 async fn stage_download_part_digest_locked(
769 &self,
770 offset: u64,
771 digest: [u8; 32],
772 arm_timer: bool,
773 ) -> Result<(), crate::MeowError> {
774 let progress = Arc::clone(&self.download_progress);
775 let timer_progress = Arc::clone(&progress);
776 let outcome = tokio::task::spawn_blocking(move || {
777 let mut guard = progress.lock().map_err(|_| {
778 crate::MeowError::from_code_str(
779 crate::InnerErrorCode::LockPoisoned,
780 "download checkpoint lock poisoned",
781 )
782 })?;
783 let state = guard.as_mut().ok_or_else(|| {
784 crate::MeowError::from_code_str(
785 crate::InnerErrorCode::InvalidTaskState,
786 "download checkpoint state missing",
787 )
788 })?;
789 state
790 .stage_done_with_digest_deferred(offset, digest)
791 .map_err(|e| {
792 crate::MeowError::from_io("stage .rcdl checkpoint failed".to_owned(), e)
793 })
794 })
795 .await
796 .map_err(|e| {
797 crate::MeowError::from_code(
798 crate::InnerErrorCode::IoError,
799 format!("download checkpoint worker failed: {e}"),
800 )
801 })??;
802 if outcome.checkpoint_due {
803 self.checkpoint_locked_download_target().await?;
804 } else if outcome.arm_timer && arm_timer {
805 drop(arm_download_checkpoint_timer(
809 Arc::downgrade(&timer_progress),
810 Arc::downgrade(&self.download_file_slot),
811 Arc::downgrade(&self.download_checkpoint_barrier),
812 ));
813 }
814 Ok(())
815 }
816
817 async fn checkpoint_locked_download_target(&self) -> Result<(), crate::MeowError> {
820 let progress = Arc::clone(&self.download_progress);
821 let checkpoint_due = tokio::task::spawn_blocking(move || {
822 let mut guard = progress.lock().map_err(|_| {
823 crate::MeowError::from_code_str(
824 crate::InnerErrorCode::LockPoisoned,
825 "download checkpoint lock poisoned",
826 )
827 })?;
828 let Some(state) = guard.as_mut() else {
829 return Ok(false);
830 };
831 let due = state.begin_external_checkpoint().map_err(|e| {
832 crate::MeowError::from_io("begin .rcdl checkpoint failed".to_owned(), e)
833 })?;
834 Ok(due)
835 })
836 .await
837 .map_err(|e| {
838 crate::MeowError::from_code(
839 crate::InnerErrorCode::IoError,
840 format!("download checkpoint worker failed: {e}"),
841 )
842 })??;
843 if !checkpoint_due {
844 return Ok(());
845 }
846
847 {
848 let mut slot = self.download_file_slot.lock().await;
849 let file = slot.as_mut().ok_or_else(|| {
850 crate::MeowError::from_code_str(
851 crate::InnerErrorCode::InvalidTaskState,
852 "locked download target missing during checkpoint",
853 )
854 })?;
855 file.sync_data().await.map_err(|e| {
856 crate::MeowError::from_io("sync download checkpoint data failed".to_owned(), e)
857 })?;
858 }
859
860 let progress = Arc::clone(&self.download_progress);
861 tokio::task::spawn_blocking(move || {
862 let mut guard = progress.lock().map_err(|_| {
863 crate::MeowError::from_code_str(
864 crate::InnerErrorCode::LockPoisoned,
865 "download checkpoint lock poisoned",
866 )
867 })?;
868 if let Some(state) = guard.as_mut() {
869 state.commit_checkpoint_after_data_sync().map_err(|e| {
870 crate::MeowError::from_io("commit .rcdl checkpoint failed".to_owned(), e)
871 })?;
872 }
873 Ok(())
874 })
875 .await
876 .map_err(|e| {
877 crate::MeowError::from_code(
878 crate::InnerErrorCode::IoError,
879 format!("download checkpoint worker failed: {e}"),
880 )
881 })?
882 }
883
884 pub(crate) async fn force_download_checkpoint(&self) -> Result<(), crate::MeowError> {
885 let _barrier = self.download_checkpoint_barrier.lock().await;
886 self.checkpoint_locked_download_target().await
887 }
888
889 pub(crate) async fn take_download_progress_after_checkpoint(
890 &self,
891 ) -> Result<Option<crate::dflt::download_progress::DownloadProgress>, crate::MeowError> {
892 #[cfg(windows)]
893 {
894 let file = {
900 let mut slot = self.download_file_slot.lock().await;
901 slot.take().ok_or_else(|| {
902 crate::MeowError::from_code_str(
903 crate::InnerErrorCode::InvalidTaskState,
904 "locked download target missing during final validation",
905 )
906 })?
907 };
908 let file = file.into_std().await;
909 let progress = Arc::clone(&self.download_progress);
910 let worker = tokio::task::spawn_blocking(move || {
911 let mut file = file;
912 let result = (|| {
913 let mut guard = progress.lock().map_err(|_| {
914 crate::MeowError::from_code_str(
915 crate::InnerErrorCode::LockPoisoned,
916 "download checkpoint lock poisoned",
917 )
918 })?;
919 if let Some(state) = guard.as_ref() {
920 state
921 .validate_committed_content_on_locked_file(&mut file)
922 .map_err(map_download_validation_error)?;
923 }
924 Ok(guard.take())
925 })();
926 (result, file)
927 })
928 .await;
929 return match worker {
930 Ok((result, file)) => {
931 let mut slot = self.download_file_slot.lock().await;
932 if slot.is_some() {
933 return Err(crate::MeowError::from_code_str(
934 crate::InnerErrorCode::InvalidTaskState,
935 "download target slot changed during final validation",
936 ));
937 }
938 *slot = Some(tokio::fs::File::from_std(file));
939 result
940 }
941 Err(error) => Err(crate::MeowError::from_code(
942 crate::InnerErrorCode::IoError,
943 format!("download validation worker failed: {error}"),
944 )),
945 };
946 }
947
948 #[cfg(not(windows))]
949 {
950 let progress = Arc::clone(&self.download_progress);
951 tokio::task::spawn_blocking(move || {
952 let mut guard = progress.lock().map_err(|_| {
953 crate::MeowError::from_code_str(
954 crate::InnerErrorCode::LockPoisoned,
955 "download checkpoint lock poisoned",
956 )
957 })?;
958 if let Some(state) = guard.as_ref() {
959 state
960 .validate_committed_content()
961 .map_err(map_download_validation_error)?;
962 }
963 Ok(guard.take())
964 })
965 .await
966 .map_err(|e| {
967 crate::MeowError::from_code(
968 crate::InnerErrorCode::IoError,
969 format!("download checkpoint worker failed: {e}"),
970 )
971 })?
972 }
973 }
974
975 pub(crate) async fn finalize_download_content(
979 &self,
980 expected_total: u64,
981 ) -> Result<(), crate::MeowError> {
982 let _barrier = self.download_checkpoint_barrier.lock().await;
983 let finalization = async {
984 self.checkpoint_locked_download_target().await?;
985 let visible_len = tokio::fs::metadata(&self.file_path)
986 .await
987 .map(|metadata| metadata.len())
988 .map_err(|e| {
989 crate::MeowError::from_io(
990 format!("stat completed download failed: {}", self.file_path.display()),
991 e,
992 )
993 })?;
994 if visible_len != expected_total {
995 return Err(crate::MeowError::from_code(
996 crate::InnerErrorCode::LocalFileRemoved,
997 format!(
998 "download length changed before complete: expected={expected_total} actual={visible_len}"
999 ),
1000 ));
1001 }
1002 let progress = self
1003 .take_download_progress_after_checkpoint()
1004 .await?
1005 .ok_or_else(|| {
1006 crate::MeowError::from_code_str(
1007 crate::InnerErrorCode::InvalidTaskState,
1008 "download progress missing during final validation",
1009 )
1010 })?;
1011 if progress.total() != expected_total {
1012 return Err(crate::MeowError::from_code(
1013 crate::InnerErrorCode::InvalidRange,
1014 format!(
1015 "download progress total mismatch: expected={expected_total} progress={}",
1016 progress.total()
1017 ),
1018 ));
1019 }
1020 if !progress.all_done() {
1021 return Err(crate::MeowError::from_code_str(
1022 crate::InnerErrorCode::InvalidRange,
1023 "download complete called before all parts recorded done",
1024 ));
1025 }
1026 if let Err(error) = progress.delete() {
1027 crate::meow_warn_log!("download_complete", "sidecar delete failed: {}", error);
1028 }
1029 Ok(())
1030 }
1031 .await;
1032
1033 let file_release = self.release_download_target_file_lock().await;
1036 let lease_release = self.release_download_target_lease();
1037 match (finalization, file_release, lease_release) {
1038 (Err(error), _, _) => Err(error),
1039 (Ok(()), Err(error), _) => Err(error),
1040 (Ok(()), Ok(()), Err(error)) => Err(error),
1041 (Ok(()), Ok(()), Ok(())) => Ok(()),
1042 }
1043 }
1044}
1045
1046#[cfg(test)]
1047mod target_lock_error_tests {
1048 use super::map_download_target_lock_error;
1049
1050 #[test]
1051 fn contention_is_an_invalid_task_state() {
1052 let error = map_download_target_lock_error(
1053 "lock target",
1054 std::io::Error::new(std::io::ErrorKind::WouldBlock, "owned elsewhere"),
1055 );
1056
1057 assert_eq!(error.code(), crate::InnerErrorCode::InvalidTaskState as i32);
1058 }
1059
1060 #[test]
1061 fn non_contention_uses_normal_io_classification() {
1062 let permission = map_download_target_lock_error(
1063 "open target",
1064 std::io::Error::new(std::io::ErrorKind::PermissionDenied, "denied"),
1065 );
1066 let missing_parent = map_download_target_lock_error(
1067 "open target",
1068 std::io::Error::new(std::io::ErrorKind::NotFound, "gone"),
1069 );
1070
1071 assert_eq!(permission.code(), crate::InnerErrorCode::IoError as i32);
1072 assert_eq!(
1073 missing_parent.code(),
1074 crate::InnerErrorCode::LocalFileRemoved as i32
1075 );
1076 }
1077
1078 #[cfg(any(unix, windows))]
1079 #[test]
1080 fn out_of_space_preserves_disk_full() {
1081 #[cfg(unix)]
1082 let raw_error = 28;
1083 #[cfg(windows)]
1084 let raw_error = 112;
1085
1086 let error = map_download_target_lock_error(
1087 "create target ownership file",
1088 std::io::Error::from_raw_os_error(raw_error),
1089 );
1090
1091 assert_eq!(error.code(), crate::InnerErrorCode::DiskFull as i32);
1092 }
1093}
1094
1095#[cfg(test)]
1096mod checkpoint_timer_tests {
1097 use super::{arm_download_checkpoint_timer, TransferTask};
1098 use crate::dflt::download_progress::{sidecar_path, DownloadProgress};
1099 use crate::direction::Direction;
1100 use crate::http_breakpoint::{
1101 BreakpointDownloadHttpConfig, DefaultStyleUpload, StandardRangeDownload,
1102 };
1103 use reqwest::{header::HeaderMap, Method};
1104 use std::sync::{Arc, Mutex};
1105
1106 fn download_task_for(target: &std::path::Path) -> TransferTask {
1107 TransferTask {
1108 file_sign: Arc::<str>::from("positioned-write-test"),
1109 file_name: Arc::<str>::from("positioned-write-test.bin"),
1110 file_path: target.to_path_buf(),
1111 upload_source: None,
1112 upload_file_snapshot: None,
1113 direction: Direction::Download,
1114 total_size: 16,
1115 chunk_size: 8,
1116 url: "http://127.0.0.1/positioned-write-test".to_owned(),
1117 method: Method::GET,
1118 headers: HeaderMap::new(),
1119 breakpoint_download_http: BreakpointDownloadHttpConfig::default(),
1120 breakpoint_upload: Arc::new(DefaultStyleUpload::default()),
1121 breakpoint_download: Arc::new(StandardRangeDownload),
1122 http_client: None,
1123 download_file_slot: Arc::new(tokio::sync::Mutex::new(None)),
1124 download_checkpoint_barrier: Arc::new(tokio::sync::Mutex::new(())),
1125 transfer_lifecycle: Arc::new(crate::inner::inner_task::TransferLifecycle::new()),
1126 target_lease: Arc::new(Mutex::new(None)),
1127 max_parts_in_flight: 2,
1128 download_progress: Arc::new(Mutex::new(None)),
1129 max_upload_prepare_retries: 0,
1130 }
1131 }
1132
1133 #[tokio::test]
1134 async fn lone_staged_part_is_checkpointed_by_wall_clock_timer() {
1135 let target = std::env::temp_dir().join(format!(
1136 "rusty_cat_checkpoint_timer_{}_{}",
1137 std::process::id(),
1138 std::time::SystemTime::now()
1139 .duration_since(std::time::UNIX_EPOCH)
1140 .expect("clock")
1141 .as_nanos()
1142 ));
1143 std::fs::write(&target, vec![7_u8; 20]).expect("target");
1144 let mut progress =
1145 DownloadProgress::load_or_create(&target, 20, 10, 4, "timer-test").expect("progress");
1146 assert!(!progress.stage_done(0).expect("stage"));
1147 assert!(!progress.is_done(0), "the batch threshold was not reached");
1148
1149 let slot = Arc::new(Mutex::new(Some(progress)));
1150 let target_file = std::fs::OpenOptions::new()
1151 .read(true)
1152 .write(true)
1153 .open(&target)
1154 .expect("target handle");
1155 let file_slot = Arc::new(tokio::sync::Mutex::new(Some(tokio::fs::File::from_std(
1156 target_file,
1157 ))));
1158 let barrier = Arc::new(tokio::sync::Mutex::new(()));
1159 arm_download_checkpoint_timer(
1160 Arc::downgrade(&slot),
1161 Arc::downgrade(&file_slot),
1162 Arc::downgrade(&barrier),
1163 )
1164 .await
1165 .expect("timer task");
1166 assert!(
1167 slot.lock()
1168 .expect("progress lock")
1169 .as_ref()
1170 .expect("progress")
1171 .is_done(0),
1172 "the 250 ms wall-clock wake-up must publish the staged part"
1173 );
1174
1175 let _ = std::fs::remove_file(sidecar_path(&target));
1176 let _ = std::fs::remove_file(target);
1177 }
1178
1179 #[tokio::test]
1180 async fn positioned_write_reports_local_file_removed_when_locked_target_was_truncated() {
1181 let target = std::env::temp_dir().join(format!(
1182 "rusty_cat_positioned_write_truncated_{}_{}",
1183 std::process::id(),
1184 std::time::SystemTime::now()
1185 .duration_since(std::time::UNIX_EPOCH)
1186 .expect("clock")
1187 .as_nanos()
1188 ));
1189 std::fs::write(&target, vec![0_u8; 16]).expect("preallocate target");
1190 let task = download_task_for(&target);
1191 assert_eq!(
1192 task.ensure_download_target_file_locked()
1193 .await
1194 .expect("lock actual target"),
1195 16
1196 );
1197 {
1198 let mut slot = task.download_file_slot.lock().await;
1199 slot.as_mut()
1200 .expect("locked target handle")
1201 .set_len(4)
1202 .await
1203 .expect("truncate actual target before positioned write");
1204 }
1205
1206 let error = task
1207 .write_and_stage_download_part(8, b"abcdefgh", [0_u8; 32])
1208 .await
1209 .expect_err("a part extending beyond the truncated target must fail");
1210 assert_eq!(error.code(), crate::InnerErrorCode::LocalFileRemoved as i32);
1211
1212 task.release_download_target_file_lock()
1213 .await
1214 .expect("unlock actual target");
1215 let _ = std::fs::remove_file(target);
1216 }
1217
1218 #[test]
1219 fn final_content_validation_uses_bytes_not_file_identity() {
1220 let target = std::env::temp_dir().join(format!(
1221 "rusty_cat_download_content_validation_{}_{}",
1222 std::process::id(),
1223 std::time::SystemTime::now()
1224 .duration_since(std::time::UNIX_EPOCH)
1225 .expect("clock")
1226 .as_nanos()
1227 ));
1228 let original = b"same bytes across replacement";
1229 std::fs::write(&target, original).expect("target");
1230 let mut progress = DownloadProgress::load_or_create(
1231 &target,
1232 original.len() as u64,
1233 original.len() as u64,
1234 1,
1235 "content-identity-test",
1236 )
1237 .expect("progress");
1238 progress
1239 .mark_done_and_persist(0)
1240 .expect("persist content digest");
1241
1242 let replacement = target.with_extension("replacement");
1243 std::fs::write(&replacement, original).expect("same-content replacement");
1244 std::fs::rename(&replacement, &target).expect("replace target");
1245 progress
1246 .validate_committed_content()
1247 .expect("identical bytes are the same content generation");
1248
1249 std::fs::write(&target, vec![b'X'; original.len()]).expect("different bytes");
1250 let error = progress
1251 .validate_committed_content()
1252 .expect_err("same-length different content must fail");
1253 assert_eq!(error.kind(), std::io::ErrorKind::InvalidData);
1254 let _ = std::fs::remove_file(sidecar_path(&target));
1255 let _ = std::fs::remove_file(target);
1256 }
1257}