pub mod source;
use self::source::{ArchiveBucketSource, BucketInput, BucketSourceAdapter, RepositoryBucketSource};
use crate::io::{
bagit::Bag,
config::transfer::{transfer_bucket_files, TransferManifest, TransferPolicy},
http::download_with_progress,
ApiResult, ArchiveCandidate,
};
use crate::util::constants::app::APPLICATION;
use acorn_core::{options::RuntimeOptions, util::MimeType};
use acorn_host::fs::{file_uri_to_path, TemporaryDirectory};
use bon::Builder;
use color_eyre::eyre::eyre;
use futures::{future, stream, StreamExt, TryStreamExt};
use serde::{Deserialize, Serialize};
use serde_with::skip_serializing_none;
pub use source::{BagArchive, BucketSource};
use std::path::{Path, PathBuf};
pub type BucketOptions = RuntimeOptions<BucketExtension, PathBuf>;
#[skip_serializing_none]
#[derive(Builder, Clone, Debug, Deserialize, Serialize)]
#[serde(rename_all = "camelCase", try_from = "BucketInput")]
#[builder(start_fn = init)]
pub struct Bucket {
pub name: Option<String>,
pub description: Option<String>,
#[serde(rename = "codeRepository")]
pub source: BucketSource,
}
#[derive(Clone, Debug, Default, Deserialize, Serialize)]
#[serde(default, rename_all = "camelCase")]
pub struct BucketExtension {
pub canonicalize: bool,
pub clobber: bool,
pub database_path: Option<PathBuf>,
pub flatten: bool,
pub no_local_database: bool,
pub strict: bool,
}
#[derive(Clone, Debug, Default)]
pub struct Buckets(pub Vec<Bucket>);
#[derive(Clone, Debug)]
pub enum SourceInput {
Bucket(Bucket),
Uri(String),
}
impl From<Bucket> for SourceInput {
fn from(value: Bucket) -> Self {
Self::Bucket(value)
}
}
impl From<String> for SourceInput {
fn from(value: String) -> Self {
Self::Uri(value)
}
}
#[derive(Debug)]
pub struct ResolvedBucket {
bucket: Bucket,
provenance: Option<String>,
_temporary: Option<TemporaryDirectory>,
}
impl ResolvedBucket {
pub async fn transfer(self, options: &BucketOptions, policy: TransferPolicy) -> ApiResult<usize> {
self.bucket.transfer(options, policy, self.provenance.as_deref()).await
}
}
impl SourceInput {
pub async fn resolve(self, format: Option<MimeType>, offline: bool) -> ApiResult<ResolvedBucket> {
match self {
| Self::Bucket(bucket) => Ok(ResolvedBucket {
bucket,
provenance: None,
_temporary: None,
}),
| Self::Uri(source) => {
let path = file_uri_to_path(&source).unwrap_or_else(|_| PathBuf::from(&source));
let inferred = MimeType::from(source.as_str());
let forced = format.is_some();
match (path.is_dir(), path.is_file(), forced, inferred) {
| (true, _, false, _) => match path.canonicalize() {
| Ok(path) => Ok(ResolvedBucket {
bucket: Bucket::from(path.as_path()),
provenance: Some(source),
_temporary: None,
}),
| Err(why) => Err(why.into()),
},
| (_, true, _, _) | (_, _, true, _) | (_, _, _, MimeType::Gzip | MimeType::SevenZip | MimeType::Tar | MimeType::Zip) => {
match TemporaryDirectory::create(&format!("{APPLICATION}-bagit")) {
| Ok(temporary) => match Bucket::archive(&source, format, &temporary, offline).await {
| Ok((bucket, provenance)) => Ok(ResolvedBucket {
bucket,
provenance: Some(provenance),
_temporary: Some(temporary),
}),
| Err(why) => Err(why),
},
| Err(why) => Err(why.into()),
}
}
| (_, _, _, _) if source.contains("://") => match Bucket::try_from(source.as_str()) {
| Ok(bucket) => Ok(ResolvedBucket {
bucket,
provenance: Some(source),
_temporary: None,
}),
| Err(why) => Err(why),
},
| _ => Err(eyre!("Source does not exist or has an unsupported format: {source}")),
}
}
}
}
}
impl Bucket {
pub async fn archive(source: &str, format: Option<MimeType>, temporary: &TemporaryDirectory, offline: bool) -> ApiResult<(Self, String)> {
let local = file_uri_to_path(source).unwrap_or_else(|_| PathBuf::from(source));
let archive_path = match local.is_file() {
| true => Ok(local),
| false if offline => Err(eyre!("Offline mode cannot transfer remote archive: {source}")),
| false => {
let path = temporary.path().join("archive");
download_with_progress(source, &path, |_, _| {}, None, None, None, None)
.await
.map(|_| path)
}
};
archive_path
.and_then(|path| ArchiveCandidate::from(path).extract(Some(temporary.path().join("extracted")), format))
.and_then(|root| bag_root(&root))
.and_then(Bag::verify_payload)
.map(|payload| (Self::from_payload(source, payload), source.to_string()))
}
pub async fn transfer(self, options: &BucketOptions, policy: TransferPolicy, provenance: Option<&str>) -> ApiResult<usize> {
let local = self.source.is_local();
match (policy, &self.source, options.common.offline) {
| (TransferPolicy::Download, source, _) if source.is_local_repository() => {
Err(eyre!("Bucket download requires a remote repository source"))
}
| (_, source, true) if !source.is_local() => Err(eyre!("Offline mode cannot transfer remote bucket sources")),
| _ => match local {
| true => self.copy_files(options).await,
| false => self.download_files(options).await,
}
.and_then(|manifest| {
let manifest = match provenance {
| Some(source) => TransferManifest {
repository: source.to_string(),
..manifest
},
| None => manifest,
};
let count = manifest.count();
manifest
.canonicalize(options)
.and_then(|manifest| manifest.ingest(options))
.map(|()| count)
}),
}
}
pub async fn import(self, options: &BucketOptions, provenance: Option<&str>) -> ApiResult<usize> {
self.transfer(options, TransferPolicy::Import, provenance).await
}
pub async fn download(self, options: &BucketOptions, provenance: Option<&str>) -> ApiResult<usize> {
self.transfer(options, TransferPolicy::Download, provenance).await
}
#[cfg(test)]
pub(crate) fn domain(&self) -> ApiResult<String> {
match &self.source {
| BucketSource::Repository(repository) => RepositoryBucketSource::new(repository).domain(),
| BucketSource::Archive(_) => Err(eyre!("Archive payloads have no repository domain")),
}
}
pub async fn copy_files(self, options: &BucketOptions) -> ApiResult<TransferManifest> {
let Self { name, source, .. } = self;
match source {
| BucketSource::Repository(repository) => {
let adapter = RepositoryBucketSource::new(&repository);
match adapter.is_local() {
| true => transfer_bucket_files(name, &adapter, options).await,
| false => Ok(TransferManifest {
bucket: name,
repository: adapter.manifest_identity(),
files: Vec::new(),
}),
}
}
| BucketSource::Archive(archive) => {
let adapter = ArchiveBucketSource::new(&archive);
transfer_bucket_files(name, &adapter, options).await
}
}
}
pub async fn download_files(self, options: &BucketOptions) -> ApiResult<TransferManifest> {
let Self { name, source, .. } = self;
match source {
| BucketSource::Repository(repository) => {
let adapter = RepositoryBucketSource::new(&repository);
transfer_bucket_files(name, &adapter, options).await
}
| BucketSource::Archive(archive) => {
let adapter = ArchiveBucketSource::new(&archive);
transfer_bucket_files(name, &adapter, options).await
}
}
}
}
impl Buckets {
pub async fn transfer(self, options: &BucketOptions, policy: TransferPolicy) -> ApiResult<usize> {
match options.output {
| Some(_) => {
stream::iter(self.0)
.then(|bucket| bucket.transfer(options, policy, None))
.try_fold(0_usize, |total, count| future::ready(Ok(total.saturating_add(count))))
.await
}
| None => Ok(0),
}
}
}
impl From<Vec<Bucket>> for Buckets {
fn from(value: Vec<Bucket>) -> Self {
Self(value)
}
}
fn bag_root(root: &Path) -> ApiResult<PathBuf> {
match root.join("bagit.txt").is_file() {
| true => Ok(root.to_path_buf()),
| false => root
.read_dir()
.map_err(Into::into)
.map(|entries| {
entries
.filter_map(Result::ok)
.filter_map(|entry| entry.file_type().ok().filter(|kind| kind.is_dir()).map(|_| entry.path()))
.collect::<Vec<_>>()
})
.and_then(|directories| match directories.as_slice() {
| [directory] if directory.join("bagit.txt").is_file() => Ok(directory.clone()),
| _ => Err(eyre!("BagIt metadata must be at the archive root or beneath one enclosing directory")),
}),
}
}