use std::collections::HashSet;
use std::fs;
use std::path::Path;
use crate::commands::keyring::bucket::client;
use crate::commands::keyring::bucket::store::BucketConfig;
pub(crate) use crate::core::data::extension_of;
pub(crate) struct PullTask {
pub key: String,
pub size: u64,
}
#[derive(Debug, Clone)]
pub(crate) struct TypeSummary {
pub extension: String,
pub count: usize,
pub total_bytes: u64,
}
pub(crate) struct PullTransformPlan {
pub tasks: Vec<PullTask>,
pub type_summary: Vec<TypeSummary>,
}
pub(crate) const PROCESSED_FILE_NAME: &str = ".processed";
pub(crate) fn load_checkpoint(local_output: &Path) -> Result<HashSet<String>, String> {
let path = local_output.join(PROCESSED_FILE_NAME);
let contents = match fs::read_to_string(&path) {
Ok(contents) => contents,
Err(err) if err.kind() == std::io::ErrorKind::NotFound => return Ok(HashSet::new()),
Err(err) => return Err(format!("failed to read {}: {err}", path.display())),
};
Ok(contents.lines().map(str::to_string).collect())
}
pub(crate) fn append_checkpoint(local_output: &Path, key: &str) -> Result<(), String> {
use std::io::Write;
let path = local_output.join(PROCESSED_FILE_NAME);
let mut file = fs::OpenOptions::new()
.create(true)
.append(true)
.open(&path)
.map_err(|err| format!("failed to open {}: {err}", path.display()))?;
writeln!(file, "{key}").map_err(|err| format!("failed to write {}: {err}", path.display()))
}
pub(crate) async fn gather_pending(
bucket_config: &BucketConfig,
secret: &str,
local_output: &Path,
) -> Result<PullTransformPlan, String> {
let done = load_checkpoint(local_output)?;
let entries = client::list_objects(bucket_config, secret, "", true).await?;
let mut tasks = Vec::with_capacity(entries.len());
let mut by_extension: std::collections::BTreeMap<String, (usize, u64)> =
std::collections::BTreeMap::new();
for entry in entries {
if entry.is_prefix || done.contains(&entry.key) {
continue;
}
let extension = extension_of(&entry.key);
let bucket = by_extension.entry(extension).or_insert((0, 0));
bucket.0 += 1;
bucket.1 += entry.size;
tasks.push(PullTask {
key: entry.key,
size: entry.size,
});
}
let type_summary = by_extension
.into_iter()
.map(|(extension, (count, total_bytes))| TypeSummary {
extension,
count,
total_bytes,
})
.collect();
Ok(PullTransformPlan {
tasks,
type_summary,
})
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn load_checkpoint_is_empty_for_a_fresh_directory() {
let dir = tempfile::tempdir().unwrap();
assert!(load_checkpoint(dir.path()).unwrap().is_empty());
}
#[test]
fn append_checkpoint_then_load_round_trips() {
let dir = tempfile::tempdir().unwrap();
append_checkpoint(dir.path(), "a.jpg").unwrap();
append_checkpoint(dir.path(), "b.zip").unwrap();
let loaded = load_checkpoint(dir.path()).unwrap();
assert!(loaded.contains("a.jpg"));
assert!(loaded.contains("b.zip"));
assert_eq!(loaded.len(), 2);
}
}