use super::backend::{self, BackendRow, Connection};
use super::schema::{ResearchActivityRow, Table};
use super::{transaction, Database, Row};
use crate::io::api::handle::Record;
use crate::io::ApiResult;
use acorn_core::util::to_rfc3339;
use acorn_schema::discovery::RemoteEntity;
use acorn_schema::pid;
use alloc::collections::BTreeSet;
use color_eyre::eyre::eyre;
use jiff::Timestamp;
use nanoid::nanoid;
use serde::{Deserialize, Serialize};
use serde_json::{Map, Value};
#[derive(Serialize)]
#[serde(tag = "provider", rename_all = "lowercase")]
pub(crate) enum Provenance {
Bucket {
bucket: Option<String>,
repository: String,
relative_path: String,
observed_at: String,
},
Citeas {
title: String,
doi: String,
project_url: String,
},
Discovery {
source: String,
source_format: String,
identifier: String,
resolution_status: String,
observed_at: String,
},
Handle {
record: Record,
},
Osti {
provider_identifier: String,
entity: RemoteEntity,
metadata: Value,
observed_at: String,
},
Raid {
records: Vec<pid::raid::Metadata>,
},
}
#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
pub struct ResearchActivityCandidate {
pub rad_json: Value,
pub identity_keys: Vec<String>,
pub provenance: Vec<Value>,
}
impl ResearchActivityCandidate {
pub fn new(rad_json: Value, identity_keys: Vec<String>, provenance: Vec<Value>) -> Self {
Self {
rad_json: canonicalize(rad_json),
identity_keys: normalize_identity_keys(identity_keys),
provenance: deduplicate_provenance(provenance),
}
}
}
#[derive(Clone, Copy, Debug, Deserialize, Eq, PartialEq, Serialize)]
#[serde(rename_all = "lowercase")]
pub enum CandidateAction {
Created,
Enriched,
Unchanged,
Conflict,
}
#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
pub struct CandidatePersistence {
pub iid: Option<String>,
pub action: CandidateAction,
}
#[derive(Clone, Debug, Eq, PartialEq)]
struct MergeResult {
rad_json: Value,
identity_keys: Vec<String>,
provenance: Vec<Value>,
conflicts: Vec<Value>,
changed: bool,
}
trait ConnectionExt {
fn persist_candidate(&self, candidate: ResearchActivityCandidate) -> ApiResult<CandidatePersistence>;
fn create_candidate(&self, candidate: ResearchActivityCandidate) -> ApiResult<CandidatePersistence>;
fn merge_candidate(&self, row: &ResearchActivityRow, candidate: ResearchActivityCandidate) -> ApiResult<CandidatePersistence>;
fn research_activities(&self) -> ApiResult<Vec<ResearchActivityRow>>;
fn unique_research_activity_iid(&self) -> ApiResult<String>;
}
impl Database<Table> {
pub fn research_activity_by_iid(&self, iid: &str) -> ApiResult<Option<ResearchActivityRow>> {
self.migrate_table(Table::ResearchActivities).and_then(|()| {
self.with_connection(|conn| {
conn.query_row(
"SELECT id, iid, rad_json, identity_keys_json, provenance_json, created_at, updated_at FROM research_activities WHERE iid = ?",
backend::params![iid],
|row| Ok(ResearchActivityRow::from(row)),
)
.map(Some)
.or_else(|why| match is_not_found(&why) {
| true => Ok(None),
| false => Err(eyre!("Failed to select research activity {iid} — {why}")),
})
})
})
}
pub fn research_activities_by_identity(&self, identity_keys: &[String]) -> ApiResult<Vec<ResearchActivityRow>> {
let keys = normalize_identity_keys(identity_keys.to_vec());
self.migrate_table(Table::ResearchActivities)
.and_then(|()| self.with_connection(|conn| conn.research_activities().and_then(|rows| matching_rows(rows, &keys))))
}
pub fn create_or_enrich(&self, candidate: ResearchActivityCandidate) -> ApiResult<CandidatePersistence> {
let candidate = ResearchActivityCandidate::new(candidate.rad_json, candidate.identity_keys, candidate.provenance);
match candidate.identity_keys.is_empty() {
| true => Err(eyre!("Research activity candidate has no project identity keys")),
| false => self
.migrate_table(Table::ResearchActivities)
.and_then(|()| self.with_connection(|connection| transaction(connection, |connection| connection.persist_candidate(candidate)))),
}
}
}
impl ConnectionExt for Connection {
fn persist_candidate(&self, candidate: ResearchActivityCandidate) -> ApiResult<CandidatePersistence> {
let matches = matching_rows(self.research_activities()?, &candidate.identity_keys)?;
match matches.as_slice() {
| [] => self.create_candidate(candidate),
| [row] => self.merge_candidate(row, candidate),
| _ => Ok(CandidatePersistence {
iid: None,
action: CandidateAction::Conflict,
}),
}
}
fn create_candidate(&self, candidate: ResearchActivityCandidate) -> ApiResult<CandidatePersistence> {
let iid = self.unique_research_activity_iid()?;
let now = Timestamp::now();
ResearchActivityRow::init()
.iid(iid.clone())
.rad_json(serialize(&candidate.rad_json)?)
.identity_keys_json(serialize(&candidate.identity_keys)?)
.provenance_json(serialize(&candidate.provenance)?)
.created_at(now)
.updated_at(now)
.build()
.insert(self)
.map(|_| CandidatePersistence {
iid: Some(iid),
action: CandidateAction::Created,
})
}
fn merge_candidate(&self, row: &ResearchActivityRow, candidate: ResearchActivityCandidate) -> ApiResult<CandidatePersistence> {
let iid = row.iid.clone().ok_or_else(|| eyre!("Research activity row is missing iid"))?;
let merged = merge_candidate(row, candidate)?;
let action = match (merged.conflicts.is_empty(), merged.changed) {
| (false, _) => CandidateAction::Conflict,
| (true, true) => CandidateAction::Enriched,
| (true, false) => CandidateAction::Unchanged,
};
match merged.changed {
| false => Ok(CandidatePersistence { iid: Some(iid), action }),
| true => self
.execute(
"UPDATE research_activities SET rad_json = ?, identity_keys_json = ?, provenance_json = ?, updated_at = ? WHERE iid = ?",
backend::params![
serialize(&merged.rad_json)?,
serialize(&merged.identity_keys)?,
serialize(&merged.provenance)?,
to_rfc3339(Timestamp::now()),
iid,
],
)
.map_err(|why| eyre!("Failed to enrich research activity — {why}"))
.map(|_| CandidatePersistence {
iid: row.iid.clone(),
action,
}),
}
}
fn research_activities(&self) -> ApiResult<Vec<ResearchActivityRow>> {
self.prepare("SELECT id, iid, rad_json, identity_keys_json, provenance_json, created_at, updated_at FROM research_activities")
.map_err(|why| eyre!("Failed to prepare research activity lookup — {why}"))
.and_then(|mut statement| {
statement
.query_map(backend::params![], |row: &BackendRow<'_>| Ok(ResearchActivityRow::from(row)))
.map_err(|why| eyre!("Failed to query research activities — {why}"))
.and_then(|rows| {
rows.collect::<core::result::Result<Vec<_>, _>>()
.map_err(|why| eyre!("Failed to read research activities — {why}"))
})
})
}
fn unique_research_activity_iid(&self) -> ApiResult<String> {
(0..8)
.map(|_| nanoid!())
.find_map(|iid| {
self.query_row(
"SELECT COUNT(*) FROM research_activities WHERE iid = ?",
backend::params![iid.clone()],
|row| row.get::<_, i64>(0),
)
.ok()
.filter(|count| *count == 0)
.map(|_| iid)
})
.ok_or_else(|| eyre!("Failed to generate a unique research activity iid"))
}
}
fn merge_candidate(row: &ResearchActivityRow, candidate: ResearchActivityCandidate) -> ApiResult<MergeResult> {
let existing_rad = parse_value(row.rad_json.as_deref(), "rad_json")?;
let existing_keys = parse_strings(row.identity_keys_json.as_deref(), "identity_keys_json")?;
let existing_provenance = parse_values(row.provenance_json.as_deref(), "provenance_json")?;
let (rad_json, conflicts) = merge_values(existing_rad.clone(), candidate.rad_json, String::new());
let identity_keys = normalize_identity_keys(existing_keys.clone().into_iter().chain(candidate.identity_keys).collect());
let conflict_provenance = conflicts.iter().cloned();
let provenance = deduplicate_provenance(
existing_provenance
.clone()
.into_iter()
.chain(candidate.provenance)
.chain(conflict_provenance)
.collect(),
);
let changed = rad_json != existing_rad || identity_keys != existing_keys || provenance != existing_provenance;
Ok(MergeResult {
rad_json,
identity_keys,
provenance,
conflicts,
changed,
})
}
fn merge_values(existing: Value, incoming: Value, path: String) -> (Value, Vec<Value>) {
match (existing, incoming) {
| (Value::Null, incoming) => (incoming, Vec::new()),
| (existing, Value::Null) => (existing, Vec::new()),
| (Value::Object(existing), Value::Object(incoming)) => {
let (value, conflicts) = incoming
.into_iter()
.fold((existing, Vec::new()), |(mut merged, conflicts), (key, value)| {
let child_path = if path.is_empty() { key.clone() } else { format!("{path}.{key}") };
let (value, child_conflicts) = match merged.remove(&key) {
| Some(current) => merge_values(current, value, child_path),
| None => (value, Vec::new()),
};
merged.insert(key, value);
(merged, conflicts.into_iter().chain(child_conflicts).collect())
});
(canonicalize(Value::Object(value)), conflicts)
}
| (Value::Array(existing), Value::Array(incoming)) => {
let merged = incoming.into_iter().fold(existing, |mut values, value| {
if !values.contains(&value) {
values.push(value);
}
values
});
(Value::Array(merged), Vec::new())
}
| (existing, incoming) if existing == incoming => (existing, Vec::new()),
| (existing, incoming) => {
let conflict = serde_json::json!({
"kind": "field-conflict",
"path": path,
"existing": existing,
"observed": incoming,
});
(conflict.get("existing").cloned().unwrap_or(Value::Null), vec![conflict])
}
}
}
fn canonicalize(value: Value) -> Value {
match value {
| Value::Object(values) => Value::Object(
values
.into_iter()
.map(|(key, value)| (key, canonicalize(value)))
.collect::<alloc::collections::BTreeMap<_, _>>()
.into_iter()
.collect::<Map<_, _>>(),
),
| Value::Array(values) => Value::Array(values.into_iter().map(canonicalize).collect()),
| value => value,
}
}
fn normalize_identity_keys(values: Vec<String>) -> Vec<String> {
values
.into_iter()
.map(|value| normalize_identity_key(&value))
.filter(|value| !value.is_empty())
.collect::<BTreeSet<_>>()
.into_iter()
.collect()
}
fn normalize_identity_key(value: &str) -> String {
let value = value.trim();
match value.split_once(':') {
| Some((prefix, identifier)) if matches!(prefix.to_ascii_lowercase().as_str(), "doi" | "raid" | "isbn" | "patent") => {
format!("{}:{}", prefix.to_ascii_lowercase(), identifier.trim().to_ascii_lowercase())
}
| Some((prefix, identifier)) => format!("{}:{}", prefix.to_ascii_lowercase(), identifier.trim()),
| None => value.to_string(),
}
}
fn deduplicate_provenance(values: Vec<Value>) -> Vec<Value> {
values.into_iter().map(canonicalize).fold(Vec::new(), |mut unique, value| {
let fingerprint = provenance_fingerprint(&value);
if !unique.iter().any(|existing| provenance_fingerprint(existing) == fingerprint) {
unique.push(value);
}
unique
})
}
fn provenance_fingerprint(value: &Value) -> Value {
match value {
| Value::Object(values) => Value::Object(
values
.iter()
.filter(|(key, _)| !matches!(key.as_str(), "timestamp" | "observed_at" | "discovered_at"))
.map(|(key, value)| (key.clone(), provenance_fingerprint(value)))
.collect(),
),
| Value::Array(values) => Value::Array(values.iter().map(provenance_fingerprint).collect()),
| value => value.clone(),
}
}
fn matching_rows(rows: Vec<ResearchActivityRow>, keys: &[String]) -> ApiResult<Vec<ResearchActivityRow>> {
rows.into_iter().try_fold(Vec::new(), |mut matches, row| {
parse_strings(row.identity_keys_json.as_deref(), "identity_keys_json").map(|row_keys| {
if row_keys.iter().any(|key| keys.contains(key)) {
matches.push(row);
}
matches
})
})
}
fn parse_value(value: Option<&str>, field: &str) -> ApiResult<Value> {
value
.ok_or_else(|| eyre!("Research activity row is missing {field}"))
.and_then(|value| serde_json::from_str(value).map_err(|why| eyre!("Research activity {field} is invalid — {why}")))
}
fn parse_strings(value: Option<&str>, field: &str) -> ApiResult<Vec<String>> {
value
.ok_or_else(|| eyre!("Research activity row is missing {field}"))
.and_then(|value| serde_json::from_str(value).map_err(|why| eyre!("Research activity {field} is invalid — {why}")))
}
fn parse_values(value: Option<&str>, field: &str) -> ApiResult<Vec<Value>> {
value
.ok_or_else(|| eyre!("Research activity row is missing {field}"))
.and_then(|value| serde_json::from_str(value).map_err(|why| eyre!("Research activity {field} is invalid — {why}")))
}
fn serialize(value: &impl Serialize) -> ApiResult<String> {
serde_json::to_string(value).map_err(|why| eyre!("Failed to serialize research activity candidate — {why}"))
}
fn is_not_found(error: &backend::Error) -> bool {
error.to_string().contains("Query returned no rows") || error.to_string().contains("no rows")
}
#[cfg(test)]
mod tests;