use crate::{
io::{
adapter::migrate::Migration,
config::{
bucket::{
source::{BucketSourceAdapter, BucketSourceItem},
BucketExtension, BucketOptions,
},
is_filtered_path, is_ignored_path, FilterSet,
},
database::{schema::Table, Database, Provenance, ResearchActivityCandidate},
read_file, with_progress, ApiResult, FromPath, InputOutput, ProgressType,
},
util::{constants::app::SUPPORTED_RAD_FILETYPES, is_filetype},
};
use acorn_core::{
options::{CommonRuntime, RuntimeOptions},
prelude::HashMap,
util::{suffix, MimeType},
};
use acorn_host::{
fs::{FsError, SafePath, TemporaryDirectory},
terminal::Label,
};
use acorn_schema::{
pid::{Identifier, PID},
research_activity::ResearchActivity,
};
use color_eyre::eyre::{eyre, Report};
use core::iter::once;
use jiff::Timestamp;
use owo_colors::OwoColorize;
use std::{
fs::read,
path::{Path, PathBuf},
sync::Mutex,
};
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub enum TransferPolicy {
Download,
Import,
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub(crate) struct TransferItem {
pub(crate) source: BucketSourceItem,
pub(crate) destination: SafePath,
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct TransferManifest {
pub bucket: Option<String>,
pub repository: String,
pub files: Vec<PathBuf>,
}
impl TransferItem {
async fn apply<A>(self, adapter: &A, output: &Path, staging: Option<&TemporaryDirectory>, clobber: bool, write_lock: &Mutex<()>) -> ApiResult<()>
where
A: BucketSourceAdapter + Sync,
{
match (adapter.is_local(), clobber, staging) {
| (true, true, Some(staging)) => match self.destination.materialize_under(staging.path()) {
| Ok(staged) => match adapter.retrieve(&self.source, &staged).await {
| Ok(()) => self
.destination
.write_file(output, true, write_lock, || async {
read(&staged).map_err(|source| FsError::Io { path: staged, source })
})
.await
.map_err(Report::from),
| Err(why) => Err(why),
},
| Err(why) => Err(why.into()),
},
| (false, true, _) => {
let target = match write_lock.lock() {
| Ok(_guard) => self.destination.prepare(output).map_err(Report::from),
| Err(why) => Err(eyre!("Failed to lock bucket output — {why}")),
};
match target {
| Ok(target) => adapter.retrieve(&self.source, &target).await,
| Err(why) => Err(why),
}
}
| _ => match self.destination.materialize_under(output) {
| Ok(target) => adapter.retrieve(&self.source, &target).await,
| Err(why) => Err(why.into()),
},
}
}
pub(crate) fn collect(items: Vec<BucketSourceItem>, flatten: bool) -> ApiResult<Vec<Self>> {
let items = items
.into_iter()
.map(|source| Self {
destination: source.path.clone(),
source,
})
.map(|item| match flatten {
| true => item.flatten(),
| false => Ok(item),
})
.collect::<ApiResult<Vec<_>>>();
match items {
| Ok(items) => items
.iter()
.try_fold(HashMap::<SafePath, SafePath>::new(), |mut destinations, item| {
let destination = &item.destination;
match destinations.insert(destination.clone(), item.source.path.clone()) {
| Some(source) => Err(eyre!(
"Output path collision for '{}' — '{}' and '{}'",
destination.as_path().display(),
source.as_path().display(),
item.source.path.as_path().display()
)),
| None => Ok(destinations),
}
})
.map(|_| items),
| Err(why) => Err(why),
}
}
fn flatten(self) -> ApiResult<Self> {
match self.source.path.as_path().file_name().map(PathBuf::from).map(SafePath::new) {
| Some(Ok(destination)) => Ok(Self { destination, ..self }),
| Some(Err(why)) => Err(why.into()),
| None => Err(eyre!(
"Cannot flatten bucket path without a filename — {}",
self.source.path.as_path().display()
)),
}
}
}
impl TransferManifest {
pub fn count(&self) -> usize {
let paths = self.files.iter().map(|path| path.display().to_string()).collect::<Vec<_>>();
count_json_files(&paths).saturating_add(count_image_files(&paths))
}
pub(crate) fn canonicalize(self, options: &BucketOptions) -> ApiResult<Self> {
let RuntimeOptions { extension, output, .. } = options;
let BucketExtension { canonicalize, .. } = extension;
match canonicalize {
| false => Ok(self),
| true => {
let output = output.as_deref().unwrap_or_else(|| Path::new(""));
let canonicalized = self
.files
.iter()
.filter(is_filetype(SUPPORTED_RAD_FILETYPES))
.try_fold((), |(), relative| {
let path = output.join(relative);
ResearchActivity::migrate(path.clone())
.map(|_| ())
.map_err(|why| eyre!("Failed to canonicalize transferred RAD {} — {why}", path.display()))
});
canonicalized.map(|()| self)
}
}
}
pub fn ingest(&self, options: &BucketOptions) -> ApiResult<()> {
let RuntimeOptions { extension, output, .. } = options;
let BucketExtension {
database_path,
no_local_database,
strict,
..
} = extension;
match no_local_database {
| true => Ok(()),
| false => {
let output = output.as_deref().unwrap_or_else(|| Path::new(""));
let database = Database::<Table>::from_path(database_path.as_ref());
self.files
.iter()
.filter(is_filetype(SUPPORTED_RAD_FILETYPES))
.filter(|relative| {
let path = output.join(relative);
MimeType::from_path(&path) != MimeType::Markdown || read_file(&path).is_ok_and(ResearchActivity::is_markdown)
})
.try_fold((), |(), relative| {
let path = output.join(relative);
match strict {
| true => ResearchActivity::read_strict(path.clone()),
| false => ResearchActivity::read(path.clone()),
}
.and_then(|rad| {
serde_json::to_value(&rad)
.map_err(Report::from)
.and_then(|rad_json| database.create_or_enrich(self.candidate(&rad, rad_json, relative)).map(|_| ()))
})
.map_err(|why| eyre!("Failed to ingest transferred RAD {} — {why}", path.display()))
})
}
}
}
fn candidate(&self, rad: &ResearchActivity, rad_json: serde_json::Value, relative: &Path) -> ResearchActivityCandidate {
let pairs = [
(PID::DOI, rad.meta.doi.as_ref()),
(PID::Handle, rad.meta.handle.as_ref()),
(PID::ISBN, rad.meta.books.as_ref()),
(PID::Patent, rad.meta.patents.as_ref()),
(PID::RAID, rad.meta.raid.as_ref()),
(PID::SWHID, rad.meta.swhid.as_ref()),
];
let pid_keys = pairs.into_iter().flat_map(|(kind, values)| {
values.into_iter().flatten().filter_map(move |value| {
Identifier::init()
.kind(kind.clone())
.value(value)
.build()
.normalized()
.map(|identifier| format!("{}:{}", identifier.kind.as_str(), identifier.value))
})
});
let rad_key = format!("rad:{}:{}", self.repository, rad.meta.identifier);
let prov = Provenance::Bucket {
bucket: self.bucket.clone(),
repository: self.repository.clone(),
relative_path: relative.display().to_string(),
observed_at: Timestamp::now().to_string(),
};
ResearchActivityCandidate::new(
rad_json,
pid_keys.chain(once(rad_key)).collect(),
vec![serde_json::to_value(prov).unwrap_or_default()],
)
}
}
pub(crate) fn count_image_files(paths: &[String]) -> usize {
paths.iter().filter(|&x| has_image_extension(x)).count()
}
pub(crate) fn count_json_files(paths: &[String]) -> usize {
paths.iter().filter(|&path| path.to_lowercase().ends_with(".json")).count()
}
#[allow(clippy::ptr_arg)]
pub(crate) fn has_image_extension(path: &String) -> bool {
path.to_lowercase().ends_with(".png") || path.to_lowercase().ends_with(".jpg")
}
pub(crate) fn operations_complete_message(name: Option<String>, json_count: usize, image_count: usize) -> String {
let total = json_count.saturating_add(image_count);
let message = if json_count != image_count {
let recommendation = if json_count > image_count {
"Do you need to add some images?"
} else {
"Do you need to add some JSON files?"
};
format!(
" ({} data file{}, {} image{} - {})",
json_count.yellow(),
suffix(json_count),
image_count.yellow(),
suffix(image_count),
recommendation.italic(),
)
} else {
"".to_string()
};
let bucket_description = match name {
| Some(value) => format!("{} bucket", value.to_uppercase().cyan()),
| None => "<URL>".cyan().to_string(),
};
format!(
"{}Obtained {} file{} from {bucket_description}{}",
if total > 0 { Label::CHECKMARK } else { Label::CAUTION },
if total > 0 {
total.green().to_string()
} else {
total.yellow().to_string()
},
suffix(total),
message,
)
}
pub(crate) async fn transfer_bucket_files<A>(name: Option<String>, adapter: &A, options: &BucketOptions) -> ApiResult<TransferManifest>
where
A: BucketSourceAdapter + Sync,
{
match (FilterSet::compile(options), adapter.items().await) {
| (Ok(filters), Ok(items)) => {
let selected = items
.into_iter()
.filter(|item| {
let path = item.path.as_path().display().to_string();
!is_ignored_path(&path, &filters.ignore) && is_filtered_path(&path, &filters.filter)
})
.collect::<Vec<_>>();
match TransferItem::collect(selected, options.extension.flatten) {
| Ok(items) => transfer_selected(name, adapter, items, options).await,
| Err(why) => Err(why),
}
}
| (Err(why), _) | (_, Err(why)) => Err(why),
}
}
async fn transfer_selected<A>(name: Option<String>, adapter: &A, items: Vec<TransferItem>, options: &BucketOptions) -> ApiResult<TransferManifest>
where
A: BucketSourceAdapter + Sync,
{
let CommonRuntime { quiet, threads, .. } = &options.common;
let source_paths = items
.iter()
.map(|item| item.source.path.as_path().display().to_string())
.collect::<Vec<_>>();
let total_data = count_json_files(&source_paths);
let total_images = count_image_files(&source_paths);
let verb = match adapter.is_local() {
| true => "Copying",
| false => "Downloading",
};
let message = move |item: &TransferItem| format!("{verb} {}", item.source.path.as_path().display());
let finish_name = name.clone();
let finish_message = |_| operations_complete_message(finish_name, total_data, total_images);
let progress_type = match quiet {
| true => ProgressType::Silent,
| false => ProgressType::Bar,
};
let files = items.iter().map(|item| item.destination.clone().into_path_buf()).collect::<Vec<_>>();
let output = options.output.clone().unwrap_or_default();
let staging = match (adapter.is_local(), options.extension.clobber, items.is_empty()) {
| (true, true, false) => TemporaryDirectory::create("acorn-bucket-transfer").map(Some).map_err(Report::from),
| _ => Ok(None),
};
match staging {
| Ok(staging) => {
let write_lock = Mutex::new(());
let operation = |item: TransferItem| item.apply(adapter, &output, staging.as_ref(), options.extension.clobber, &write_lock);
with_progress(items, message, operation, finish_message, Some(*threads), progress_type)
.await
.map(|_| TransferManifest {
bucket: name,
repository: adapter.manifest_identity(),
files,
})
}
| Err(why) => Err(why),
}
}