use crate::io::api::{Configuration, Endpoint, Param, Params, RemoteResource};
use crate::io::enrichment::{Enrich, ResearchOutputMetadata};
use crate::io::{workflow, ApiResult};
use crate::util::constants::app::{DEFAULT_OPENALEX_DOMAIN, OPENALEX_WORK_FIELDS};
use crate::util::constants::env::{OPENALEX_API_HOST, OPENALEX_API_KEY};
use acorn_core::options::{ApiExtension, ApiOptions};
use acorn_core::prelude::{format, String, ToString, Vec};
use acorn_core::util::merge_unique_by;
use acorn_core::Location;
use acorn_schema::pid::{PersistentIdentifierParse, DOI};
use acorn_schema::research_activity::output::{Access, Affiliation, Award, Contributor, Funding};
use acorn_schema::research_activity::ResearchOutput;
use acorn_schema::Website;
use async_trait::async_trait;
use color_eyre::eyre::eyre;
use futures::stream::{self, StreamExt};
use secrecy::ExposeSecret;
use serde::{Deserialize, Serialize};
pub type Options = ApiOptions<Extension, Param>;
#[derive(Clone, Debug, Deserialize, Serialize)]
pub struct Author {
pub id: String,
pub display_name: String,
pub orcid: Option<String>,
}
#[derive(Clone, Debug, Deserialize, Serialize)]
pub struct Authorship {
pub author: Author,
pub author_position: String,
pub is_corresponding: bool,
#[serde(default)]
pub institutions: Vec<Institution>,
}
#[derive(Clone, Debug, Default)]
pub struct Extension;
#[derive(Clone, Debug, Deserialize, Serialize)]
pub struct Grant {
pub funder: Option<String>,
pub funder_display_name: Option<String>,
pub award_id: Option<String>,
}
#[derive(Clone, Debug, Deserialize, Serialize)]
pub struct Institution {
pub id: String,
pub display_name: String,
pub ror: Option<String>,
}
#[derive(Clone, Debug, Deserialize, Serialize)]
pub struct Keyword {
pub display_name: String,
pub score: f64,
}
#[derive(Clone, Debug, Deserialize, Serialize)]
pub struct OpenAccess {
pub is_oa: bool,
pub oa_status: String,
pub oa_url: Option<String>,
}
#[derive(Clone, Debug)]
pub struct Provider {
options: Options,
}
#[derive(Clone, Debug, Deserialize, Serialize)]
pub struct Work {
pub id: String,
pub doi: Option<String>,
pub title: String,
#[serde(rename = "type")]
pub kind: String,
pub publication_date: Option<String>,
#[serde(default)]
pub authorships: Vec<Authorship>,
#[serde(default)]
pub grants: Vec<Grant>,
pub open_access: Option<OpenAccess>,
pub primary_location: Option<WorkLocation>,
pub best_oa_location: Option<WorkLocation>,
#[serde(default)]
pub keywords: Vec<Keyword>,
}
#[derive(Clone, Debug, Deserialize, Serialize)]
pub struct WorkLocation {
pub landing_page_url: Option<String>,
pub pdf_url: Option<String>,
pub license: Option<String>,
}
impl From<Institution> for Affiliation {
fn from(institution: Institution) -> Self {
Affiliation::init().name(institution.display_name).maybe_ror(institution.ror).build()
}
}
impl From<Authorship> for Contributor {
fn from(authorship: Authorship) -> Self {
Contributor::init()
.name(authorship.author.display_name)
.maybe_orcid(
authorship
.author
.orcid
.map(|value| value.rsplit('/').next().unwrap_or(&value).to_string()),
)
.corresponding(authorship.is_corresponding)
.position(authorship.author_position)
.affiliations(authorship.institutions.into_iter().map(Affiliation::from).collect())
.build()
}
}
impl ApiExtension for Extension {
fn default_domain() -> String {
DEFAULT_OPENALEX_DOMAIN.to_string()
}
fn env_token_var() -> &'static str {
OPENALEX_API_KEY
}
fn env_domain_var() -> &'static str {
OPENALEX_API_HOST
}
}
impl From<Grant> for Option<Funding> {
fn from(grant: Grant) -> Self {
grant.funder_display_name.map(|name| {
let awards = grant
.award_id
.filter(|value| !value.trim().is_empty())
.map(|identifier| vec![Award::init().identifier(identifier).build()])
.unwrap_or_default();
Funding::init().name(name).awards(awards).build()
})
}
}
impl Provider {
pub fn new(options: Options) -> Self {
Self { options }
}
pub fn from_env() -> Self {
Self::new(Options::from_env())
}
}
#[async_trait]
impl<T> Enrich<(Location, T)> for Provider
where
T: ResearchOutputMetadata + 'static,
{
type Output = ApiResult<workflow::EnrichmentResult<T>>;
async fn enrich(&self, (_, data): (Location, T)) -> Self::Output {
let result = stream::iter(data.publication_identifiers())
.fold(
workflow::EnrichmentResult {
data: data.research_outputs(),
conflicts: Vec::new(),
failures: Vec::new(),
},
|result, identifier| async move {
let normalized = DOI::format(&identifier);
match normalized.is_empty() {
| true => result,
| false => {
let record = work(&self.options.clone().with_identifier(normalized.clone())).await;
merge_work_result(result, normalized, record)
}
}
},
)
.await;
Ok(workflow::EnrichmentResult {
data: data.with_research_outputs(result.data),
conflicts: result.conflicts,
failures: result.failures,
})
}
}
impl From<Work> for ResearchOutput {
fn from(work: Work) -> Self {
let contributors = work.authorships.into_iter().map(Contributor::from).collect::<Vec<_>>();
let funding = work.grants.into_iter().filter_map(Option::<Funding>::from).collect::<Vec<_>>();
let websites = merge_unique_by(
work.primary_location
.as_ref()
.into_iter()
.flat_map(|location| location.websites("Primary location")),
work.best_oa_location
.as_ref()
.into_iter()
.flat_map(|location| location.websites("Open-access location")),
|left, right| left.url == right.url,
);
let access = work.open_access.map(|value| {
Access::init()
.open(value.is_oa)
.status(value.oa_status)
.maybe_url(value.oa_url)
.maybe_license(work.best_oa_location.and_then(|location| location.license))
.build()
});
ResearchOutput::init()
.identifier(work.id)
.maybe_doi(work.doi.map(|value| DOI::format(&value)))
.title(work.title)
.kind(work.kind)
.maybe_publication_year(
work.publication_date
.as_ref()
.and_then(|value| value.split('-').next()?.parse::<i32>().ok()),
)
.maybe_contributors((!contributors.is_empty()).then_some(contributors))
.maybe_funding((!funding.is_empty()).then_some(funding))
.maybe_websites((!websites.is_empty()).then_some(websites))
.keywords(work.keywords.into_iter().map(|value| value.display_name).collect())
.maybe_access(access)
.build()
}
}
impl WorkLocation {
fn websites(&self, description: &str) -> Vec<Website> {
let Self {
landing_page_url, pdf_url, ..
} = self;
merge_unique_by(
[],
[landing_page_url, pdf_url]
.into_iter()
.flatten()
.map(|url| Website::at(url).description(description).build()),
|left, right| left.url == right.url,
)
}
}
pub(crate) fn merge_work_result(
result: workflow::EnrichmentResult<Vec<ResearchOutput>>,
identifier: String,
record: ApiResult<Work>,
) -> workflow::EnrichmentResult<Vec<ResearchOutput>> {
match record {
| Ok(record) => {
let candidate = ResearchOutput::from(record);
let exists = result
.data
.iter()
.filter_map(|output| output.doi.as_ref())
.any(|doi| DOI::format(doi).eq_ignore_ascii_case(&identifier));
let data = match exists {
| true => result
.data
.into_iter()
.map(
|output| match output.doi.as_ref().is_some_and(|doi| DOI::format(doi).eq_ignore_ascii_case(&identifier)) {
| true => output.merge(candidate.clone()),
| false => output,
},
)
.collect(),
| false => result.data.into_iter().chain([candidate]).collect(),
};
workflow::EnrichmentResult { data, ..result }
}
| Err(why) => workflow::EnrichmentResult {
failures: result.failures.into_iter().chain([format!("{identifier}: {why}")]).collect(),
..result
},
}
}
pub async fn work(options: &Options) -> ApiResult<Work> {
let template = "openalex";
let action = "work";
let identifier = options.identifier.as_ref().map(DOI::format).filter(|value| !value.is_empty());
match identifier.ok_or_else(|| eyre!("An OpenAlex work lookup requires a DOI")) {
| Ok(identifier) => {
let params = Params::new()
.with_template("identifier", Some(identifier.as_str()))
.with_keyvalue("select", Some(OPENALEX_WORK_FIELDS))
.with_api_key(ExposeSecret::expose_secret(&options.token))
.with_custom(options.params())
.build();
match Endpoint::from_template(template).map(|endpoint| endpoint.with_domain(options.domain())) {
| Ok(endpoint) => endpoint
.handle::<Work>(endpoint.invoke(action, Some(params)).await)
.map_err(|why| eyre!("Failed to retrieve OpenAlex work — {why}")),
| Err(why) => Err(why),
}
}
| Err(why) => Err(why),
}
}