1use std::collections::{BTreeSet, HashSet};
28use std::ops::ControlFlow;
29
30use auths_keri::{
31 Acdc, DelegatorKelLookup, KelSealIndex, Prefix, Said, SignedEvent, SourceSeal, TelEvent,
32 parse_delegated_attachment, validate_signed_kel, validate_tel,
33};
34
35use super::kel_resolver::{KelResolveError, collect_kel_capped, verify_prefix_binding};
36use crate::ports::registry::{RegistryBackend, RegistryError};
37
38#[derive(Debug, Clone, Copy)]
40pub struct KelCaps {
41 pub max_events: usize,
43 pub max_bytes: usize,
45}
46
47#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize)]
49#[serde(rename_all = "snake_case", tag = "outcome")]
50pub enum MergeOutcome {
51 Imported {
54 events: usize,
56 },
57 Advanced {
60 events: usize,
62 },
63 AlreadyCurrent,
66}
67
68#[derive(Debug, Clone, serde::Serialize)]
70pub struct MergedKel {
71 pub prefix: Prefix,
73 #[serde(flatten)]
75 pub outcome: MergeOutcome,
76}
77
78#[derive(Debug, thiserror::Error)]
81#[non_exhaustive]
82pub enum RegistryMergeError {
83 #[error("source registry lists an invalid identity prefix '{id}': {reason}")]
85 InvalidPrefix {
86 id: String,
88 reason: String,
90 },
91
92 #[error("source KEL for {prefix} was refused: {source}")]
95 SourceKel {
96 prefix: Prefix,
98 #[source]
100 source: KelResolveError,
101 },
102
103 #[error("source KEL for {prefix} has no signature attachment at sequence {sequence}")]
106 MissingSignature {
107 prefix: Prefix,
109 sequence: u128,
111 },
112
113 #[error("source KEL for {prefix} failed authentication: {reason}")]
116 Unauthenticated {
117 prefix: Prefix,
119 reason: String,
121 },
122
123 #[error(
126 "KEL fork for {prefix} at sequence {sequence}: \
127 destination has {destination}, source has {incoming}"
128 )]
129 Forked {
130 prefix: Prefix,
132 sequence: u128,
134 destination: Said,
136 incoming: Said,
138 },
139
140 #[error("registry backend error: {0}")]
142 Storage(#[from] RegistryError),
143
144 #[error("source credential {credential} under {issuer} was refused: {reason}")]
149 CredentialRefused {
150 issuer: Prefix,
152 credential: Said,
154 reason: String,
156 },
157
158 #[error("source credential {credential} names unknown issuer {issuer} (no authenticated KEL)")]
161 CredentialOrphan {
162 issuer: Prefix,
164 credential: Said,
166 },
167
168 #[error("source TEL {credential} under {issuer}/{registry} was refused: {reason}")]
172 TelRefused {
173 issuer: Prefix,
175 registry: Said,
177 credential: Said,
179 reason: String,
181 },
182}
183
184pub fn merge_registries(
201 source: &dyn RegistryBackend,
202 dest: &dyn RegistryBackend,
203 caps: &KelCaps,
204) -> Result<Vec<MergedKel>, RegistryMergeError> {
205 let mut ids: Vec<String> = Vec::new();
206 source.visit_identities(&mut |id| {
207 ids.push(id.to_string());
208 ControlFlow::Continue(())
209 })?;
210 ids.sort();
211 ids.dedup();
212
213 let mut report = Vec::with_capacity(ids.len());
214 for id in ids {
215 let prefix = Prefix::new(id.clone()).map_err(|e| RegistryMergeError::InvalidPrefix {
216 id,
217 reason: e.to_string(),
218 })?;
219 let outcome = merge_kel(source, dest, &prefix, caps)?;
220 report.push(MergedKel { prefix, outcome });
221 }
222 Ok(report)
223}
224
225fn merge_kel(
227 source: &dyn RegistryBackend,
228 dest: &dyn RegistryBackend,
229 prefix: &Prefix,
230 caps: &KelCaps,
231) -> Result<MergeOutcome, RegistryMergeError> {
232 let refused = |source: KelResolveError| RegistryMergeError::SourceKel {
233 prefix: prefix.clone(),
234 source,
235 };
236
237 let events =
238 collect_kel_capped(source, prefix, caps.max_events, caps.max_bytes).map_err(refused)?;
239 verify_prefix_binding(prefix, &events).map_err(refused)?;
240
241 let mut attachments = Vec::with_capacity(events.len());
244 let mut signed = Vec::with_capacity(events.len());
245 for event in &events {
246 let sequence = event.sequence().value();
247 let attachment = source.get_attachment(prefix, sequence)?.ok_or_else(|| {
248 RegistryMergeError::MissingSignature {
249 prefix: prefix.clone(),
250 sequence,
251 }
252 })?;
253 let (sigs, _seals) = parse_delegated_attachment(&attachment).map_err(|e| {
254 RegistryMergeError::Unauthenticated {
255 prefix: prefix.clone(),
256 reason: format!("unparseable attachment at sequence {sequence}: {e}"),
257 }
258 })?;
259 signed.push(SignedEvent::new(event.clone(), sigs));
262 attachments.push(attachment);
263 }
264
265 let lookup = BackendSealLookup {
266 backends: [source, dest],
267 caps: *caps,
268 };
269 validate_signed_kel(&signed, Some(&lookup)).map_err(|e| {
270 RegistryMergeError::Unauthenticated {
271 prefix: prefix.clone(),
272 reason: e.to_string(),
273 }
274 })?;
275
276 let dest_tip = match dest.get_tip(prefix) {
277 Ok(tip) => Some(tip.sequence),
278 Err(RegistryError::NotFound { .. }) => None,
279 Err(e) => return Err(e.into()),
280 };
281
282 match dest_tip {
283 None => {
284 for (event, attachment) in events.iter().zip(&attachments) {
285 dest.append_signed_event(prefix, event, attachment)?;
286 }
287 Ok(MergeOutcome::Imported {
288 events: events.len(),
289 })
290 }
291 Some(dest_tip) => {
292 for event in events.iter().filter(|e| e.sequence().value() <= dest_tip) {
294 let sequence = event.sequence().value();
295 let local = dest.get_event(prefix, sequence)?;
296 if local.said() != event.said() {
297 return Err(RegistryMergeError::Forked {
298 prefix: prefix.clone(),
299 sequence,
300 destination: local.said().clone(),
301 incoming: event.said().clone(),
302 });
303 }
304 }
305 let newer: Vec<_> = events
306 .iter()
307 .zip(&attachments)
308 .filter(|(event, _)| event.sequence().value() > dest_tip)
309 .collect();
310 if newer.is_empty() {
311 return Ok(MergeOutcome::AlreadyCurrent);
312 }
313 let appended = newer.len();
314 for (event, attachment) in newer {
315 dest.append_signed_event(prefix, event, attachment)?;
316 }
317 Ok(MergeOutcome::Advanced { events: appended })
318 }
319 }
320}
321
322#[derive(Debug, Clone, Default, PartialEq, Eq, serde::Serialize)]
329pub struct MergedCredentials {
330 pub credentials_imported: usize,
332 pub credentials_already_present: usize,
334 pub tel_events_imported: usize,
336 pub tel_events_already_present: usize,
338}
339
340pub fn merge_credentials_and_tel(
379 source: &dyn RegistryBackend,
380 dest: &dyn RegistryBackend,
381 authenticated_issuers: &HashSet<Prefix>,
382) -> Result<MergedCredentials, RegistryMergeError> {
383 let mut report = MergedCredentials::default();
384 merge_credential_bodies(source, dest, authenticated_issuers, &mut report)?;
385 merge_tel_chains(source, dest, authenticated_issuers, &mut report)?;
386 Ok(report)
387}
388
389fn merge_credential_bodies(
392 source: &dyn RegistryBackend,
393 dest: &dyn RegistryBackend,
394 authenticated_issuers: &HashSet<Prefix>,
395 report: &mut MergedCredentials,
396) -> Result<(), RegistryMergeError> {
397 let mut credentials: Vec<(Prefix, Said, Vec<u8>)> = Vec::new();
398 source.visit_credentials(&mut |issuer, credential, bytes| {
399 credentials.push((issuer.clone(), credential.clone(), bytes.to_vec()));
400 ControlFlow::Continue(())
401 })?;
402
403 for (issuer, credential, bytes) in credentials {
404 if !authenticated_issuers.contains(&issuer) {
405 return Err(RegistryMergeError::CredentialOrphan { issuer, credential });
406 }
407 verify_credential_body(&issuer, &credential, &bytes)?;
408
409 if dest.load_credential(&issuer, &credential)?.is_some() {
410 report.credentials_already_present += 1;
411 continue;
412 }
413 dest.store_credential(&issuer, &credential, &bytes)?;
414 report.credentials_imported += 1;
415 }
416 Ok(())
417}
418
419fn merge_tel_chains(
422 source: &dyn RegistryBackend,
423 dest: &dyn RegistryBackend,
424 authenticated_issuers: &HashSet<Prefix>,
425 report: &mut MergedCredentials,
426) -> Result<(), RegistryMergeError> {
427 let mut coordinates: Vec<(Prefix, Said, Said)> = Vec::new();
428 source.visit_tel_registries(&mut |issuer, registry, credential| {
429 coordinates.push((issuer.clone(), registry.clone(), credential.clone()));
430 ControlFlow::Continue(())
431 })?;
432
433 for (issuer, registry, credential) in coordinates {
434 if !authenticated_issuers.contains(&issuer) {
435 return Err(RegistryMergeError::CredentialOrphan { issuer, credential });
436 }
437 merge_one_tel(source, dest, &issuer, ®istry, &credential, report)?;
438 }
439 Ok(())
440}
441
442fn merge_one_tel(
444 source: &dyn RegistryBackend,
445 dest: &dyn RegistryBackend,
446 issuer: &Prefix,
447 registry: &Said,
448 credential: &Said,
449 report: &mut MergedCredentials,
450) -> Result<(), RegistryMergeError> {
451 let refused = |reason: String| RegistryMergeError::TelRefused {
452 issuer: issuer.clone(),
453 registry: registry.clone(),
454 credential: credential.clone(),
455 reason,
456 };
457
458 let mut raw: Vec<(u128, Vec<u8>)> = Vec::new();
461 read_tel_raw(source, issuer, registry, credential, &mut raw).map_err(&refused)?;
462 if raw.is_empty() {
463 return Ok(());
464 }
465
466 let mut chain: Vec<TelEvent> = Vec::new();
472 if credential != registry {
473 collect_tel(source, issuer, registry, registry, &mut chain).map_err(&refused)?;
474 }
475 collect_tel(source, issuer, registry, credential, &mut chain).map_err(&refused)?;
476 validate_tel(&chain).map_err(|e| refused(e.to_string()))?;
477
478 let mut present: BTreeSet<u128> = BTreeSet::new();
480 read_tel_sns(dest, issuer, registry, credential, &mut present)
481 .map_err(|reason| refused(format!("local TEL event did not parse: {reason}")))?;
482 for (sn, bytes) in raw {
483 if present.contains(&sn) {
484 report.tel_events_already_present += 1;
485 continue;
486 }
487 dest.append_tel_event(issuer, registry, credential, sn, &bytes)?;
488 report.tel_events_imported += 1;
489 }
490 Ok(())
491}
492
493fn read_tel_raw(
496 source: &dyn RegistryBackend,
497 issuer: &Prefix,
498 registry: &Said,
499 credential: &Said,
500 raw: &mut Vec<(u128, Vec<u8>)>,
501) -> Result<(), String> {
502 let mut parse_err: Option<String> = None;
503 source
504 .visit_tel_events(
505 issuer,
506 registry,
507 credential,
508 &mut |bytes| match TelEvent::from_wire_bytes(bytes) {
509 Ok(event) => {
510 raw.push((tel_event_sn(&event), bytes.to_vec()));
511 ControlFlow::Continue(())
512 }
513 Err(e) => {
514 parse_err = Some(format!("TEL event did not parse: {e}"));
515 ControlFlow::Break(())
516 }
517 },
518 )
519 .map_err(|e| e.to_string())?;
520 match parse_err {
521 Some(reason) => Err(reason),
522 None => Ok(()),
523 }
524}
525
526fn read_tel_sns(
528 dest: &dyn RegistryBackend,
529 issuer: &Prefix,
530 registry: &Said,
531 credential: &Said,
532 present: &mut BTreeSet<u128>,
533) -> Result<(), String> {
534 let mut parse_err: Option<String> = None;
535 dest.visit_tel_events(
536 issuer,
537 registry,
538 credential,
539 &mut |bytes| match TelEvent::from_wire_bytes(bytes) {
540 Ok(event) => {
541 present.insert(tel_event_sn(&event));
542 ControlFlow::Continue(())
543 }
544 Err(e) => {
545 parse_err = Some(e.to_string());
546 ControlFlow::Break(())
547 }
548 },
549 )
550 .map_err(|e| e.to_string())?;
551 match parse_err {
552 Some(reason) => Err(reason),
553 None => Ok(()),
554 }
555}
556
557fn tel_event_sn(event: &TelEvent) -> u128 {
559 match event {
560 TelEvent::Vcp(vcp) => vcp.s.value(),
561 TelEvent::Iss(iss) => iss.s.value(),
562 TelEvent::Rev(rev) => rev.s.value(),
563 }
564}
565
566fn verify_credential_body(
569 issuer: &Prefix,
570 credential: &Said,
571 bytes: &[u8],
572) -> Result<(), RegistryMergeError> {
573 let refused = |reason: String| RegistryMergeError::CredentialRefused {
574 issuer: issuer.clone(),
575 credential: credential.clone(),
576 reason,
577 };
578
579 let value: serde_json::Value =
583 serde_json::from_slice(bytes).map_err(|e| refused(format!("blob is not JSON: {e}")))?;
584 let acdc_value = value
585 .get("acdc")
586 .ok_or_else(|| refused("blob has no `acdc` body".to_string()))?;
587 let acdc: Acdc = serde_json::from_value(acdc_value.clone())
588 .map_err(|e| refused(format!("acdc body did not parse: {e}")))?;
589
590 acdc.verify_said()
591 .map_err(|e| refused(format!("acdc SAID does not recompute: {e}")))?;
592 if acdc.d.as_str() != credential.as_str() {
593 return Err(refused(format!(
594 "acdc SAID {} does not match the path SAID {}",
595 acdc.d.as_str(),
596 credential.as_str()
597 )));
598 }
599 if acdc.i.as_str() != issuer.as_str() {
600 return Err(refused(format!(
601 "acdc issuer {} does not match the path issuer {}",
602 acdc.i.as_str(),
603 issuer.as_str()
604 )));
605 }
606 Ok(())
607}
608
609fn collect_tel(
612 source: &dyn RegistryBackend,
613 issuer: &Prefix,
614 registry: &Said,
615 credential: &Said,
616 chain: &mut Vec<TelEvent>,
617) -> Result<(), String> {
618 let mut parse_err: Option<String> = None;
619 source
620 .visit_tel_events(
621 issuer,
622 registry,
623 credential,
624 &mut |bytes| match TelEvent::from_wire_bytes(bytes) {
625 Ok(event) => {
626 chain.push(event);
627 ControlFlow::Continue(())
628 }
629 Err(e) => {
630 parse_err = Some(format!("TEL event did not parse: {e}"));
631 ControlFlow::Break(())
632 }
633 },
634 )
635 .map_err(|e| e.to_string())?;
636 if let Some(reason) = parse_err {
637 return Err(reason);
638 }
639 Ok(())
640}
641
642struct BackendSealLookup<'a> {
646 backends: [&'a dyn RegistryBackend; 2],
647 caps: KelCaps,
648}
649
650impl DelegatorKelLookup for BackendSealLookup<'_> {
651 fn find_seal(&self, delegator_aid: &Prefix, seal_said: &Said) -> Option<SourceSeal> {
652 for backend in self.backends {
653 let Ok(events) = collect_kel_capped(
654 backend,
655 delegator_aid,
656 self.caps.max_events,
657 self.caps.max_bytes,
658 ) else {
659 continue;
660 };
661 if let Some(seal) =
662 KelSealIndex::from_events(&events).find_seal(delegator_aid, seal_said)
663 {
664 return Some(seal);
665 }
666 }
667 None
668 }
669}