Skip to main content

rusty_cat/
transfer_task.rs

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
37/// Maps failures while acquiring either layer of download-target ownership.
38///
39/// `target_lease` normalizes every genuine advisory-lock conflict to
40/// `WouldBlock`. Only that condition is a task-state conflict; failures while
41/// creating/opening the lease or target are ordinary local I/O failures and
42/// must retain `MeowError::from_io` classifications such as `DiskFull` and
43/// `LocalFileRemoved`.
44fn 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            // The transfer ended before the timer fired. Weak references ensure
67            // a stale timer never extends the actual target lock lifetime.
68            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/// Immutable task snapshot exposed to transfer executor implementations.
121///
122/// This type is constructed from the crate-internal scheduler task state and
123/// intentionally exposes read-only accessors.
124#[derive(Clone)]
125pub struct TransferTask {
126    /// Stable file signature.
127    file_sign: Arc<str>,
128    /// Display file name.
129    file_name: Arc<str>,
130    /// Local file path.
131    file_path: PathBuf,
132    /// Upload-only source descriptor.
133    upload_source: Option<UploadSource>,
134    upload_file_snapshot: Option<UploadFileSnapshot>,
135    /// Transfer direction.
136    direction: Direction,
137    /// Total file size in bytes.
138    total_size: u64,
139    /// Chunk size in bytes.
140    chunk_size: u64,
141    /// Request URL.
142    url: String,
143    /// Request HTTP method.
144    method: Method,
145    /// Base request headers.
146    headers: HeaderMap,
147    /// HTTP config for breakpoint download behavior.
148    breakpoint_download_http: BreakpointDownloadHttpConfig,
149    /// Upload breakpoint protocol implementation.
150    breakpoint_upload: Arc<dyn BreakpointUpload + Send + Sync>,
151    /// Download breakpoint protocol implementation.
152    breakpoint_download: Arc<dyn BreakpointDownload + Send + Sync>,
153    /// Optional per-task custom HTTP client.
154    http_client: Option<reqwest::Client>,
155    /// Task-level download file handle slot to avoid reopening per chunk.
156    download_file_slot: Arc<Mutex<Option<tokio::fs::File>>>,
157    /// Serializes writes and checkpoint data barriers for the locked target.
158    download_checkpoint_barrier: Arc<Mutex<()>>,
159    /// Shared terminal arbitration and abort idempotence for every task view
160    /// derived from the same scheduler entry.
161    transfer_lifecycle: Arc<crate::inner::inner_task::TransferLifecycle>,
162    /// Cross-client/process ownership of the visible download target.
163    target_lease: Arc<StdMutex<Option<crate::target_lease::TargetLease>>>,
164    /// Max parts of this file transferred concurrently (intra-file parallel).
165    /// `1` means the strict-serial legacy path.
166    max_parts_in_flight: usize,
167    /// Shared progress bitmap for the concurrent download path (None until the
168    /// parallel `download_prepare` initializes it). Guarded so concurrent parts
169    /// can flip their bit without racing.
170    download_progress: Arc<StdMutex<Option<crate::dflt::download_progress::DownloadProgress>>>,
171    /// Max retries after first failed upload prepare (`BreakpointUpload::prepare`).
172    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    /// Creates a transfer snapshot from an internal runtime task.
201    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    /// Returns transfer direction.
229    ///
230    /// # Examples
231    ///
232    /// ```no_run
233    /// use rusty_cat::api::TransferTask;
234    ///
235    /// fn inspect(task: &TransferTask) {
236    ///     let _ = task.direction();
237    /// }
238    /// ```
239    pub fn direction(&self) -> Direction {
240        self.direction
241    }
242
243    /// Returns total file size in bytes.
244    ///
245    /// # Examples
246    ///
247    /// ```no_run
248    /// use rusty_cat::api::TransferTask;
249    ///
250    /// fn inspect(task: &TransferTask) {
251    ///     let _ = task.total_size();
252    /// }
253    /// ```
254    pub fn total_size(&self) -> u64 {
255        self.total_size
256    }
257
258    /// Returns chunk size in bytes.
259    ///
260    /// # Examples
261    ///
262    /// ```no_run
263    /// use rusty_cat::api::TransferTask;
264    ///
265    /// fn inspect(task: &TransferTask) {
266    ///     let _ = task.chunk_size();
267    /// }
268    /// ```
269    pub fn chunk_size(&self) -> u64 {
270        self.chunk_size
271    }
272
273    /// Returns file signature.
274    ///
275    /// # Examples
276    ///
277    /// ```no_run
278    /// use rusty_cat::api::TransferTask;
279    ///
280    /// fn inspect(task: &TransferTask) {
281    ///     let _ = task.file_sign();
282    /// }
283    /// ```
284    pub fn file_sign(&self) -> &str {
285        &self.file_sign
286    }
287
288    /// Returns display file name.
289    ///
290    /// # Examples
291    ///
292    /// ```no_run
293    /// use rusty_cat::api::TransferTask;
294    ///
295    /// fn inspect(task: &TransferTask) {
296    ///     let _ = task.file_name();
297    /// }
298    /// ```
299    pub fn file_name(&self) -> &str {
300        &self.file_name
301    }
302
303    /// Returns local file path.
304    ///
305    /// # Examples
306    ///
307    /// ```no_run
308    /// use rusty_cat::api::TransferTask;
309    ///
310    /// fn inspect(task: &TransferTask) {
311    ///     let _ = task.file_path();
312    /// }
313    /// ```
314    pub fn file_path(&self) -> &Path {
315        &self.file_path
316    }
317
318    /// Returns upload source for upload tasks.
319    pub(crate) fn upload_source(&self) -> Option<&UploadSource> {
320        self.upload_source.as_ref()
321    }
322
323    /// Returns request URL.
324    ///
325    /// # Examples
326    ///
327    /// ```no_run
328    /// use rusty_cat::api::TransferTask;
329    ///
330    /// fn inspect(task: &TransferTask) {
331    ///     let _ = task.url();
332    /// }
333    /// ```
334    pub fn url(&self) -> &str {
335        &self.url
336    }
337
338    /// Returns request HTTP method.
339    ///
340    /// # Examples
341    ///
342    /// ```no_run
343    /// use rusty_cat::api::TransferTask;
344    ///
345    /// fn inspect(task: &TransferTask) {
346    ///     let _ = task.method();
347    /// }
348    /// ```
349    pub fn method(&self) -> Method {
350        self.method.clone()
351    }
352
353    /// Returns base request headers.
354    ///
355    /// # Examples
356    ///
357    /// ```no_run
358    /// use rusty_cat::api::TransferTask;
359    ///
360    /// fn inspect(task: &TransferTask) {
361    ///     let _ = task.headers();
362    /// }
363    /// ```
364    pub fn headers(&self) -> &HeaderMap {
365        &self.headers
366    }
367
368    /// Returns task-level breakpoint download HTTP configuration.
369    ///
370    /// Custom [`crate::download_trait::BreakpointDownload`] implementations can
371    /// read values such as `range_accept`.
372    ///
373    /// # Examples
374    ///
375    /// ```no_run
376    /// use rusty_cat::api::TransferTask;
377    ///
378    /// fn inspect(task: &TransferTask) {
379    ///     let _ = task.breakpoint_download_http();
380    /// }
381    /// ```
382    pub fn breakpoint_download_http(&self) -> Option<&BreakpointDownloadHttpConfig> {
383        Some(&self.breakpoint_download_http)
384    }
385
386    /// Returns task-level upload protocol implementation.
387    pub(crate) fn breakpoint_upload(&self) -> Option<&Arc<dyn BreakpointUpload + Send + Sync>> {
388        Some(&self.breakpoint_upload)
389    }
390
391    /// Returns task-level download protocol implementation.
392    pub(crate) fn breakpoint_download(&self) -> Option<&Arc<dyn BreakpointDownload + Send + Sync>> {
393        Some(&self.breakpoint_download)
394    }
395
396    /// Returns max retries after the first failed upload prepare.
397    pub(crate) fn max_upload_prepare_retries(&self) -> u32 {
398        self.max_upload_prepare_retries
399    }
400
401    /// Returns task-level custom HTTP client, if configured.
402    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    /// Opens the actual target and acquires its cross-platform file lock.
470    /// Every transfer write must reuse the returned task slot: on Windows a
471    /// second handle cannot access a range locked through the first handle.
472    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    /// Flushes, explicitly unlocks and closes the actual target handle. The
522    /// caller releases the path lease next, after completion validation or
523    /// failure checkpointing and before publishing a terminal event.
524    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    /// Releases the normalized path/inode lease after the actual target handle
554    /// has been closed. Terminal events must not become observable before this
555    /// succeeds, otherwise an immediate same-target enqueue can spuriously see
556    /// the completed task as still owning the path.
557    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    /// Returns download file handle slot used by executor.
573    pub(crate) fn download_file_slot(&self) -> &Arc<Mutex<Option<tokio::fs::File>>> {
574        &self.download_file_slot
575    }
576
577    /// Returns the configured max concurrent parts for this task.
578    pub(crate) fn max_parts_in_flight(&self) -> usize {
579        self.max_parts_in_flight
580    }
581
582    /// Returns the shared concurrent-download progress slot.
583    pub(crate) fn download_progress(
584        &self,
585    ) -> &Arc<StdMutex<Option<crate::dflt::download_progress::DownloadProgress>>> {
586        &self.download_progress
587    }
588
589    /// Writes a fully received parallel part through the unique locked target
590    /// handle, then stages its digest. If the checkpoint batch is due, the same
591    /// handle is synced before the sidecar publishes the part bit.
592    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    /// Stages a serial part after its streaming write released the file mutex.
667    /// Dropping that mutex first keeps the global order `barrier -> file` and
668    /// prevents the timer from deadlocking with a serial network response.
669    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    /// Establishes the serial layout invariant: the visible target contains
680    /// exactly the longest committed prefix, and the sidecar contains no bits
681    /// beyond it. Target truncation is synced before the compacted snapshot is
682    /// published, so a crash can never expose bits for missing bytes.
683    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            // The progress object itself suppresses duplicate timers. The task
806            // may complete before this wake-up; then the shared slot is empty
807            // and the timer exits without touching the completed file.
808            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    /// Runs an externally coordinated checkpoint while the caller already owns
818    /// `download_checkpoint_barrier`.
819    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            // Windows whole-file locks can reject reads through a second handle
895            // in the same process. Temporarily move the exact locked handle to
896            // the blocking validator without closing or unlocking it. The
897            // handle was opened without FILE_SHARE_DELETE, so the visible path
898            // cannot be replaced while this content check runs.
899            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    /// Commits all pending digests, validates every committed range through the
976    /// visible path, removes the sidecar, and only then releases the actual
977    /// target lock. Both serial and parallel downloads use this exact sequence.
978    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        // Always release the OS lock, but never let an unlock error mask the
1034        // content/checkpoint failure that made completion unsafe.
1035        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}