use std::collections::{HashMap, HashSet};
use std::path::{Path, PathBuf};
use std::sync::Arc;
use std::time::Duration;
use crate::downloads::sources::{
prepare_staging_directory, DownloadSink, DownloadStart, SourceHandle,
};
use crate::downloads::store::DownloadSource;
use crate::downloads::DownloadError;
pub const DEFAULT_POLL_INTERVAL: Duration = Duration::from_millis(250);
pub const IN_PROGRESS_SUFFIXES: [&str; 4] = ["crdownload", "tmp", "partial", "part"];
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct StagedDownload {
pub path: PathBuf,
pub name: String,
}
#[derive(Debug)]
pub struct DirectoryWatcher {
directory: PathBuf,
sizes: HashMap<String, u64>,
reported: HashSet<String>,
}
impl DirectoryWatcher {
pub fn new(directory: impl Into<PathBuf>) -> Self {
let directory = directory.into();
let reported = Self::candidates(&directory).into_iter().collect();
Self {
directory,
sizes: HashMap::new(),
reported,
}
}
fn candidates(directory: &Path) -> Vec<String> {
let Ok(entries) = std::fs::read_dir(directory) else {
return Vec::new();
};
let mut names: Vec<String> = entries
.filter_map(Result::ok)
.filter(|entry| entry.path().is_file())
.filter_map(|entry| entry.file_name().into_string().ok())
.filter(|name| !name.starts_with('.'))
.filter(|name| {
!IN_PROGRESS_SUFFIXES
.iter()
.any(|suffix| name.ends_with(&format!(".{suffix}")))
})
.collect();
names.sort();
names
}
pub fn poll_once(&mut self) -> Vec<StagedDownload> {
let mut claimed = Vec::new();
for name in Self::candidates(&self.directory) {
if self.reported.contains(&name) {
continue;
}
let path = self.directory.join(&name);
let Ok(size) = std::fs::metadata(&path).map(|metadata| metadata.len()) else {
continue;
};
if self.sizes.get(&name) != Some(&size) {
self.sizes.insert(name, size);
continue;
}
self.sizes.remove(&name);
self.reported.insert(name.clone());
claimed.push(StagedDownload { path, name });
}
claimed
}
}
pub fn attach_filesystem_watcher(
root: &Path,
sink: Arc<dyn DownloadSink>,
interval: Duration,
) -> Result<SourceHandle, DownloadError> {
let staging_directory = prepare_staging_directory(root)?;
let mut watcher = DirectoryWatcher::new(&staging_directory);
let (stop, mut stopped) = tokio::sync::oneshot::channel::<()>();
let task = tokio::spawn({
let sink = Arc::clone(&sink);
async move {
loop {
tokio::select! {
_ = &mut stopped => break,
_ = tokio::time::sleep(interval) => {}
}
report(&mut watcher, &sink).await;
}
report(&mut watcher, &sink).await;
}
});
Ok(SourceHandle {
staging_directory,
detach: Box::new(move || {
Box::pin(async move {
let _ = stop.send(());
let _ = task.await;
})
}),
})
}
async fn report(watcher: &mut DirectoryWatcher, sink: &Arc<dyn DownloadSink>) {
for staged in watcher.poll_once() {
let id = sink.started(DownloadStart {
engine_handle: staged.path.to_string_lossy().into_owned(),
url: Some(format!("file://{}", staged.path.to_string_lossy())),
suggested_filename: Some(staged.name),
mime_type: None,
});
sink.finished(id, DownloadSource::staged(staged.path)).await;
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::downloads::test_support::{RecordingSink, TempDir};
fn write(directory: &Path, name: &str, contents: &[u8]) {
std::fs::write(directory.join(name), contents).unwrap();
}
#[test]
fn claims_a_file_only_once_it_has_stopped_growing() {
let temp = TempDir::new("bc-watch-growing");
let mut watcher = DirectoryWatcher::new(temp.path());
write(temp.path(), "report.pdf", b"half");
assert!(watcher.poll_once().is_empty(), "claimed a growing file");
write(temp.path(), "report.pdf", b"half and the rest");
assert!(watcher.poll_once().is_empty(), "claimed a growing file");
let claimed = watcher.poll_once();
assert_eq!(claimed.len(), 1);
assert_eq!(claimed[0].name, "report.pdf");
assert_eq!(claimed[0].path, temp.path().join("report.pdf"));
}
#[test]
fn claims_a_download_exactly_once() {
let temp = TempDir::new("bc-watch-once");
let mut watcher = DirectoryWatcher::new(temp.path());
write(temp.path(), "report.pdf", b"body");
watcher.poll_once();
assert_eq!(watcher.poll_once().len(), 1);
assert!(
watcher.poll_once().is_empty(),
"claimed the same file twice"
);
}
#[test]
fn ignores_files_chromium_is_still_writing() {
let temp = TempDir::new("bc-watch-inprogress");
let mut watcher = DirectoryWatcher::new(temp.path());
write(temp.path(), "report.pdf.crdownload", b"body");
write(temp.path(), "notes.txt.tmp", b"body");
write(temp.path(), ".hidden", b"body");
watcher.poll_once();
assert!(watcher.poll_once().is_empty());
}
#[test]
fn ignores_files_that_were_there_before_the_session() {
let temp = TempDir::new("bc-watch-existing");
write(temp.path(), "from-yesterday.pdf", b"body");
let mut watcher = DirectoryWatcher::new(temp.path());
watcher.poll_once();
assert!(watcher.poll_once().is_empty());
}
#[test]
fn survives_a_directory_that_is_not_there() {
let temp = TempDir::new("bc-watch-missing");
let mut watcher = DirectoryWatcher::new(temp.path().join("never-created"));
assert!(watcher.poll_once().is_empty());
}
#[tokio::test]
async fn saves_what_the_browser_left_in_the_staging_directory() {
let temp = TempDir::new("bc-watch-attach");
let root = temp.path().join("downloads");
let sink = Arc::new(RecordingSink::default());
let handle = attach_filesystem_watcher(
&root,
Arc::clone(&sink) as Arc<dyn DownloadSink>,
Duration::from_millis(10),
)
.unwrap();
write(&handle.staging_directory, "report.pdf", b"%PDF-1.7");
tokio::time::sleep(Duration::from_millis(60)).await;
(handle.detach)().await;
let seen = sink.started_names();
assert_eq!(seen, vec!["report.pdf".to_string()]);
assert_eq!(sink.finished_count(), 1);
}
#[tokio::test]
async fn takes_one_last_look_when_it_is_detached() {
let temp = TempDir::new("bc-watch-last-look");
let root = temp.path().join("downloads");
let sink = Arc::new(RecordingSink::default());
let handle = attach_filesystem_watcher(
&root,
Arc::clone(&sink) as Arc<dyn DownloadSink>,
Duration::from_millis(50),
)
.unwrap();
write(&handle.staging_directory, "late.pdf", b"body");
tokio::time::sleep(Duration::from_millis(80)).await;
(handle.detach)().await;
assert_eq!(sink.started_names(), vec!["late.pdf".to_string()]);
}
}