use std::collections::{BTreeSet, HashSet};
use std::ops::ControlFlow;
use auths_keri::{
Acdc, DelegatorKelLookup, KelSealIndex, Prefix, Said, SignedEvent, SourceSeal, TelEvent,
parse_delegated_attachment, validate_signed_kel, validate_tel,
};
use super::kel_resolver::{KelResolveError, collect_kel_capped, verify_prefix_binding};
use crate::ports::registry::{RegistryBackend, RegistryError};
#[derive(Debug, Clone, Copy)]
pub struct KelCaps {
pub max_events: usize,
pub max_bytes: usize,
}
#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize)]
#[serde(rename_all = "snake_case", tag = "outcome")]
pub enum MergeOutcome {
Imported {
events: usize,
},
Advanced {
events: usize,
},
AlreadyCurrent,
}
#[derive(Debug, Clone, serde::Serialize)]
pub struct MergedKel {
pub prefix: Prefix,
#[serde(flatten)]
pub outcome: MergeOutcome,
}
#[derive(Debug, thiserror::Error)]
#[non_exhaustive]
pub enum RegistryMergeError {
#[error("source registry lists an invalid identity prefix '{id}': {reason}")]
InvalidPrefix {
id: String,
reason: String,
},
#[error("source KEL for {prefix} was refused: {source}")]
SourceKel {
prefix: Prefix,
#[source]
source: KelResolveError,
},
#[error("source KEL for {prefix} has no signature attachment at sequence {sequence}")]
MissingSignature {
prefix: Prefix,
sequence: u128,
},
#[error("source KEL for {prefix} failed authentication: {reason}")]
Unauthenticated {
prefix: Prefix,
reason: String,
},
#[error(
"KEL fork for {prefix} at sequence {sequence}: \
destination has {destination}, source has {incoming}"
)]
Forked {
prefix: Prefix,
sequence: u128,
destination: Said,
incoming: Said,
},
#[error("registry backend error: {0}")]
Storage(#[from] RegistryError),
#[error("source credential {credential} under {issuer} was refused: {reason}")]
CredentialRefused {
issuer: Prefix,
credential: Said,
reason: String,
},
#[error("source credential {credential} names unknown issuer {issuer} (no authenticated KEL)")]
CredentialOrphan {
issuer: Prefix,
credential: Said,
},
#[error("source TEL {credential} under {issuer}/{registry} was refused: {reason}")]
TelRefused {
issuer: Prefix,
registry: Said,
credential: Said,
reason: String,
},
}
pub fn merge_registries(
source: &dyn RegistryBackend,
dest: &dyn RegistryBackend,
caps: &KelCaps,
) -> Result<Vec<MergedKel>, RegistryMergeError> {
let mut ids: Vec<String> = Vec::new();
source.visit_identities(&mut |id| {
ids.push(id.to_string());
ControlFlow::Continue(())
})?;
ids.sort();
ids.dedup();
let mut report = Vec::with_capacity(ids.len());
for id in ids {
let prefix = Prefix::new(id.clone()).map_err(|e| RegistryMergeError::InvalidPrefix {
id,
reason: e.to_string(),
})?;
let outcome = merge_kel(source, dest, &prefix, caps)?;
report.push(MergedKel { prefix, outcome });
}
Ok(report)
}
fn merge_kel(
source: &dyn RegistryBackend,
dest: &dyn RegistryBackend,
prefix: &Prefix,
caps: &KelCaps,
) -> Result<MergeOutcome, RegistryMergeError> {
let refused = |source: KelResolveError| RegistryMergeError::SourceKel {
prefix: prefix.clone(),
source,
};
let events =
collect_kel_capped(source, prefix, caps.max_events, caps.max_bytes).map_err(refused)?;
verify_prefix_binding(prefix, &events).map_err(refused)?;
let mut attachments = Vec::with_capacity(events.len());
let mut signed = Vec::with_capacity(events.len());
for event in &events {
let sequence = event.sequence().value();
let attachment = source.get_attachment(prefix, sequence)?.ok_or_else(|| {
RegistryMergeError::MissingSignature {
prefix: prefix.clone(),
sequence,
}
})?;
let (sigs, _seals) = parse_delegated_attachment(&attachment).map_err(|e| {
RegistryMergeError::Unauthenticated {
prefix: prefix.clone(),
reason: format!("unparseable attachment at sequence {sequence}: {e}"),
}
})?;
signed.push(SignedEvent::new(event.clone(), sigs));
attachments.push(attachment);
}
let lookup = BackendSealLookup {
backends: [source, dest],
caps: *caps,
};
validate_signed_kel(&signed, Some(&lookup)).map_err(|e| {
RegistryMergeError::Unauthenticated {
prefix: prefix.clone(),
reason: e.to_string(),
}
})?;
let dest_tip = match dest.get_tip(prefix) {
Ok(tip) => Some(tip.sequence),
Err(RegistryError::NotFound { .. }) => None,
Err(e) => return Err(e.into()),
};
match dest_tip {
None => {
for (event, attachment) in events.iter().zip(&attachments) {
dest.append_signed_event(prefix, event, attachment)?;
}
Ok(MergeOutcome::Imported {
events: events.len(),
})
}
Some(dest_tip) => {
for event in events.iter().filter(|e| e.sequence().value() <= dest_tip) {
let sequence = event.sequence().value();
let local = dest.get_event(prefix, sequence)?;
if local.said() != event.said() {
return Err(RegistryMergeError::Forked {
prefix: prefix.clone(),
sequence,
destination: local.said().clone(),
incoming: event.said().clone(),
});
}
}
let newer: Vec<_> = events
.iter()
.zip(&attachments)
.filter(|(event, _)| event.sequence().value() > dest_tip)
.collect();
if newer.is_empty() {
return Ok(MergeOutcome::AlreadyCurrent);
}
let appended = newer.len();
for (event, attachment) in newer {
dest.append_signed_event(prefix, event, attachment)?;
}
Ok(MergeOutcome::Advanced { events: appended })
}
}
}
#[derive(Debug, Clone, Default, PartialEq, Eq, serde::Serialize)]
pub struct MergedCredentials {
pub credentials_imported: usize,
pub credentials_already_present: usize,
pub tel_events_imported: usize,
pub tel_events_already_present: usize,
}
pub fn merge_credentials_and_tel(
source: &dyn RegistryBackend,
dest: &dyn RegistryBackend,
authenticated_issuers: &HashSet<Prefix>,
) -> Result<MergedCredentials, RegistryMergeError> {
let mut report = MergedCredentials::default();
merge_credential_bodies(source, dest, authenticated_issuers, &mut report)?;
merge_tel_chains(source, dest, authenticated_issuers, &mut report)?;
Ok(report)
}
fn merge_credential_bodies(
source: &dyn RegistryBackend,
dest: &dyn RegistryBackend,
authenticated_issuers: &HashSet<Prefix>,
report: &mut MergedCredentials,
) -> Result<(), RegistryMergeError> {
let mut credentials: Vec<(Prefix, Said, Vec<u8>)> = Vec::new();
source.visit_credentials(&mut |issuer, credential, bytes| {
credentials.push((issuer.clone(), credential.clone(), bytes.to_vec()));
ControlFlow::Continue(())
})?;
for (issuer, credential, bytes) in credentials {
if !authenticated_issuers.contains(&issuer) {
return Err(RegistryMergeError::CredentialOrphan { issuer, credential });
}
verify_credential_body(&issuer, &credential, &bytes)?;
if dest.load_credential(&issuer, &credential)?.is_some() {
report.credentials_already_present += 1;
continue;
}
dest.store_credential(&issuer, &credential, &bytes)?;
report.credentials_imported += 1;
}
Ok(())
}
fn merge_tel_chains(
source: &dyn RegistryBackend,
dest: &dyn RegistryBackend,
authenticated_issuers: &HashSet<Prefix>,
report: &mut MergedCredentials,
) -> Result<(), RegistryMergeError> {
let mut coordinates: Vec<(Prefix, Said, Said)> = Vec::new();
source.visit_tel_registries(&mut |issuer, registry, credential| {
coordinates.push((issuer.clone(), registry.clone(), credential.clone()));
ControlFlow::Continue(())
})?;
for (issuer, registry, credential) in coordinates {
if !authenticated_issuers.contains(&issuer) {
return Err(RegistryMergeError::CredentialOrphan { issuer, credential });
}
merge_one_tel(source, dest, &issuer, ®istry, &credential, report)?;
}
Ok(())
}
fn merge_one_tel(
source: &dyn RegistryBackend,
dest: &dyn RegistryBackend,
issuer: &Prefix,
registry: &Said,
credential: &Said,
report: &mut MergedCredentials,
) -> Result<(), RegistryMergeError> {
let refused = |reason: String| RegistryMergeError::TelRefused {
issuer: issuer.clone(),
registry: registry.clone(),
credential: credential.clone(),
reason,
};
let mut raw: Vec<(u128, Vec<u8>)> = Vec::new();
read_tel_raw(source, issuer, registry, credential, &mut raw).map_err(&refused)?;
if raw.is_empty() {
return Ok(());
}
let mut chain: Vec<TelEvent> = Vec::new();
if credential != registry {
collect_tel(source, issuer, registry, registry, &mut chain).map_err(&refused)?;
}
collect_tel(source, issuer, registry, credential, &mut chain).map_err(&refused)?;
validate_tel(&chain).map_err(|e| refused(e.to_string()))?;
let mut present: BTreeSet<u128> = BTreeSet::new();
read_tel_sns(dest, issuer, registry, credential, &mut present)
.map_err(|reason| refused(format!("local TEL event did not parse: {reason}")))?;
for (sn, bytes) in raw {
if present.contains(&sn) {
report.tel_events_already_present += 1;
continue;
}
dest.append_tel_event(issuer, registry, credential, sn, &bytes)?;
report.tel_events_imported += 1;
}
Ok(())
}
fn read_tel_raw(
source: &dyn RegistryBackend,
issuer: &Prefix,
registry: &Said,
credential: &Said,
raw: &mut Vec<(u128, Vec<u8>)>,
) -> Result<(), String> {
let mut parse_err: Option<String> = None;
source
.visit_tel_events(
issuer,
registry,
credential,
&mut |bytes| match TelEvent::from_wire_bytes(bytes) {
Ok(event) => {
raw.push((tel_event_sn(&event), bytes.to_vec()));
ControlFlow::Continue(())
}
Err(e) => {
parse_err = Some(format!("TEL event did not parse: {e}"));
ControlFlow::Break(())
}
},
)
.map_err(|e| e.to_string())?;
match parse_err {
Some(reason) => Err(reason),
None => Ok(()),
}
}
fn read_tel_sns(
dest: &dyn RegistryBackend,
issuer: &Prefix,
registry: &Said,
credential: &Said,
present: &mut BTreeSet<u128>,
) -> Result<(), String> {
let mut parse_err: Option<String> = None;
dest.visit_tel_events(
issuer,
registry,
credential,
&mut |bytes| match TelEvent::from_wire_bytes(bytes) {
Ok(event) => {
present.insert(tel_event_sn(&event));
ControlFlow::Continue(())
}
Err(e) => {
parse_err = Some(e.to_string());
ControlFlow::Break(())
}
},
)
.map_err(|e| e.to_string())?;
match parse_err {
Some(reason) => Err(reason),
None => Ok(()),
}
}
fn tel_event_sn(event: &TelEvent) -> u128 {
match event {
TelEvent::Vcp(vcp) => vcp.s.value(),
TelEvent::Iss(iss) => iss.s.value(),
TelEvent::Rev(rev) => rev.s.value(),
}
}
fn verify_credential_body(
issuer: &Prefix,
credential: &Said,
bytes: &[u8],
) -> Result<(), RegistryMergeError> {
let refused = |reason: String| RegistryMergeError::CredentialRefused {
issuer: issuer.clone(),
credential: credential.clone(),
reason,
};
let value: serde_json::Value =
serde_json::from_slice(bytes).map_err(|e| refused(format!("blob is not JSON: {e}")))?;
let acdc_value = value
.get("acdc")
.ok_or_else(|| refused("blob has no `acdc` body".to_string()))?;
let acdc: Acdc = serde_json::from_value(acdc_value.clone())
.map_err(|e| refused(format!("acdc body did not parse: {e}")))?;
acdc.verify_said()
.map_err(|e| refused(format!("acdc SAID does not recompute: {e}")))?;
if acdc.d.as_str() != credential.as_str() {
return Err(refused(format!(
"acdc SAID {} does not match the path SAID {}",
acdc.d.as_str(),
credential.as_str()
)));
}
if acdc.i.as_str() != issuer.as_str() {
return Err(refused(format!(
"acdc issuer {} does not match the path issuer {}",
acdc.i.as_str(),
issuer.as_str()
)));
}
Ok(())
}
fn collect_tel(
source: &dyn RegistryBackend,
issuer: &Prefix,
registry: &Said,
credential: &Said,
chain: &mut Vec<TelEvent>,
) -> Result<(), String> {
let mut parse_err: Option<String> = None;
source
.visit_tel_events(
issuer,
registry,
credential,
&mut |bytes| match TelEvent::from_wire_bytes(bytes) {
Ok(event) => {
chain.push(event);
ControlFlow::Continue(())
}
Err(e) => {
parse_err = Some(format!("TEL event did not parse: {e}"));
ControlFlow::Break(())
}
},
)
.map_err(|e| e.to_string())?;
if let Some(reason) = parse_err {
return Err(reason);
}
Ok(())
}
struct BackendSealLookup<'a> {
backends: [&'a dyn RegistryBackend; 2],
caps: KelCaps,
}
impl DelegatorKelLookup for BackendSealLookup<'_> {
fn find_seal(&self, delegator_aid: &Prefix, seal_said: &Said) -> Option<SourceSeal> {
for backend in self.backends {
let Ok(events) = collect_kel_capped(
backend,
delegator_aid,
self.caps.max_events,
self.caps.max_bytes,
) else {
continue;
};
if let Some(seal) =
KelSealIndex::from_events(&events).find_seal(delegator_aid, seal_said)
{
return Some(seal);
}
}
None
}
}