use async_trait::async_trait;
use crate::domain::{
error::DppError,
passport::{Passport, PassportId},
product_identity::ProductIdentity,
status::PassportStatus,
};
#[async_trait]
pub trait PassportRepository: Send + Sync {
async fn create(&self, passport: Passport) -> Result<Passport, DppError>;
async fn find_by_id(&self, id: PassportId) -> Result<Option<Passport>, DppError>;
async fn find_published_by_id(&self, id: PassportId) -> Result<Option<Passport>, DppError>;
async fn find_published_by_gtin(&self, gtin: &str) -> Result<Option<Passport>, DppError>;
async fn find_by_id_any_status(&self, id: PassportId) -> Result<Option<Passport>, DppError>;
async fn find_by_identity(
&self,
identity: &ProductIdentity,
) -> Result<Option<Passport>, DppError> {
let drafts = self
.list(Some(PassportStatus::Draft), None, None, u32::MAX, 0)
.await?;
let published = self
.list(Some(PassportStatus::Published), None, None, u32::MAX, 0)
.await?;
Ok(drafts
.into_iter()
.chain(published)
.find(|p| ProductIdentity::from_passport(p).as_ref() == Some(identity)))
}
async fn update(&self, passport: Passport) -> Result<Passport, DppError>;
async fn patch_fields(
&self,
id: PassportId,
delta: serde_json::Value,
) -> Result<Passport, DppError> {
let Some(mut passport) = self.find_by_id(id).await? else {
return Err(DppError::NotFound(id.to_string()));
};
let mut p_val = serde_json::to_value(&passport)
.map_err(|e| DppError::Internal(format!("serialize: {e}")))?;
if let (serde_json::Value::Object(pm), serde_json::Value::Object(dm)) = (&mut p_val, delta)
{
pm.extend(dm);
}
passport = serde_json::from_value(p_val)
.map_err(|e| DppError::Internal(format!("deserialize: {e}")))?;
self.update(passport).await
}
async fn update_status(
&self,
id: PassportId,
status: PassportStatus,
) -> Result<Passport, DppError>;
async fn list(
&self,
status: Option<PassportStatus>,
q: Option<&str>,
facility_id: Option<&str>,
limit: u32,
offset: u32,
) -> Result<Vec<Passport>, DppError>;
async fn count(
&self,
status: Option<PassportStatus>,
facility_id: Option<&str>,
) -> Result<u64, DppError>;
async fn create_batch(&self, passports: Vec<Passport>) -> Vec<Result<Passport, DppError>> {
let mut results = Vec::with_capacity(passports.len());
for passport in passports {
results.push(self.create(passport).await);
}
results
}
async fn update_batch(&self, passports: Vec<Passport>) -> Vec<Result<Passport, DppError>> {
let mut results = Vec::with_capacity(passports.len());
for passport in passports {
results.push(self.update(passport).await);
}
results
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::domain::passport::ManufacturerInfo;
use crate::domain::sector::Sector;
use std::collections::HashMap;
use std::sync::Mutex;
#[derive(Default)]
struct InMemoryRepo {
store: Mutex<HashMap<PassportId, Passport>>,
}
#[async_trait]
impl PassportRepository for InMemoryRepo {
async fn create(&self, passport: Passport) -> Result<Passport, DppError> {
self.store
.lock()
.unwrap()
.insert(passport.id, passport.clone());
Ok(passport)
}
async fn find_by_id(&self, id: PassportId) -> Result<Option<Passport>, DppError> {
Ok(self.store.lock().unwrap().get(&id).cloned())
}
async fn find_published_by_id(&self, id: PassportId) -> Result<Option<Passport>, DppError> {
self.find_by_id(id).await
}
async fn find_published_by_gtin(&self, _gtin: &str) -> Result<Option<Passport>, DppError> {
Ok(None)
}
async fn find_by_id_any_status(
&self,
id: PassportId,
) -> Result<Option<Passport>, DppError> {
self.find_by_id(id).await
}
async fn update(&self, passport: Passport) -> Result<Passport, DppError> {
self.store
.lock()
.unwrap()
.insert(passport.id, passport.clone());
Ok(passport)
}
async fn update_status(
&self,
id: PassportId,
status: PassportStatus,
) -> Result<Passport, DppError> {
let mut g = self.store.lock().unwrap();
let mut p = g
.get(&id)
.cloned()
.ok_or(DppError::NotFound(id.to_string()))?;
p.status = status;
g.insert(id, p.clone());
Ok(p)
}
async fn list(
&self,
_status: Option<PassportStatus>,
_q: Option<&str>,
_facility_id: Option<&str>,
_limit: u32,
_offset: u32,
) -> Result<Vec<Passport>, DppError> {
Ok(self.store.lock().unwrap().values().cloned().collect())
}
async fn count(
&self,
_status: Option<PassportStatus>,
_facility_id: Option<&str>,
) -> Result<u64, DppError> {
Ok(self.store.lock().unwrap().len() as u64)
}
}
fn draft_passport(name: &str) -> Passport {
Passport {
id: PassportId::new(),
batch_id: None,
product_name: name.into(),
sector: Sector::Textile,
product_category: None,
manufacturer: ManufacturerInfo {
name: "Brand".into(),
address: "Berlin, DE".into(),
did_web_url: None,
},
materials: vec![],
co2e_per_unit: None,
repairability_score: None,
compliance_result: None,
lint_result: None,
sector_data: None,
status: PassportStatus::Draft,
qr_code_url: None,
jws_signature: None,
public_jws_signature: None,
created_at: chrono::Utc::now(),
updated_at: chrono::Utc::now(),
published_at: None,
schema_version: "1.1.0".into(),
retention_locked: false,
version: 1,
supersedes_id: None,
parent_passport_ref: None,
component_refs: Vec::new(),
retention_until: None,
product_id: None,
operator_identifier: None,
facility: None,
seal: None,
}
}
#[tokio::test]
async fn default_patch_fields_merges_delta() {
let repo = InMemoryRepo::default();
let p = repo.create(draft_passport("Original")).await.unwrap();
let patched = repo
.patch_fields(p.id, serde_json::json!({ "productName": "Renamed" }))
.await
.unwrap();
assert_eq!(patched.product_name, "Renamed");
assert_eq!(patched.id, p.id);
}
#[tokio::test]
async fn default_patch_fields_unknown_id_is_not_found() {
let repo = InMemoryRepo::default();
let err = repo
.patch_fields(PassportId::new(), serde_json::json!({}))
.await
.unwrap_err();
assert!(matches!(err, DppError::NotFound(_)));
}
#[tokio::test]
async fn default_find_by_identity_matches_across_draft_and_published() {
use crate::domain::gtin::Gtin;
use crate::domain::sector::{BatteryChemistry, BatteryData, SectorData};
let repo = InMemoryRepo::default();
let mut p = draft_passport("Battery A");
p.sector = Sector::Battery;
p.sector_data = Some(SectorData::Battery(BatteryData {
gtin: Gtin::parse("09506000134352").unwrap(),
battery_chemistry: BatteryChemistry::Lfp,
nominal_voltage_v: 3.2,
nominal_capacity_ah: 100.0,
expected_lifetime_cycles: 3000,
co2e_per_unit_kg: 85.4,
recycled_content_cobalt_pct: None,
recycled_content_lithium_pct: None,
recycled_content_nickel_pct: None,
state_of_health_pct: None,
rated_capacity_kwh: None,
carbon_footprint_class: None,
due_diligence_url: None,
cathode_material: None,
anode_material: None,
electrolyte_material: None,
critical_raw_materials: None,
disassembly_instructions_url: None,
soh_methodology: None,
operating_temp_min_c: None,
operating_temp_max_c: None,
rated_energy_wh: None,
recycled_content_lead_pct: None,
battery_weight_kg: None,
battery_type: None,
round_trip_efficiency_pct: None,
internal_resistance_mohm: None,
manufacturing_date: None,
manufacturing_place: None,
battery_model_id: None,
battery_passport_number: None,
}));
p.batch_id = Some("BATCH-1".into());
let created = repo.create(p).await.unwrap();
let identity = ProductIdentity {
sector: Sector::Battery,
gtin: "09506000134352".into(),
batch_id: Some("BATCH-1".into()),
};
let found = repo.find_by_identity(&identity).await.unwrap();
assert_eq!(found.map(|p| p.id), Some(created.id));
let no_match = ProductIdentity {
sector: Sector::Battery,
gtin: "00000000000000".into(),
batch_id: None,
};
assert!(repo.find_by_identity(&no_match).await.unwrap().is_none());
}
#[tokio::test]
async fn default_create_and_update_batch_run_sequentially() {
let repo = InMemoryRepo::default();
let created = repo
.create_batch(vec![draft_passport("A"), draft_passport("B")])
.await;
assert_eq!(created.len(), 2);
assert!(created.iter().all(|r| r.is_ok()));
let mut a = created[0].as_ref().unwrap().clone();
a.product_name = "A2".into();
let updated = repo.update_batch(vec![a]).await;
assert_eq!(updated.len(), 1);
assert_eq!(updated[0].as_ref().unwrap().product_name, "A2");
}
}