//! DKG/reshare crash-recovery storage. //! //! This store exists only to recover state a restarted actor cannot otherwise //! re-obtain. After a crash the actor does not re-receive P2P messages and //! marshal does not re-deliver finalized blocks, so this store keeps a plaintext //! journal of the public messages a restarted node would otherwise lose: //! //! - public dealer messages and player acknowledgements, so a player can rebuild //! the acks it already emitted (its private dealings are recovered from //! [`SecretStore`]); //! - finalized dealer logs observed during inclusion. //! //! The current epoch's public state is not persisted here: ordinary recovery //! re-derives it from the finalized boundary block, while state-sync recovery //! uses the shared [`state_sync::Plan`](crate::dkg::state_sync::Plan). Everything //! secret stays out of this plaintext store and is held only through //! [`SecretStore`]: shares, private dealings, and the dealer RNG seed (which //! seeds the dealer polynomial and so reveals every share that dealer //! distributes). use crate::dkg::{SecretStore, network::Directory, types::EpochInfo}; use bytes::{Buf, BufMut}; use commonware_codec::{EncodeSize, Error as CodecError, Read, ReadExt, Write}; use commonware_consensus::types::Epoch; use commonware_cryptography::{ BatchVerifier, PublicKey, Signer, bls12381::{ dkg::feldman_desmedt::{ Dealer as CryptoDealer, DealerLog, DealerMessageError as DkgDealerMessageError, DealerPrivMsg, DealerPubMsg, Error as DkgError, FinalizeError as DkgFinalizeError, Info, Logs, Output, Player as CryptoPlayer, PlayerAck, PlayerAckError as DkgAckError, SignedDealerLog, }, primitives::{group, variant::Variant}, }, transcript::{Summary, Transcript, Version}, }; use commonware_math::algebra::Random; use commonware_parallel::Strategy; use commonware_runtime::{ BufferPooler, Clock, Metrics, ReadOptions, Storage as RuntimeStorage, buffer::paged::CacheRef, }; use commonware_storage::journal::{ self, segmented::variable::{Config as JournalConfig, Journal}, }; use commonware_utils::{Faults, N3f1, NZU16, NZUsize, futures::rebind, sequence::Unit}; use rand_core::CryptoRng; use std::{ collections::BTreeMap, num::{NonZeroU16, NonZeroU32, NonZeroUsize}, }; use tracing::{debug, warn}; const PAGE_SIZE: NonZeroU16 = NZU16!(1 << 12); // 4 KiB const PAGE_CACHE_CAPACITY: NonZeroUsize = NZUsize!(1 << 13); // 8 KiB const WRITE_BUFFER: NonZeroUsize = NZUsize!(1 << 12); // 4 KiB const READ_BUFFER: NonZeroUsize = NZUsize!(1 << 20); // 1 MiB enum Event { Dealing(P, DealerPubMsg), Ack(P, PlayerAck

), Log(P, DealerLog), } impl EncodeSize for Event { fn encode_size(&self) -> usize { 1 + match self { Self::Dealing(dealer, public) => dealer.encode_size() + public.encode_size(), Self::Ack(player, ack) => player.encode_size() + ack.encode_size(), Self::Log(dealer, log) => dealer.encode_size() + log.encode_size(), } } } impl Write for Event { fn write(&self, writer: &mut impl BufMut) { match self { Self::Dealing(dealer, public) => { 0u8.write(writer); dealer.write(writer); public.write(writer); } Self::Ack(player, ack) => { 1u8.write(writer); player.write(writer); ack.write(writer); } Self::Log(dealer, log) => { 2u8.write(writer); dealer.write(writer); log.write(writer); } } } } impl Read for Event { type Cfg = NonZeroU32; fn read_cfg(reader: &mut impl Buf, cfg: &Self::Cfg) -> Result { match u8::read(reader)? { 0 => Ok(Self::Dealing( ReadExt::read(reader)?, Read::read_cfg(reader, cfg)?, )), 1 => Ok(Self::Ack(ReadExt::read(reader)?, ReadExt::read(reader)?)), 2 => Ok(Self::Log( ReadExt::read(reader)?, Read::read_cfg(reader, cfg)?, )), tag => Err(CodecError::InvalidEnum(tag)), } } } #[cfg(feature = "arbitrary")] impl arbitrary::Arbitrary<'_> for Event where P: for<'a> arbitrary::Arbitrary<'a>, DealerPubMsg: for<'a> arbitrary::Arbitrary<'a>, DealerLog: for<'a> arbitrary::Arbitrary<'a>, PlayerAck

: for<'a> arbitrary::Arbitrary<'a>, { fn arbitrary(u: &mut arbitrary::Unstructured<'_>) -> arbitrary::Result { Ok(match u.int_in_range(0..=2)? { 0 => Self::Dealing(u.arbitrary()?, u.arbitrary()?), 1 => Self::Ack(u.arbitrary()?, u.arbitrary()?), _ => Self::Log(u.arbitrary()?, u.arbitrary()?), }) } } struct EpochCache { dealings: BTreeMap, DealerPrivMsg)>, acks: BTreeMap>, logs: BTreeMap>, } impl Default for EpochCache { fn default() -> Self { Self { dealings: BTreeMap::new(), acks: BTreeMap::new(), logs: BTreeMap::new(), } } } /// DKG/reshare crash-recovery store. /// /// The plaintext side holds only the dealer-message, acknowledgement, and /// finalized-log journal. The current epoch's public state comes from finalized /// boundary blocks or the separate state-sync store, not this journal. All /// secret material (shares, private dealings, and the dealer RNG seed) is held only through /// [`SecretStore`], never in plaintext. pub struct Store where E: BufferPooler + Clock + RuntimeStorage + Metrics, SS: SecretStore, V: Variant, P: PublicKey, D: Directory

, { secret_store: SS, events: Option>>, current: Option>, epochs: BTreeMap>, } impl Store where E: BufferPooler + Clock + RuntimeStorage + Metrics, SS: SecretStore, V: Variant, P: PublicKey, D: Directory

, { /// Initializes the store and replays durable crash-recovery state. pub async fn init( context: E, partition_prefix: &str, max_participants: NonZeroU32, mut secret_store: SS, ) -> Self { let page_cache = CacheRef::from_pooler(&context, PAGE_SIZE, PAGE_CACHE_CAPACITY); let events = Journal::init( context.child("events"), JournalConfig { partition: format!("{partition_prefix}_events"), compression: None, codec_config: max_participants, page_cache, write_buffer: WRITE_BUFFER, }, ) .await .expect("failed to initialize reshare event journal"); // The current epoch is not persisted: it is re-derived from finalized // boundary blocks by the setup state, so a restarted store starts with no // current epoch. let current = None; let mut epochs = BTreeMap::>::new(); let events = { // Replay rebuilds the epoch caches in memory, so journal pages need // not remain in the OS page cache. let mut replay = events .replay(0, 0, READ_BUFFER, ReadOptions::DONT_CACHE) .await .expect("failed to replay reshare events"); while let Some(result) = replay.next().await { let (section, _, _, event) = result.expect("failed to read reshare event"); let epoch = Epoch::new(section); let cache = epochs.entry(epoch).or_default(); match event { Event::Dealing(dealer, public) => { let private = secret_store.get_dealing(epoch, &dealer).await; if let Some(private) = private { cache.dealings.insert(dealer, (public, private)); } } Event::Ack(player, ack) => { cache.acks.insert(player, ack); } Event::Log(dealer, log) => { cache.logs.insert(dealer, log); } } } replay.finish().expect("failed to replay reshare events") }; Self { secret_store, events: Some(events), current, epochs, } } /// Returns the current epoch state, if one has been entered. pub fn current(&self) -> Option> { self.current.clone() } /// Returns the share for `epoch`, if any. pub async fn share(&mut self, epoch: Epoch) -> Option { self.secret_store.get_share(epoch).await } /// Returns the persisted dealer RNG seed for `epoch`, if any. /// /// A seed exists only for an epoch this node previously entered. Reusing it /// keeps dealer randomness identical across a restart. pub async fn seed(&mut self, epoch: Epoch) -> Option

{ self.secret_store.get_seed(epoch).await } /// Returns the persisted dealer RNG seed for `epoch`, generating a fresh /// random seed from `rng` if none exists. pub(crate) async fn seed_or_random(&mut self, epoch: Epoch, rng: impl CryptoRng) -> Summary { self.seed(epoch) .await .unwrap_or_else(|| Summary::random(rng)) } /// Persists the dealer RNG seed for `epoch`. pub(crate) async fn put_seed(&mut self, epoch: Epoch, rng_seed: Summary) { self.secret_store.put_seed(epoch, rng_seed).await; } /// Advances to `info`, persisting its secrets before the current epoch moves. /// /// The public artifact in `info` is read from the finalized boundary block, /// not from local state, and is held only in memory. The share is persisted /// only when it matches that finalized truth; otherwise the node continues as /// an observer. pub async fn commit_epoch( &mut self, info: EpochInfo, rng_seed: Summary, share: Option, ) { let epoch = info.epoch; if let Some(share) = share { self.secret_store.put_share(epoch, share).await; } self.secret_store.put_seed(epoch, rng_seed).await; self.current = Some(info); } /// Prunes public recovery data and secret material older than `min`. pub async fn prune(&mut self, min: Epoch) { self.epochs.retain(|epoch, _| *epoch >= min); // Prune the recovery journal and the secret store concurrently; they are // independent backends. let secret = &mut self.secret_store; futures::join!( async { rebind(&mut self.events, |events| events.prune(min.get())) .await .expect("failed to prune reshare events"); }, secret.prune(min), ); } fn cache(&mut self, epoch: Epoch) -> &mut EpochCache { self.epochs.entry(epoch).or_default() } /// Returns finalized dealer logs for `epoch`. pub fn logs(&self, epoch: Epoch) -> BTreeMap> { self.epochs .get(&epoch) .map(|cache| cache.logs.clone()) .unwrap_or_default() } /// Returns true if `dealer` already has a finalized log recorded. pub fn has_log(&self, epoch: Epoch, dealer: &P) -> bool { self.epochs .get(&epoch) .is_some_and(|cache| cache.logs.contains_key(dealer)) } fn dealings(&self, epoch: Epoch) -> Vec<(P, DealerPubMsg, DealerPrivMsg)> { self.epochs .get(&epoch) .map(|cache| { cache .dealings .iter() .map(|(dealer, (public, private))| { (dealer.clone(), public.clone(), private.clone()) }) .collect() }) .unwrap_or_default() } fn acks(&self, epoch: Epoch) -> Vec<(P, PlayerAck

)> { self.epochs .get(&epoch) .map(|cache| { cache .acks .iter() .map(|(player, ack)| (player.clone(), ack.clone())) .collect() }) .unwrap_or_default() } /// Persists a public/private dealing already accepted by the cryptographic player. async fn append_dealing( &mut self, epoch: Epoch, dealer: P, public: DealerPubMsg, private: DealerPrivMsg, ) -> bool { if self .epochs .get(&epoch) .is_some_and(|cache| cache.dealings.contains_key(&dealer)) { return false; } // Persist the private dealing (secret store) and the public dealer message // (recovery journal) concurrently. Both are durable before this returns, so // the ack the caller emits next is always backed by recoverable state. A // crash mid-write is safe: replay loads a dealing only when both its public // and private parts survived. let event = Event::Dealing(dealer.clone(), public.clone()); let secret = &mut self.secret_store; futures::join!( secret.put_dealing(epoch, dealer.clone(), private.clone()), async { rebind(&mut self.events, |events| { append_synced(events, epoch, &event) }) .await .expect("failed to record reshare dealing"); }, ); self.cache(epoch).dealings.insert(dealer, (public, private)); true } async fn append_ack(&mut self, epoch: Epoch, player: P, ack: PlayerAck

) -> bool { if self .epochs .get(&epoch) .is_some_and(|cache| cache.acks.contains_key(&player)) { return false; } let event = Event::Ack(player.clone(), ack.clone()); rebind(&mut self.events, |events| { append_synced(events, epoch, &event) }) .await .expect("failed to record reshare ack"); self.cache(epoch).acks.insert(player, ack); true } /// Records a finalized dealer log. pub async fn append_log(&mut self, epoch: Epoch, dealer: P, log: DealerLog) -> bool { if self.has_log(epoch, &dealer) { return false; } let event = Event::Log(dealer.clone(), log.clone()); rebind(&mut self.events, |events| { append_synced(events, epoch, &event) }) .await .expect("failed to record reshare log"); self.cache(epoch).logs.insert(dealer, log); true } /// Replays dealer state for `epoch`. pub fn create_dealer( &self, epoch: Epoch, signer: C, info: Info, share: Option, rng_seed: Summary, ) -> Option> where C: Signer, M: Faults, { if self.has_log(epoch, &signer.public_key()) { return None; } let (mut dealer, public, private) = CryptoDealer::start::( Transcript::resume(rng_seed, Version::V1).noise(b"dealer-rng"), info, signer, share, ) .expect("failed to create reshare dealer"); let mut unsent: BTreeMap = private.into_iter().collect(); for (player, ack) in self.acks(epoch) { if unsent.contains_key(&player) && dealer.receive_player_ack(player.clone(), ack).is_ok() { unsent.remove(&player); debug!(?epoch, ?player, "replayed reshare ack"); } } Some(Dealer { dealer: Some(dealer), public, unsent, finalized: None, }) } /// Replays player state for `epoch`. /// /// Returns `None` when this node observes an on-chain dealer log carrying its /// own acknowledgement but has lost the matching private dealing; see /// [`resume_player`](Self::resume_player). pub fn create_player( &self, epoch: Epoch, signer: C, info: Info, ) -> Option> where C: Signer, M: Faults, { self.resume_player::(epoch, signer, info, &self.logs(epoch)) } /// Replays player state using a supplied, non-durable log view. /// /// Returns `None` under the same missing-dealing condition as /// [`create_player`](Self::create_player). pub fn create_player_with_logs( &self, epoch: Epoch, signer: C, info: Info, logs: &BTreeMap>, ) -> Option> where C: Signer, M: Faults, { self.resume_player::(epoch, signer, info, logs) } /// Resumes player state for `epoch` from `logs`. /// /// Degrades to observer mode (returns `None`) when [`CryptoPlayer::resume`] /// reports [`MissingPlayerDealing`](DkgError::MissingPlayerDealing): a /// finalized dealer log carries this node's valid acknowledgement, but the /// matching private dealing is absent from the secret store (for example, a /// secret store restored from a backup taken before the ack was emitted). /// The ceremony has otherwise succeeded, so the node commits the epoch as a /// share-less verifier rather than crash-looping. Any other error signals /// genuine corruption and stays fatal. fn resume_player( &self, epoch: Epoch, signer: C, info: Info, logs: &BTreeMap>, ) -> Option> where C: Signer, M: Faults, { match CryptoPlayer::resume::(info, signer, logs, self.dealings(epoch)) { Ok((player, acks)) => Some(Player { player, acks }), Err(DkgError::MissingPlayerDealing) => { warn!( ?epoch, "missing private dealing on resume; entering epoch as observer" ); None } Err(error) => panic!("failed to resume reshare player: {error:?}"), } } } /// Appends `event` to the recovery journal for `epoch` and flushes it durably. async fn append_synced( events: Journal>, epoch: Epoch, event: &Event, ) -> Result>, journal::Error> where E: BufferPooler + Clock + RuntimeStorage + Metrics, V: Variant, P: PublicKey, { let section = epoch.get(); let (events, _, _) = events.append(section, event).await?; events.sync(section).await } /// Dealer state for one epoch. pub struct Dealer { dealer: Option>, public: DealerPubMsg, unsent: BTreeMap, finalized: Option>, } /// Successful outcome of processing a player acknowledgement. #[derive(Clone, Copy, Debug, Eq, PartialEq)] pub(crate) enum AckOutcome { /// The acknowledgement was newly persisted. Recorded, /// The acknowledgement was already recorded. Duplicate, } impl Dealer { /// Returns whether the recorded acknowledgements form a quorum. pub fn has_acknowledgement_quorum(&self, players: usize) -> bool { self.unsent.len() <= M::max_faults(players) as usize } /// Records a player ack. /// /// Returns [`AckOutcome::Recorded`] for a new acknowledgement and /// [`AckOutcome::Duplicate`] for one already recorded. /// [`DkgAckError::UnexpectedPlayer`] identifies a sender outside the round's /// player set, while [`DkgAckError::InvalidAck`] identifies a signature that /// does not match this dealer's transcript. /// /// # Panics /// /// Panics if called after the dealer has finalized. pub async fn handle( &mut self, store: &mut Store, epoch: Epoch, player: C::PublicKey, ack: PlayerAck, ) -> Result where E: BufferPooler + Clock + RuntimeStorage + Metrics, SS: SecretStore, D: Directory, { let dealer = self .dealer .as_mut() .expect("acknowledgements are handled only before dealer finalization"); dealer.receive_player_ack(player.clone(), ack.clone())?; if self.unsent.remove(&player).is_none() { return Ok(AckOutcome::Duplicate); } assert!( store.append_ack(epoch, player, ack).await, "pending acknowledgement must not already be persisted" ); Ok(AckOutcome::Recorded) } /// Finalizes once and returns true if a new log became available. pub fn finalize(&mut self) -> bool { if self.finalized.is_some() { return false; } let Some(dealer) = self.dealer.take() else { return false; }; self.finalized = Some(dealer.finalize::()); true } /// Returns a cloned finalized log. pub fn finalized(&self) -> Option> { self.finalized.clone() } /// Clears a finalized log after it is observed in a finalized block. pub fn clear_finalized(&mut self) { self.finalized = None; } /// Returns private dealings that still need to be sent. pub fn shares_to_distribute( &self, ) -> impl Iterator, DealerPrivMsg)> + '_ { self.unsent .iter() .map(|(player, private)| (player.clone(), self.public.clone(), private.clone())) } } /// Player state for one epoch. pub struct Player { player: CryptoPlayer, acks: BTreeMap>, } impl Player { /// Handles a dealer message, persisting it before returning the ack. /// /// A previously processed dealer returns its cached acknowledgement. A first /// message returns the exact [`DkgDealerMessageError`] for an unexpected /// dealer, invalid commitment degree, mismatched reshare commitment, or /// invalid private share. Rejected messages are not persisted. pub async fn handle( &mut self, store: &mut Store, epoch: Epoch, dealer: C::PublicKey, public: DealerPubMsg, private: DealerPrivMsg, ) -> Result, DkgDealerMessageError> where E: BufferPooler + Clock + RuntimeStorage + Metrics, SS: SecretStore, D: Directory, { if let Some(ack) = self.acks.get(&dealer) { return Ok(ack.clone()); } let ack = self .player .dealer_message::(dealer.clone(), public.clone(), private.clone())? .expect("processed dealer must have a cached acknowledgement"); store .append_dealing(epoch, dealer.clone(), public, private) .await; self.acks.insert(dealer, ack.clone()); Ok(ack) } /// Finalizes and returns the public output plus local share. #[allow(clippy::type_complexity)] pub fn finalize( self, rng: &mut impl CryptoRng, logs: Logs, strategy: &impl Strategy, ) -> Result<(Output, group::Share), DkgFinalizeError> where M: Faults, B: BatchVerifier, { self.player.finalize::(rng, logs, strategy) } } #[cfg(test)] mod tests { use super::*; use crate::dkg::{ tests::mocks::MemorySecretStore, types::{EpochInfo, EpochOutcome}, }; use commonware_codec::FixedSize; use commonware_consensus::types::Epoch; use commonware_cryptography::{ Signer, bls12381::{ dkg::feldman_desmedt::{Info, Output, Reveal, deal}, primitives::{sharing::Mode, variant::MinPk}, }, ed25519::{PrivateKey, PublicKey}, }; use commonware_runtime::{Runner, Supervisor as _, deterministic}; use commonware_utils::{N3f1, NZU32, TestRng, ordered::Set}; type TestStore = Store; fn summary(seed: u8) -> Summary { let bytes = [seed; Summary::SIZE]; Summary::read(&mut bytes.as_ref()).expect("valid summary") } fn output(seed: u64) -> Output { let (output, _) = deal::( TestRng::new(seed), Mode::NonZeroCounter, players(&signers()), ) .expect("trusted deal"); output } fn epoch_info(epoch: Epoch, output: Output) -> EpochInfo { EpochInfo { outcome: EpochOutcome::Success, epoch, output, players: Set::default(), next_players: Set::default(), directory: Unit, } } fn signers() -> Vec { (0..4).map(PrivateKey::from_seed).collect() } fn players(signers: &[PrivateKey]) -> Set { Set::from_iter_dedup(signers.iter().map(Signer::public_key)) } async fn init_store( context: E, partition: &str, secret_store: MemorySecretStore, ) -> TestStore where E: BufferPooler + Clock + RuntimeStorage + Metrics, { Store::init(context, partition, NZU32!(16), secret_store).await } #[test] fn commit_epoch_seeds_configured_epoch() { let executor = deterministic::Runner::default(); executor.start(|context| async move { let secret_store = MemorySecretStore::default(); let mut store = init_store(context.child("store"), "configured-start", secret_store).await; store .commit_epoch(epoch_info(Epoch::new(7), output(1)), summary(1), None) .await; let info = store.current().expect("current epoch"); assert_eq!(info.epoch, Epoch::new(7)); }); } #[test] fn replay_restores_dealings_acks_and_logs() { let executor = deterministic::Runner::default(); executor.start(|context| async move { let secret_store = MemorySecretStore::default(); let mut store = init_store(context.child("store"), "replay", secret_store.clone()).await; let signers = signers(); let players = players(&signers); let info = Info::new::( b"_COMMONWARE_GLUE_DKG_RESHARE_STORE_TEST", 0, None, Mode::NonZeroCounter, Reveal::V1, players.clone(), players.clone(), ) .expect("valid info"); store .commit_epoch(epoch_info(Epoch::zero(), output(1)), summary(1), None) .await; let dealer_pk = signers[0].public_key(); let player_pk = signers[1].public_key(); let mut dealer = store .create_dealer::( Epoch::zero(), signers[0].clone(), info.clone(), None, summary(2), ) .expect("dealer"); let mut player = store .create_player::(Epoch::zero(), signers[1].clone(), info.clone()) .expect("player"); let (_, public, private) = dealer .shares_to_distribute() .find(|(recipient, _, _)| *recipient == player_pk) .expect("dealing for player"); let duplicate_public = public.clone(); let duplicate_private = private.clone(); let ack = player .handle( &mut store, Epoch::zero(), dealer_pk.clone(), public, private, ) .await .expect("valid dealing"); let duplicate_ack = player .handle( &mut store, Epoch::zero(), dealer_pk.clone(), duplicate_public.clone(), duplicate_private.clone(), ) .await .expect("cached dealing"); assert_eq!(duplicate_ack, ack); // Sender membership and signature failures remain visible to the // authenticated transport policy without consuming pending state. let stranger = PrivateKey::from_seed(u64::MAX).public_key(); assert!(matches!( dealer .handle(&mut store, Epoch::zero(), stranger, ack.clone()) .await, Err(DkgAckError::UnexpectedPlayer) )); assert!(matches!( dealer .handle( &mut store, Epoch::zero(), signers[2].public_key(), ack.clone(), ) .await, Err(DkgAckError::InvalidAck) )); // Acknowledgements distinguish new and duplicate outcomes. assert!(matches!( dealer .handle(&mut store, Epoch::zero(), player_pk.clone(), ack.clone()) .await, Ok(AckOutcome::Recorded) )); assert!(matches!( dealer .handle(&mut store, Epoch::zero(), player_pk.clone(), ack.clone()) .await, Ok(AckOutcome::Duplicate) )); assert!(dealer.finalize::()); let signed = dealer.finalized().expect("signed log"); let (dealer, log) = signed.check(&info).expect("valid log"); store.append_log(Epoch::zero(), dealer, log).await; drop(store); let mut store = init_store(context.child("restart"), "replay", secret_store).await; // The current epoch is not persisted; the setup state re-derives it from // finalized blocks, so a restarted store has no current epoch on its own. assert!(store.current().is_none()); // The public recovery journal is replayed. assert_eq!(store.dealings(Epoch::zero()).len(), 1); assert_eq!(store.acks(Epoch::zero()).len(), 1); assert_eq!(store.logs(Epoch::zero()).len(), 1); let mut replayed_player = store .create_player::(Epoch::zero(), signers[1].clone(), info) .expect("player"); assert_eq!(replayed_player.acks.len(), 1); let replayed_ack = replayed_player .handle( &mut store, Epoch::zero(), dealer_pk, duplicate_public, duplicate_private, ) .await .expect("replayed dealing"); assert_eq!(replayed_ack, ack); }); } #[test] fn resume_with_missing_dealing_degrades_to_observer() { let executor = deterministic::Runner::default(); executor.start(|context| async move { let secret_store = MemorySecretStore::default(); let mut store = init_store( context.child("store"), "missing-dealing", secret_store.clone(), ) .await; let signers = signers(); let players = players(&signers); let info = Info::new::( b"_COMMONWARE_GLUE_DKG_RESHARE_STORE_TEST", 0, None, Mode::NonZeroCounter, Reveal::V1, players.clone(), players.clone(), ) .expect("valid info"); store .commit_epoch(epoch_info(Epoch::zero(), output(1)), summary(1), None) .await; let dealer_pk = signers[0].public_key(); let mut dealer = store .create_dealer::( Epoch::zero(), signers[0].clone(), info.clone(), None, summary(2), ) .expect("dealer"); // Collect enough acknowledgements (quorum) that the dealer log records // an Ack for each acking player rather than degrading to reveals. With // n=4, f=1 the log tolerates at most one reveal, so three players ack. for idx in [1usize, 2, 3] { let player_pk = signers[idx].public_key(); let mut player = Player { player: CryptoPlayer::new(info.clone(), signers[idx].clone()).expect("player"), acks: BTreeMap::new(), }; let (_, public, private) = dealer .shares_to_distribute() .find(|(recipient, _, _)| *recipient == player_pk) .expect("dealing for player"); let ack = player .handle( &mut store, Epoch::zero(), dealer_pk.clone(), public, private, ) .await .expect("valid dealing"); assert!(matches!( dealer .handle(&mut store, Epoch::zero(), player_pk, ack) .await, Ok(AckOutcome::Recorded) )); } assert!(dealer.finalize::()); let signed = dealer.finalized().expect("signed log"); let (dealer, log) = signed.check(&info).expect("valid log"); store.append_log(Epoch::zero(), dealer, log).await; drop(store); // Restart with an empty secret store: the finalized dealer log (which // records this player's acknowledgement) replays from public storage, // but the matching private dealing is gone. Resuming as that player // must degrade to an observer rather than panic. let empty_secret_store = MemorySecretStore::default(); let restarted = init_store( context.child("restart"), "missing-dealing", empty_secret_store, ) .await; assert!(restarted.dealings(Epoch::zero()).is_empty()); assert_eq!(restarted.logs(Epoch::zero()).len(), 1); let player = restarted.create_player::( Epoch::zero(), signers[1].clone(), info, ); assert!( player.is_none(), "a missing private dealing must degrade to observer, not panic" ); }); } #[test] fn protocol_storage_does_not_restore_private_dealings() { let executor = deterministic::Runner::default(); executor.start(|context| async move { let secret_store = MemorySecretStore::default(); let mut store = init_store( context.child("store"), "private-dealing-boundary", secret_store, ) .await; let signers = signers(); let players = players(&signers); let info = Info::new::( b"_COMMONWARE_GLUE_DKG_RESHARE_STORE_TEST", 0, None, Mode::NonZeroCounter, Reveal::V1, players.clone(), players, ) .expect("valid info"); store .commit_epoch(epoch_info(Epoch::zero(), output(1)), summary(1), None) .await; let dealer_pk = signers[0].public_key(); let player_pk = signers[1].public_key(); let dealer = store .create_dealer::( Epoch::zero(), signers[0].clone(), info.clone(), None, summary(2), ) .expect("dealer"); let mut player = store .create_player::(Epoch::zero(), signers[1].clone(), info.clone()) .expect("player"); let (_, public, private) = dealer .shares_to_distribute() .find(|(recipient, _, _)| *recipient == player_pk) .expect("dealing for player"); player .handle(&mut store, Epoch::zero(), dealer_pk, public, private) .await .expect("valid dealing"); assert_eq!(store.dealings(Epoch::zero()).len(), 1); drop(store); let empty_secret_store = MemorySecretStore::default(); let restarted = init_store( context.child("restart"), "private-dealing-boundary", empty_secret_store, ) .await; assert!(restarted.dealings(Epoch::zero()).is_empty()); let replayed_player = restarted .create_player::(Epoch::zero(), signers[1].clone(), info) .expect("player"); assert!(replayed_player.acks.is_empty()); }); } #[test] fn protocol_storage_does_not_restore_secret_shares() { let executor = deterministic::Runner::default(); executor.start(|context| async move { let signers = signers(); let (_, shares) = deal::(TestRng::new(9), Mode::NonZeroCounter, players(&signers)) .expect("trusted deal"); let share = shares .get_value(&signers[0].public_key()) .expect("share") .clone(); let secret_store = MemorySecretStore::default(); let mut store = init_store(context.child("store"), "secret-boundary", secret_store).await; store .commit_epoch( epoch_info(Epoch::zero(), output(1)), summary(1), Some(share), ) .await; assert!(store.share(Epoch::zero()).await.is_some()); drop(store); let empty_secret_store = MemorySecretStore::default(); let mut restarted = init_store( context.child("restart"), "secret-boundary", empty_secret_store, ) .await; // The current epoch is never persisted; setup re-derives it from // finalized boundary blocks. The share lived only in the secret store, // which is now empty, so neither a current epoch nor the share survives. assert!(restarted.current().is_none()); assert!(restarted.share(Epoch::zero()).await.is_none()); }); } #[test] fn prune_removes_old_protocol_state_and_secret_shares() { let executor = deterministic::Runner::default(); executor.start(|context| async move { let secret_store = MemorySecretStore::default(); let mut store = init_store(context.child("store"), "prune", secret_store.clone()).await; let signers = signers(); let (next_output, shares) = deal::(TestRng::new(10), Mode::NonZeroCounter, players(&signers)) .expect("trusted deal"); let share = shares .get_value(&signers[0].public_key()) .expect("share") .clone(); store .commit_epoch(epoch_info(Epoch::zero(), output(1)), summary(1), None) .await; store .commit_epoch( epoch_info(Epoch::new(1), next_output), summary(2), Some(share), ) .await; store.prune(Epoch::new(1)).await; drop(store); let store = init_store(context.child("restart"), "prune", secret_store.clone()).await; // The current epoch is not persisted (the setup state re-derives it), and // the pruned epoch's public journal and secret material stay pruned across // the restart. assert!(store.current().is_none()); assert!(!store.epochs.contains_key(&Epoch::zero())); assert_eq!(secret_store.prunes(), vec![Epoch::new(1)]); assert!(!secret_store.has_share(Epoch::zero())); }); } } #[cfg(all(test, feature = "arbitrary"))] mod conformance { use super::*; use commonware_codec::conformance::CodecConformance; use commonware_cryptography::{bls12381::primitives::variant::MinSig, ed25519}; commonware_conformance::conformance_tests! { CodecConformance>, } }