1use std::fmt;
5use std::io::Read;
6use std::sync::Arc;
7
8#[cfg(feature = "native-grpc")]
9use iota_grpc_client::Client as GrpcClient;
10use iota_sdk_types::CheckpointContents;
11use iota_types::committee::{Committee, EpochId};
12use iota_types::digests::ChainIdentifier;
13use iota_types::effects::{TransactionEffects, TransactionEvents};
14use iota_types::error::IotaError;
15use iota_types::iota_system_state::{IotaSystemStateTrait, get_iota_system_state};
16use iota_types::messages_checkpoint::CertifiedCheckpointSummary;
17use iota_types::object::Object;
18use iota_types::transaction::Transaction;
19use serde::Deserialize;
20
21use crate::{
22 BoxError, CommitteeCache, CommitteeCacheError, CommitteeCacheKey, MemoryCommitteeCache, Proof, ProofVerifier,
23 Source, VerifiedProof, VerifyError,
24};
25
26#[derive(Debug, thiserror::Error)]
28#[non_exhaustive]
29#[error("failed to resolve committee for epoch {target_epoch}")]
30pub struct CommitteeResolutionError {
31 pub target_epoch: EpochId,
33 #[source]
35 pub kind: CommitteeResolutionErrorKind,
36}
37
38impl CommitteeResolutionError {
39 fn new(target_epoch: EpochId, kind: CommitteeResolutionErrorKind) -> Self {
41 Self { target_epoch, kind }
42 }
43}
44
45#[derive(Debug, thiserror::Error)]
47#[non_exhaustive]
48pub enum CommitteeResolutionErrorKind {
49 #[error("failed to load the committee from the trusted genesis blob")]
51 LoadGenesisCommittee {
52 #[source]
54 source: BoxError,
55 },
56 #[error("trusted genesis checkpoint has epoch {epoch}, expected epoch 0")]
58 UnexpectedGenesisCheckpointEpoch {
59 epoch: EpochId,
61 },
62 #[error("trusted genesis committee has epoch {epoch}, expected epoch 0")]
64 UnexpectedGenesisCommitteeEpoch {
65 epoch: EpochId,
67 },
68 #[error("trusted genesis checkpoint failed verification")]
70 InvalidGenesisCheckpoint {
71 #[source]
73 source: BoxError,
74 },
75 #[error("failed to fetch committee for epoch {epoch} from the trusted node")]
77 FetchCommittee {
78 epoch: EpochId,
80 #[source]
82 source: BoxError,
83 },
84 #[error("target epoch is before trusted anchor epoch {anchor_epoch}")]
86 TargetBeforeAnchor {
87 anchor_epoch: EpochId,
89 },
90 #[error("failed to fetch the node's current epoch")]
92 FetchCurrentEpoch {
93 #[source]
95 source: BoxError,
96 },
97 #[error("service information is missing the current epoch")]
99 MissingCurrentEpoch,
100 #[error("target epoch is ahead of node current epoch {current_epoch}")]
102 TargetAheadOfNode {
103 current_epoch: EpochId,
105 },
106 #[error("failed to fetch end-of-epoch checkpoint information for epoch {epoch}")]
108 FetchEpochHistory {
109 epoch: EpochId,
111 #[source]
113 source: BoxError,
114 },
115 #[error("epoch {epoch} is missing its epoch-close proof")]
117 MissingEpochCloseProof {
118 epoch: EpochId,
120 },
121 #[error("failed to verify epoch {epoch} end-of-epoch checkpoint {sequence_number}")]
123 InvalidEndOfEpochCheckpoint {
124 epoch: EpochId,
126 sequence_number: u64,
128 #[source]
130 source: BoxError,
131 },
132 #[error("checkpoint {sequence_number} is not an end-of-epoch checkpoint")]
134 NotEndOfEpoch {
135 sequence_number: u64,
137 },
138 #[error("next epoch after {epoch} overflows u64")]
140 NextEpochOverflow {
141 epoch: EpochId,
143 },
144 #[error("committee cache failed at epoch {epoch}")]
146 Cache {
147 epoch: EpochId,
149 #[source]
151 source: CommitteeCacheError,
152 },
153}
154
155#[derive(Debug, thiserror::Error)]
157#[non_exhaustive]
158pub enum ProofVerificationError {
159 #[error("failed to resolve the committee required by the proof")]
161 CommitteeResolution {
162 #[source]
164 source: CommitteeResolutionError,
165 },
166 #[error("proof verification failed")]
168 Proof {
169 #[source]
171 source: VerifyError,
172 },
173}
174
175#[derive(Clone)]
177#[non_exhaustive]
178pub enum CommitteeResolution {
179 TrustedNode,
184 Anchored {
186 committee: Committee,
188 chain_identifier: Option<ChainIdentifier>,
191 cache: Arc<dyn CommitteeCache>,
193 },
194}
195
196impl fmt::Debug for CommitteeResolution {
197 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
198 match self {
199 Self::TrustedNode => formatter.write_str("TrustedNode"),
200 Self::Anchored {
201 committee,
202 chain_identifier,
203 ..
204 } => formatter
205 .debug_struct("Anchored")
206 .field("committee", committee)
207 .field("chain_identifier", chain_identifier)
208 .finish_non_exhaustive(),
209 }
210 }
211}
212
213impl CommitteeResolution {
214 pub fn anchored(committee: Committee) -> Self {
218 Self::Anchored {
219 committee,
220 chain_identifier: None,
221 cache: Arc::new(MemoryCommitteeCache::new()),
222 }
223 }
224
225 pub fn anchored_with_cache(
230 chain_identifier: ChainIdentifier,
231 committee: Committee,
232 cache: impl CommitteeCache + 'static,
233 ) -> Self {
234 Self::Anchored {
235 committee,
236 chain_identifier: Some(chain_identifier),
237 cache: Arc::new(cache),
238 }
239 }
240
241 pub fn from_genesis(reader: impl Read) -> Result<Self, CommitteeResolutionError> {
246 Self::from_genesis_with_cache(reader, MemoryCommitteeCache::new())
247 }
248
249 pub fn from_genesis_with_cache(
254 reader: impl Read,
255 cache: impl CommitteeCache + 'static,
256 ) -> Result<Self, CommitteeResolutionError> {
257 Self::load_genesis(reader, cache).map_err(|kind| CommitteeResolutionError::new(0, kind))
258 }
259
260 fn load_genesis(
261 reader: impl Read,
262 cache: impl CommitteeCache + 'static,
263 ) -> Result<Self, CommitteeResolutionErrorKind> {
264 #[allow(dead_code)]
265 #[derive(Deserialize)]
266 struct GenesisBlob {
267 checkpoint: CertifiedCheckpointSummary,
268 checkpoint_contents: CheckpointContents,
269 transaction: Transaction,
270 effects: TransactionEffects,
271 events: TransactionEvents,
272 objects: Vec<Object>,
273 }
274
275 let genesis: GenesisBlob =
276 bcs::from_reader(reader).map_err(|source| CommitteeResolutionErrorKind::LoadGenesisCommittee {
277 source: Box::new(source),
278 })?;
279 let checkpoint_epoch = genesis.checkpoint.epoch();
280 if checkpoint_epoch != 0 {
281 return Err(CommitteeResolutionErrorKind::UnexpectedGenesisCheckpointEpoch {
282 epoch: checkpoint_epoch,
283 });
284 }
285
286 let objects = genesis.objects.as_slice();
287 let system_state =
288 get_iota_system_state(&objects).map_err(|source| CommitteeResolutionErrorKind::LoadGenesisCommittee {
289 source: Box::new(source),
290 })?;
291 let committee = system_state.get_current_epoch_committee().committee().clone();
292 if committee.epoch != 0 {
293 return Err(CommitteeResolutionErrorKind::UnexpectedGenesisCommitteeEpoch { epoch: committee.epoch });
294 }
295
296 genesis
297 .checkpoint
298 .verify_with_contents(&committee, Some(&genesis.checkpoint_contents))
299 .map_err(|source| CommitteeResolutionErrorKind::InvalidGenesisCheckpoint {
300 source: Box::new(source),
301 })?;
302 let chain_identifier = ChainIdentifier::from(*genesis.checkpoint.digest());
303
304 Ok(Self::anchored_with_cache(chain_identifier, committee, cache))
305 }
306}
307
308#[derive(Clone, Debug)]
314pub struct CommitteeResolver<S> {
315 source: S,
316 mode: CommitteeResolution,
317}
318
319impl<S> CommitteeResolver<S>
320where
321 S: Source,
322{
323 pub const fn new(source: S, resolution: CommitteeResolution) -> Self {
325 Self {
326 source,
327 mode: resolution,
328 }
329 }
330
331 pub async fn resolve(&self, target_epoch: EpochId) -> Result<Committee, CommitteeResolutionError> {
337 match &self.mode {
338 CommitteeResolution::TrustedNode => self.resolve_from_node(target_epoch).await,
339 CommitteeResolution::Anchored {
340 committee,
341 chain_identifier,
342 cache,
343 } => {
344 self.resolve_from_anchor(committee, *chain_identifier, cache.as_ref(), target_epoch)
345 .await
346 }
347 }
348 }
349
350 pub async fn verify<'proof>(&self, proof: &'proof Proof) -> Result<VerifiedProof<'proof>, ProofVerificationError> {
357 let committee = self
358 .resolve(proof.checkpoint_summary().epoch())
359 .await
360 .map_err(|source| ProofVerificationError::CommitteeResolution { source })?;
361
362 ProofVerifier::new(&committee)
363 .verify(proof)
364 .map_err(|source| ProofVerificationError::Proof { source })
365 }
366
367 async fn resolve_from_node(&self, target_epoch: EpochId) -> Result<Committee, CommitteeResolutionError> {
369 self.source.committee(target_epoch).await.map_err(|source| {
370 CommitteeResolutionError::new(
371 target_epoch,
372 CommitteeResolutionErrorKind::FetchCommittee {
373 epoch: target_epoch,
374 source: Box::new(source),
375 },
376 )
377 })
378 }
379
380 async fn resolve_from_anchor(
382 &self,
383 trusted_committee: &Committee,
384 chain_identifier: Option<ChainIdentifier>,
385 cache: &dyn CommitteeCache,
386 target_epoch: EpochId,
387 ) -> Result<Committee, CommitteeResolutionError> {
388 if target_epoch < trusted_committee.epoch {
389 return Err(CommitteeResolutionError::new(
390 target_epoch,
391 CommitteeResolutionErrorKind::TargetBeforeAnchor {
392 anchor_epoch: trusted_committee.epoch,
393 },
394 ));
395 }
396
397 if target_epoch == trusted_committee.epoch {
398 return Ok(trusted_committee.clone());
399 }
400
401 let target_key = Self::cache_key(chain_identifier, target_epoch);
402 if let Some(committee) = cache.committee(target_key).await.map_err(|source| {
403 CommitteeResolutionError::new(
404 target_epoch,
405 CommitteeResolutionErrorKind::Cache {
406 epoch: target_epoch,
407 source,
408 },
409 )
410 })? {
411 if committee.epoch != target_epoch {
412 return Err(CommitteeResolutionError::new(
413 target_epoch,
414 CommitteeResolutionErrorKind::Cache {
415 epoch: target_epoch,
416 source: CommitteeCacheError::Conflict { epoch: target_epoch },
417 },
418 ));
419 }
420
421 return Ok(committee);
422 }
423
424 let mut committee = trusted_committee.clone();
425
426 while committee.epoch < target_epoch {
427 let next_epoch = committee.epoch + 1;
428 let next_key = Self::cache_key(chain_identifier, next_epoch);
429 let Some(cached) = cache.committee(next_key).await.map_err(|source| {
430 CommitteeResolutionError::new(
431 target_epoch,
432 CommitteeResolutionErrorKind::Cache {
433 epoch: next_epoch,
434 source,
435 },
436 )
437 })?
438 else {
439 break;
440 };
441
442 if cached.epoch != next_epoch {
443 return Err(CommitteeResolutionError::new(
444 target_epoch,
445 CommitteeResolutionErrorKind::Cache {
446 epoch: next_epoch,
447 source: CommitteeCacheError::Conflict { epoch: next_epoch },
448 },
449 ));
450 }
451
452 committee = cached;
453 }
454
455 if committee.epoch == target_epoch {
456 return Ok(committee);
457 }
458
459 let current_epoch = self.current_epoch(target_epoch).await?;
460 if target_epoch > current_epoch {
461 return Err(CommitteeResolutionError::new(
462 target_epoch,
463 CommitteeResolutionErrorKind::TargetAheadOfNode { current_epoch },
464 ));
465 }
466
467 while committee.epoch < target_epoch {
468 let next_committee = self
469 .fetch_next_committee(target_epoch, chain_identifier, &committee, cache)
470 .await?;
471 committee = next_committee;
472 }
473
474 Ok(committee)
475 }
476
477 async fn current_epoch(&self, target_epoch: EpochId) -> Result<EpochId, CommitteeResolutionError> {
479 self.source
480 .current_epoch()
481 .await
482 .map_err(|source| {
483 CommitteeResolutionError::new(
484 target_epoch,
485 CommitteeResolutionErrorKind::FetchCurrentEpoch {
486 source: Box::new(source),
487 },
488 )
489 })?
490 .ok_or_else(|| {
491 CommitteeResolutionError::new(target_epoch, CommitteeResolutionErrorKind::MissingCurrentEpoch)
492 })
493 }
494
495 async fn fetch_next_committee(
497 &self,
498 target_epoch: EpochId,
499 chain_identifier: Option<ChainIdentifier>,
500 current_committee: &Committee,
501 cache: &dyn CommitteeCache,
502 ) -> Result<Committee, CommitteeResolutionError> {
503 let summary = self
504 .source
505 .epoch_close_summary(current_committee.epoch)
506 .await
507 .map_err(|source| {
508 CommitteeResolutionError::new(
509 target_epoch,
510 CommitteeResolutionErrorKind::FetchEpochHistory {
511 epoch: current_committee.epoch,
512 source: Box::new(source),
513 },
514 )
515 })?
516 .ok_or_else(|| {
517 CommitteeResolutionError::new(
518 target_epoch,
519 CommitteeResolutionErrorKind::MissingEpochCloseProof {
520 epoch: current_committee.epoch,
521 },
522 )
523 })?;
524
525 let sequence_number = summary.sequence_number;
526 let summary_epoch = summary.epoch();
527 if summary_epoch != current_committee.epoch {
528 return Err(CommitteeResolutionError::new(
529 target_epoch,
530 CommitteeResolutionErrorKind::InvalidEndOfEpochCheckpoint {
531 epoch: current_committee.epoch,
532 sequence_number,
533 source: Box::new(IotaError::WrongEpoch {
534 expected_epoch: current_committee.epoch,
535 actual_epoch: summary_epoch,
536 }),
537 },
538 ));
539 }
540
541 if summary.end_of_epoch_data.is_none() {
542 return Err(CommitteeResolutionError::new(
543 target_epoch,
544 CommitteeResolutionErrorKind::NotEndOfEpoch { sequence_number },
545 ));
546 }
547
548 let next_epoch = summary_epoch.checked_add(1).ok_or_else(|| {
549 CommitteeResolutionError::new(
550 target_epoch,
551 CommitteeResolutionErrorKind::NextEpochOverflow { epoch: summary_epoch },
552 )
553 })?;
554
555 let verified = summary.try_into_verified(current_committee).map_err(|source| {
556 CommitteeResolutionError::new(
557 target_epoch,
558 CommitteeResolutionErrorKind::InvalidEndOfEpochCheckpoint {
559 epoch: current_committee.epoch,
560 sequence_number,
561 source: Box::new(source),
562 },
563 )
564 })?;
565 let next_epoch_committee = &verified
566 .end_of_epoch_data
567 .as_ref()
568 .expect("checked before signature verification")
569 .next_epoch_committee;
570 let next_committee = Committee::from_committee_members(next_epoch, next_epoch_committee);
571
572 let cache_key = Self::cache_key(chain_identifier, next_committee.epoch);
573 cache.store(cache_key, &next_committee).await.map_err(|source| {
574 CommitteeResolutionError::new(
575 target_epoch,
576 CommitteeResolutionErrorKind::Cache {
577 epoch: next_committee.epoch,
578 source,
579 },
580 )
581 })?;
582
583 Ok(next_committee)
584 }
585
586 fn cache_key(chain_identifier: Option<ChainIdentifier>, epoch: EpochId) -> CommitteeCacheKey {
587 chain_identifier.map_or_else(
588 || CommitteeCacheKey::isolated(epoch),
589 |chain| CommitteeCacheKey::new(chain, epoch),
590 )
591 }
592}
593
594#[cfg(feature = "native-grpc")]
595impl CommitteeResolver<GrpcClient> {
596 pub const fn grpc_client(&self) -> &GrpcClient {
598 &self.source
599 }
600}
601
602#[cfg(test)]
603mod tests {
604 use std::collections::BTreeMap;
605 use std::sync::Mutex;
606
607 use iota_sdk_types::gas::GasCostSummary;
608 use iota_sdk_types::{CheckpointDigest, CheckpointSummary, EndOfEpochData, ObjectId, TransactionDigest, Version};
609 use iota_types::digests::ChainIdentifier;
610 use iota_types::messages_checkpoint::CertifiedCheckpointSummary;
611 use iota_types::object::Object;
612
613 use super::*;
614 use crate::{SourceCheckpoint, SourceError, SourceTransaction};
615
616 struct StaticCache {
617 key: CommitteeCacheKey,
618 committee: Committee,
619 }
620
621 struct FailingStoreCache;
622
623 #[derive(Clone)]
624 struct EpochCloseSource {
625 summary: CertifiedCheckpointSummary,
626 }
627
628 #[derive(Clone)]
629 struct CommitteeHistorySource {
630 current_epoch: Option<EpochId>,
631 summaries: BTreeMap<EpochId, CertifiedCheckpointSummary>,
632 fail_current_epoch: bool,
633 }
634
635 #[async_trait::async_trait]
636 impl Source for EpochCloseSource {
637 async fn chain_identifier(&self) -> Result<ChainIdentifier, SourceError> {
638 unreachable!("committee transition does not resolve a chain identifier")
639 }
640
641 async fn transaction(
642 &self,
643 _transaction_digest: TransactionDigest,
644 ) -> Result<Option<SourceTransaction>, SourceError> {
645 unreachable!("committee transition does not resolve transactions")
646 }
647
648 async fn object(&self, _object_id: ObjectId, _version: Option<Version>) -> Result<Option<Object>, SourceError> {
649 unreachable!("committee transition does not resolve objects")
650 }
651
652 async fn checkpoint(&self, _sequence_number: u64) -> Result<Option<SourceCheckpoint>, SourceError> {
653 unreachable!("committee transition does not resolve checkpoints")
654 }
655
656 async fn committee(&self, _epoch: EpochId) -> Result<Committee, SourceError> {
657 unreachable!("anchored committee transition does not trust node committees")
658 }
659
660 async fn current_epoch(&self) -> Result<Option<EpochId>, SourceError> {
661 Ok(Some(self.summary.epoch().saturating_add(1)))
662 }
663
664 async fn epoch_close_summary(
665 &self,
666 _epoch: EpochId,
667 ) -> Result<Option<CertifiedCheckpointSummary>, SourceError> {
668 Ok(Some(self.summary.clone()))
669 }
670 }
671
672 #[async_trait::async_trait]
673 impl Source for CommitteeHistorySource {
674 async fn chain_identifier(&self) -> Result<ChainIdentifier, SourceError> {
675 unreachable!("committee history does not resolve a chain identifier")
676 }
677
678 async fn transaction(
679 &self,
680 _transaction_digest: TransactionDigest,
681 ) -> Result<Option<SourceTransaction>, SourceError> {
682 unreachable!("committee history does not resolve transactions")
683 }
684
685 async fn object(&self, _object_id: ObjectId, _version: Option<Version>) -> Result<Option<Object>, SourceError> {
686 unreachable!("committee history does not resolve objects")
687 }
688
689 async fn checkpoint(&self, _sequence_number: u64) -> Result<Option<SourceCheckpoint>, SourceError> {
690 unreachable!("committee history does not resolve checkpoints")
691 }
692
693 async fn committee(&self, _epoch: EpochId) -> Result<Committee, SourceError> {
694 unreachable!("anchored committee history does not trust node committees")
695 }
696
697 async fn current_epoch(&self) -> Result<Option<EpochId>, SourceError> {
698 if self.fail_current_epoch {
699 return Err(SourceError::request(std::io::Error::other("current epoch unavailable")));
700 }
701
702 Ok(self.current_epoch)
703 }
704
705 async fn epoch_close_summary(&self, epoch: EpochId) -> Result<Option<CertifiedCheckpointSummary>, SourceError> {
706 Ok(self.summaries.get(&epoch).cloned())
707 }
708 }
709
710 #[derive(Clone, Default)]
711 struct RecordingCache {
712 stored: Arc<Mutex<Vec<Committee>>>,
713 }
714
715 impl RecordingCache {
716 fn stored(&self) -> Vec<Committee> {
717 self.stored.lock().unwrap().clone()
718 }
719 }
720
721 #[async_trait::async_trait]
722 impl CommitteeCache for RecordingCache {
723 async fn committee(&self, _key: CommitteeCacheKey) -> Result<Option<Committee>, CommitteeCacheError> {
724 Ok(None)
725 }
726
727 async fn store(&self, _key: CommitteeCacheKey, committee: &Committee) -> Result<(), CommitteeCacheError> {
728 self.stored.lock().unwrap().push(committee.clone());
729 Ok(())
730 }
731 }
732
733 #[async_trait::async_trait]
734 impl CommitteeCache for StaticCache {
735 async fn committee(&self, key: CommitteeCacheKey) -> Result<Option<Committee>, CommitteeCacheError> {
736 Ok((self.key == key).then(|| self.committee.clone()))
737 }
738
739 async fn store(&self, _key: CommitteeCacheKey, _committee: &Committee) -> Result<(), CommitteeCacheError> {
740 Ok(())
741 }
742 }
743
744 #[async_trait::async_trait]
745 impl CommitteeCache for FailingStoreCache {
746 async fn committee(&self, _key: CommitteeCacheKey) -> Result<Option<Committee>, CommitteeCacheError> {
747 Ok(None)
748 }
749
750 async fn store(&self, key: CommitteeCacheKey, _committee: &Committee) -> Result<(), CommitteeCacheError> {
751 Err(CommitteeCacheError::Backend {
752 epoch: key.epoch(),
753 source: Box::new(std::io::Error::other("cache unavailable")),
754 })
755 }
756 }
757
758 fn chain_identifier(byte: u8) -> ChainIdentifier {
759 ChainIdentifier::from(CheckpointDigest::new([byte; 32]))
760 }
761
762 fn committee_with_keypairs(epoch: EpochId, size: usize) -> (Committee, Vec<iota_types::crypto::AuthorityKeyPair>) {
763 let (base_committee, keypairs) = Committee::new_simple_test_committee_of_size(size);
764 let committee = Committee::new(epoch, base_committee.voting_rights.iter().cloned().collect());
765
766 (committee, keypairs)
767 }
768
769 fn signed_committee_transition(
770 current: &Committee,
771 keypairs: &[iota_types::crypto::AuthorityKeyPair],
772 next: &Committee,
773 ) -> CertifiedCheckpointSummary {
774 let summary = CheckpointSummary {
775 epoch: current.epoch,
776 sequence_number: current.epoch,
777 network_total_transactions: 0,
778 contents_digest: Default::default(),
779 previous_digest: None,
780 epoch_rolling_gas_cost_summary: GasCostSummary::default(),
781 timestamp_ms: 0,
782 checkpoint_commitments: Vec::new(),
783 end_of_epoch_data: Some(EndOfEpochData {
784 next_epoch_committee: next.committee_members(),
785 next_epoch_protocol_version: 1,
786 epoch_commitments: Vec::new(),
787 epoch_supply_change: 0,
788 }),
789 version_specific_data: Vec::new(),
790 };
791
792 CertifiedCheckpointSummary::new_from_keypairs_for_testing(summary, keypairs, current)
793 }
794
795 fn signed_end_of_epoch_summary(
796 current_epoch: EpochId,
797 include_next_committee: bool,
798 ) -> (Committee, Committee, CertifiedCheckpointSummary) {
799 let (base_committee, keypairs) = Committee::new_simple_test_committee();
800 signed_end_of_epoch_summary_from_test_committee(
801 current_epoch,
802 include_next_committee,
803 base_committee,
804 keypairs,
805 5,
806 )
807 }
808
809 fn signed_end_of_epoch_summary_with_sizes(
810 current_epoch: EpochId,
811 include_next_committee: bool,
812 current_committee_size: usize,
813 next_committee_size: usize,
814 ) -> (Committee, Committee, CertifiedCheckpointSummary) {
815 let (base_committee, keypairs) = Committee::new_simple_test_committee_of_size(current_committee_size);
816 signed_end_of_epoch_summary_from_test_committee(
817 current_epoch,
818 include_next_committee,
819 base_committee,
820 keypairs,
821 next_committee_size,
822 )
823 }
824
825 fn signed_end_of_epoch_summary_from_test_committee(
826 current_epoch: EpochId,
827 include_next_committee: bool,
828 base_committee: Committee,
829 keypairs: Vec<iota_types::crypto::AuthorityKeyPair>,
830 next_committee_size: usize,
831 ) -> (Committee, Committee, CertifiedCheckpointSummary) {
832 let current_committee = Committee::new(current_epoch, base_committee.voting_rights.iter().cloned().collect());
833 let (next_base_committee, _) = Committee::new_simple_test_committee_of_size(next_committee_size);
834 let next_committee = Committee::new(
835 current_epoch.saturating_add(1),
836 next_base_committee.voting_rights.iter().cloned().collect(),
837 );
838 let end_of_epoch_data = include_next_committee.then(|| EndOfEpochData {
839 next_epoch_committee: next_committee.committee_members(),
840 next_epoch_protocol_version: 1,
841 epoch_commitments: Vec::new(),
842 epoch_supply_change: 0,
843 });
844 let summary = CheckpointSummary {
845 epoch: current_epoch,
846 sequence_number: 42,
847 network_total_transactions: 0,
848 contents_digest: Default::default(),
849 previous_digest: None,
850 epoch_rolling_gas_cost_summary: GasCostSummary::default(),
851 timestamp_ms: 0,
852 checkpoint_commitments: Vec::new(),
853 end_of_epoch_data,
854 version_specific_data: Vec::new(),
855 };
856 let certified_summary =
857 CertifiedCheckpointSummary::new_from_keypairs_for_testing(summary, &keypairs, ¤t_committee);
858
859 (current_committee, next_committee, certified_summary)
860 }
861
862 #[tokio::test]
863 async fn authenticated_summary_stores_exactly_the_verified_committee() {
864 let (current_committee, expected_committee, summary) = signed_end_of_epoch_summary(3, true);
865 let cache = RecordingCache::default();
866 let resolver = CommitteeResolver::new(
867 EpochCloseSource { summary },
868 CommitteeResolution::anchored(current_committee.clone()),
869 );
870
871 let committee = resolver
872 .fetch_next_committee(4, None, ¤t_committee, &cache)
873 .await
874 .unwrap();
875
876 assert_eq!(committee, expected_committee);
877 assert_eq!(cache.stored(), vec![expected_committee]);
878 }
879
880 #[tokio::test]
881 async fn invalid_checkpoint_signature_never_reaches_the_cache() {
882 let (_, _, summary) = signed_end_of_epoch_summary(3, true);
883 let (wrong_committee, _) = Committee::new_simple_test_committee_of_size(6);
884 let wrong_committee = Committee::new(3, wrong_committee.voting_rights.iter().cloned().collect());
885 let cache = RecordingCache::default();
886 let resolver = CommitteeResolver::new(
887 EpochCloseSource { summary },
888 CommitteeResolution::anchored(wrong_committee.clone()),
889 );
890
891 let error = resolver
892 .fetch_next_committee(4, None, &wrong_committee, &cache)
893 .await
894 .unwrap_err();
895
896 assert!(matches!(
897 error.kind,
898 CommitteeResolutionErrorKind::InvalidEndOfEpochCheckpoint {
899 epoch: 3,
900 sequence_number: 42,
901 ..
902 }
903 ));
904 assert!(cache.stored().is_empty());
905 }
906
907 #[tokio::test]
908 async fn checkpoint_without_end_of_epoch_data_never_reaches_the_cache() {
909 let (current_committee, _, summary) = signed_end_of_epoch_summary(3, false);
910 let cache = RecordingCache::default();
911 let resolver = CommitteeResolver::new(
912 EpochCloseSource { summary },
913 CommitteeResolution::anchored(current_committee.clone()),
914 );
915
916 let error = resolver
917 .fetch_next_committee(4, None, ¤t_committee, &cache)
918 .await
919 .unwrap_err();
920
921 assert!(matches!(
922 error.kind,
923 CommitteeResolutionErrorKind::NotEndOfEpoch { sequence_number: 42 }
924 ));
925 assert!(cache.stored().is_empty());
926 }
927
928 #[tokio::test]
929 async fn end_of_epoch_structure_is_checked_before_signatures() {
930 let (_, _, summary) = signed_end_of_epoch_summary(3, false);
931 let (wrong_committee, _) = Committee::new_simple_test_committee_of_size(6);
932 let wrong_committee = Committee::new(3, wrong_committee.voting_rights.iter().cloned().collect());
933 let cache = RecordingCache::default();
934 let resolver = CommitteeResolver::new(
935 EpochCloseSource { summary },
936 CommitteeResolution::anchored(wrong_committee.clone()),
937 );
938
939 let error = resolver
940 .fetch_next_committee(4, None, &wrong_committee, &cache)
941 .await
942 .unwrap_err();
943
944 assert!(matches!(
945 error.kind,
946 CommitteeResolutionErrorKind::NotEndOfEpoch { sequence_number: 42 }
947 ));
948 assert!(cache.stored().is_empty());
949 }
950
951 #[tokio::test]
952 async fn wrong_epoch_summary_does_not_advance_or_reach_the_cache() {
953 let (signing_committee, _, summary) = signed_end_of_epoch_summary(4, true);
954 let expected_committee = Committee::new(3, signing_committee.voting_rights.iter().cloned().collect());
955 let cache = RecordingCache::default();
956 let resolver = CommitteeResolver::new(
957 EpochCloseSource { summary },
958 CommitteeResolution::anchored(expected_committee.clone()),
959 );
960
961 let error = resolver
962 .fetch_next_committee(4, None, &expected_committee, &cache)
963 .await
964 .unwrap_err();
965
966 assert!(matches!(
967 error.kind,
968 CommitteeResolutionErrorKind::InvalidEndOfEpochCheckpoint {
969 epoch: 3,
970 sequence_number: 42,
971 ..
972 }
973 ));
974 assert!(cache.stored().is_empty());
975 }
976
977 #[tokio::test]
978 async fn overflowing_next_epoch_never_reaches_the_cache() {
979 let (current_committee, _, summary) = signed_end_of_epoch_summary(EpochId::MAX, true);
980 let cache = RecordingCache::default();
981 let resolver = CommitteeResolver::new(
982 EpochCloseSource { summary },
983 CommitteeResolution::anchored(current_committee.clone()),
984 );
985
986 let error = resolver
987 .fetch_next_committee(EpochId::MAX, None, ¤t_committee, &cache)
988 .await
989 .unwrap_err();
990
991 assert!(matches!(
992 error.kind,
993 CommitteeResolutionErrorKind::NextEpochOverflow { epoch: EpochId::MAX }
994 ));
995 assert!(cache.stored().is_empty());
996 }
997
998 #[tokio::test]
999 async fn anchor_mode_uses_a_committee_cache_by_default() {
1000 let (current_committee, next_committee, _) = signed_end_of_epoch_summary(3, true);
1001 let client = GrpcClient::new("http://127.0.0.1:1").unwrap();
1002 let resolver = CommitteeResolver::new(client, CommitteeResolution::anchored(current_committee));
1003 let CommitteeResolution::Anchored { cache, .. } = &resolver.mode else {
1004 panic!("anchor resolver must have a committee cache");
1005 };
1006 cache
1007 .store(CommitteeCacheKey::isolated(next_committee.epoch), &next_committee)
1008 .await
1009 .unwrap();
1010
1011 let resolved = resolver.resolve(4).await.unwrap();
1012
1013 assert_eq!(resolved, next_committee);
1014 }
1015
1016 #[tokio::test]
1017 async fn anchored_resolution_accepts_a_committee_from_a_trusted_cache() {
1018 let (current_committee, next_committee, _) = signed_end_of_epoch_summary(3, true);
1019 let chain_identifier = chain_identifier(1);
1020 let cache = StaticCache {
1021 key: CommitteeCacheKey::new(chain_identifier, next_committee.epoch),
1022 committee: next_committee.clone(),
1023 };
1024 let client = GrpcClient::new("http://127.0.0.1:1").unwrap();
1025 let resolver = CommitteeResolver::new(
1026 client,
1027 CommitteeResolution::anchored_with_cache(chain_identifier, current_committee, cache),
1028 );
1029
1030 let resolved = resolver.resolve(4).await.unwrap();
1031
1032 assert_eq!(resolved, next_committee);
1033 }
1034
1035 #[tokio::test]
1036 async fn target_cache_entry_must_contain_the_requested_epoch() {
1037 let (anchor, _) = committee_with_keypairs(3, 4);
1038 let (mislabeled, _) = committee_with_keypairs(5, 5);
1039 let chain_identifier = chain_identifier(1);
1040 let cache = StaticCache {
1041 key: CommitteeCacheKey::new(chain_identifier, 4),
1042 committee: mislabeled,
1043 };
1044 let resolver = CommitteeResolver::new(
1045 CommitteeHistorySource {
1046 current_epoch: Some(5),
1047 summaries: BTreeMap::new(),
1048 fail_current_epoch: false,
1049 },
1050 CommitteeResolution::anchored_with_cache(chain_identifier, anchor, cache),
1051 );
1052
1053 let error = resolver.resolve(4).await.unwrap_err();
1054
1055 assert!(matches!(
1056 error.kind,
1057 CommitteeResolutionErrorKind::Cache {
1058 epoch: 4,
1059 source: CommitteeCacheError::Conflict { epoch: 4 }
1060 }
1061 ));
1062 }
1063
1064 #[tokio::test]
1065 async fn intermediate_cache_entry_must_contain_its_key_epoch() {
1066 let (anchor, _) = committee_with_keypairs(3, 4);
1067 let (mislabeled, _) = committee_with_keypairs(5, 5);
1068 let chain_identifier = chain_identifier(1);
1069 let cache = StaticCache {
1070 key: CommitteeCacheKey::new(chain_identifier, 4),
1071 committee: mislabeled,
1072 };
1073 let resolver = CommitteeResolver::new(
1074 CommitteeHistorySource {
1075 current_epoch: Some(5),
1076 summaries: BTreeMap::new(),
1077 fail_current_epoch: false,
1078 },
1079 CommitteeResolution::anchored_with_cache(chain_identifier, anchor, cache),
1080 );
1081
1082 let error = resolver.resolve(5).await.unwrap_err();
1083
1084 assert!(matches!(
1085 error.kind,
1086 CommitteeResolutionErrorKind::Cache {
1087 epoch: 4,
1088 source: CommitteeCacheError::Conflict { epoch: 4 }
1089 }
1090 ));
1091 }
1092
1093 #[tokio::test]
1094 async fn anchored_resolution_reports_a_current_epoch_source_failure() {
1095 let (anchor, _) = committee_with_keypairs(0, 4);
1096 let resolver = CommitteeResolver::new(
1097 CommitteeHistorySource {
1098 current_epoch: None,
1099 summaries: BTreeMap::new(),
1100 fail_current_epoch: true,
1101 },
1102 CommitteeResolution::anchored(anchor),
1103 );
1104
1105 let error = resolver.resolve(1).await.unwrap_err();
1106
1107 assert!(matches!(
1108 error.kind,
1109 CommitteeResolutionErrorKind::FetchCurrentEpoch { .. }
1110 ));
1111 }
1112
1113 #[tokio::test]
1114 async fn anchored_resolution_reports_missing_epoch_close_evidence() {
1115 let (anchor, _) = committee_with_keypairs(0, 4);
1116 let resolver = CommitteeResolver::new(
1117 CommitteeHistorySource {
1118 current_epoch: Some(1),
1119 summaries: BTreeMap::new(),
1120 fail_current_epoch: false,
1121 },
1122 CommitteeResolution::anchored(anchor),
1123 );
1124
1125 let error = resolver.resolve(1).await.unwrap_err();
1126
1127 assert!(matches!(
1128 error.kind,
1129 CommitteeResolutionErrorKind::MissingEpochCloseProof { epoch: 0 }
1130 ));
1131 }
1132
1133 #[tokio::test]
1134 async fn anchored_resolution_reports_a_cache_store_failure() {
1135 let (anchor, keypairs) = committee_with_keypairs(0, 4);
1136 let (next, _) = committee_with_keypairs(1, 5);
1137 let summary = signed_committee_transition(&anchor, &keypairs, &next);
1138 let resolver = CommitteeResolver::new(
1139 CommitteeHistorySource {
1140 current_epoch: Some(1),
1141 summaries: BTreeMap::from([(0, summary)]),
1142 fail_current_epoch: false,
1143 },
1144 CommitteeResolution::anchored_with_cache(chain_identifier(1), anchor, FailingStoreCache),
1145 );
1146
1147 let error = resolver.resolve(1).await.unwrap_err();
1148
1149 assert!(matches!(
1150 error.kind,
1151 CommitteeResolutionErrorKind::Cache {
1152 epoch: 1,
1153 source: CommitteeCacheError::Backend { epoch: 1, .. }
1154 }
1155 ));
1156 }
1157
1158 #[tokio::test]
1159 async fn anchored_resolution_authenticates_multiple_epoch_transitions() {
1160 let (first, first_keypairs) = committee_with_keypairs(0, 4);
1161 let (second, second_keypairs) = committee_with_keypairs(1, 5);
1162 let (expected, _) = committee_with_keypairs(2, 6);
1163 let summaries = BTreeMap::from([
1164 (0, signed_committee_transition(&first, &first_keypairs, &second)),
1165 (1, signed_committee_transition(&second, &second_keypairs, &expected)),
1166 ]);
1167 let resolver = CommitteeResolver::new(
1168 CommitteeHistorySource {
1169 current_epoch: Some(2),
1170 summaries,
1171 fail_current_epoch: false,
1172 },
1173 CommitteeResolution::anchored(first),
1174 );
1175
1176 let resolved = resolver.resolve(2).await.unwrap();
1177
1178 assert_eq!(resolved, expected);
1179 }
1180
1181 #[tokio::test]
1182 async fn shared_cache_isolated_between_distinct_networks() {
1183 let (_, first_successor, _) = signed_end_of_epoch_summary(3, true);
1184 let (second_anchor, second_successor, second_summary) = signed_end_of_epoch_summary_with_sizes(3, true, 6, 6);
1185 let first_chain = chain_identifier(1);
1186 let second_chain = chain_identifier(2);
1187 let cache = MemoryCommitteeCache::new();
1188 cache
1189 .store(
1190 CommitteeCacheKey::new(first_chain, first_successor.epoch),
1191 &first_successor,
1192 )
1193 .await
1194 .unwrap();
1195 let resolver = CommitteeResolver::new(
1196 EpochCloseSource {
1197 summary: second_summary,
1198 },
1199 CommitteeResolution::anchored_with_cache(second_chain, second_anchor.clone(), cache.clone()),
1200 );
1201
1202 let resolved = resolver.resolve(4).await.unwrap();
1203
1204 assert_eq!(resolved, second_successor);
1205 assert_ne!(resolved, first_successor);
1206 assert_eq!(cache.len().await, 2);
1207 assert_eq!(
1208 cache
1209 .committee(CommitteeCacheKey::new(second_chain, resolved.epoch))
1210 .await
1211 .unwrap(),
1212 Some(resolved)
1213 );
1214 }
1215}