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); const PAGE_CACHE_CAPACITY: NonZeroUsize = NZUsize!(1 << 13); const WRITE_BUFFER: NonZeroUsize = NZUsize!(1 << 12); const READ_BUFFER: NonZeroUsize = NZUsize!(1 << 20);
enum Event<V: Variant, P: PublicKey> {
Dealing(P, DealerPubMsg<V>),
Ack(P, PlayerAck<P>),
Log(P, DealerLog<V, P>),
}
impl<V: Variant, P: PublicKey> EncodeSize for Event<V, P> {
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<V: Variant, P: PublicKey> Write for Event<V, P> {
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<V: Variant, P: PublicKey> Read for Event<V, P> {
type Cfg = NonZeroU32;
fn read_cfg(reader: &mut impl Buf, cfg: &Self::Cfg) -> Result<Self, CodecError> {
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<V: Variant, P: PublicKey> arbitrary::Arbitrary<'_> for Event<V, P>
where
P: for<'a> arbitrary::Arbitrary<'a>,
DealerPubMsg<V>: for<'a> arbitrary::Arbitrary<'a>,
DealerLog<V, P>: for<'a> arbitrary::Arbitrary<'a>,
PlayerAck<P>: for<'a> arbitrary::Arbitrary<'a>,
{
fn arbitrary(u: &mut arbitrary::Unstructured<'_>) -> arbitrary::Result<Self> {
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<V: Variant, P: PublicKey> {
dealings: BTreeMap<P, (DealerPubMsg<V>, DealerPrivMsg)>,
acks: BTreeMap<P, PlayerAck<P>>,
logs: BTreeMap<P, DealerLog<V, P>>,
}
impl<V: Variant, P: PublicKey> Default for EpochCache<V, P> {
fn default() -> Self {
Self {
dealings: BTreeMap::new(),
acks: BTreeMap::new(),
logs: BTreeMap::new(),
}
}
}
pub struct Store<E, SS, V, P, D = Unit>
where
E: BufferPooler + Clock + RuntimeStorage + Metrics,
SS: SecretStore,
V: Variant,
P: PublicKey,
D: Directory<P>,
{
secret_store: SS,
events: Option<Journal<E, Event<V, P>>>,
current: Option<EpochInfo<V, P, D>>,
epochs: BTreeMap<Epoch, EpochCache<V, P>>,
}
impl<E, SS, V, P, D> Store<E, SS, V, P, D>
where
E: BufferPooler + Clock + RuntimeStorage + Metrics,
SS: SecretStore,
V: Variant,
P: PublicKey,
D: Directory<P>,
{
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");
let current = None;
let mut epochs = BTreeMap::<Epoch, EpochCache<V, P>>::new();
let events = {
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,
}
}
pub fn current(&self) -> Option<EpochInfo<V, P, D>> {
self.current.clone()
}
pub async fn share(&mut self, epoch: Epoch) -> Option<group::Share> {
self.secret_store.get_share(epoch).await
}
pub async fn seed(&mut self, epoch: Epoch) -> Option<Summary> {
self.secret_store.get_seed(epoch).await
}
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))
}
pub(crate) async fn put_seed(&mut self, epoch: Epoch, rng_seed: Summary) {
self.secret_store.put_seed(epoch, rng_seed).await;
}
pub async fn commit_epoch(
&mut self,
info: EpochInfo<V, P, D>,
rng_seed: Summary,
share: Option<group::Share>,
) {
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);
}
pub async fn prune(&mut self, min: Epoch) {
self.epochs.retain(|epoch, _| *epoch >= min);
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<V, P> {
self.epochs.entry(epoch).or_default()
}
pub fn logs(&self, epoch: Epoch) -> BTreeMap<P, DealerLog<V, P>> {
self.epochs
.get(&epoch)
.map(|cache| cache.logs.clone())
.unwrap_or_default()
}
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<V>, 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<P>)> {
self.epochs
.get(&epoch)
.map(|cache| {
cache
.acks
.iter()
.map(|(player, ack)| (player.clone(), ack.clone()))
.collect()
})
.unwrap_or_default()
}
async fn append_dealing(
&mut self,
epoch: Epoch,
dealer: P,
public: DealerPubMsg<V>,
private: DealerPrivMsg,
) -> bool {
if self
.epochs
.get(&epoch)
.is_some_and(|cache| cache.dealings.contains_key(&dealer))
{
return false;
}
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<P>) -> 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
}
pub async fn append_log(&mut self, epoch: Epoch, dealer: P, log: DealerLog<V, P>) -> 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
}
pub fn create_dealer<C, M>(
&self,
epoch: Epoch,
signer: C,
info: Info<V, P>,
share: Option<group::Share>,
rng_seed: Summary,
) -> Option<Dealer<V, C>>
where
C: Signer<PublicKey = P>,
M: Faults,
{
if self.has_log(epoch, &signer.public_key()) {
return None;
}
let (mut dealer, public, private) = CryptoDealer::start::<M>(
Transcript::resume(rng_seed, Version::V1).noise(b"dealer-rng"),
info,
signer,
share,
)
.expect("failed to create reshare dealer");
let mut unsent: BTreeMap<P, DealerPrivMsg> = 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,
})
}
pub fn create_player<C, M>(
&self,
epoch: Epoch,
signer: C,
info: Info<V, P>,
) -> Option<Player<V, C>>
where
C: Signer<PublicKey = P>,
M: Faults,
{
self.resume_player::<C, M>(epoch, signer, info, &self.logs(epoch))
}
pub fn create_player_with_logs<C, M>(
&self,
epoch: Epoch,
signer: C,
info: Info<V, P>,
logs: &BTreeMap<P, DealerLog<V, P>>,
) -> Option<Player<V, C>>
where
C: Signer<PublicKey = P>,
M: Faults,
{
self.resume_player::<C, M>(epoch, signer, info, logs)
}
fn resume_player<C, M>(
&self,
epoch: Epoch,
signer: C,
info: Info<V, P>,
logs: &BTreeMap<P, DealerLog<V, P>>,
) -> Option<Player<V, C>>
where
C: Signer<PublicKey = P>,
M: Faults,
{
match CryptoPlayer::resume::<M>(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:?}"),
}
}
}
async fn append_synced<E, V, P>(
events: Journal<E, Event<V, P>>,
epoch: Epoch,
event: &Event<V, P>,
) -> Result<Journal<E, Event<V, P>>, 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
}
pub struct Dealer<V: Variant, C: Signer> {
dealer: Option<CryptoDealer<V, C>>,
public: DealerPubMsg<V>,
unsent: BTreeMap<C::PublicKey, DealerPrivMsg>,
finalized: Option<SignedDealerLog<V, C>>,
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub(crate) enum AckOutcome {
Recorded,
Duplicate,
}
impl<V: Variant, C: Signer> Dealer<V, C> {
pub fn has_acknowledgement_quorum<M: Faults>(&self, players: usize) -> bool {
self.unsent.len() <= M::max_faults(players) as usize
}
pub async fn handle<E, SS, D>(
&mut self,
store: &mut Store<E, SS, V, C::PublicKey, D>,
epoch: Epoch,
player: C::PublicKey,
ack: PlayerAck<C::PublicKey>,
) -> Result<AckOutcome, DkgAckError>
where
E: BufferPooler + Clock + RuntimeStorage + Metrics,
SS: SecretStore,
D: Directory<C::PublicKey>,
{
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)
}
pub fn finalize<M: Faults>(&mut self) -> bool {
if self.finalized.is_some() {
return false;
}
let Some(dealer) = self.dealer.take() else {
return false;
};
self.finalized = Some(dealer.finalize::<M>());
true
}
pub fn finalized(&self) -> Option<SignedDealerLog<V, C>> {
self.finalized.clone()
}
pub fn clear_finalized(&mut self) {
self.finalized = None;
}
pub fn shares_to_distribute(
&self,
) -> impl Iterator<Item = (C::PublicKey, DealerPubMsg<V>, DealerPrivMsg)> + '_ {
self.unsent
.iter()
.map(|(player, private)| (player.clone(), self.public.clone(), private.clone()))
}
}
pub struct Player<V: Variant, C: Signer> {
player: CryptoPlayer<V, C>,
acks: BTreeMap<C::PublicKey, PlayerAck<C::PublicKey>>,
}
impl<V: Variant, C: Signer> Player<V, C> {
pub async fn handle<E, SS, D>(
&mut self,
store: &mut Store<E, SS, V, C::PublicKey, D>,
epoch: Epoch,
dealer: C::PublicKey,
public: DealerPubMsg<V>,
private: DealerPrivMsg,
) -> Result<PlayerAck<C::PublicKey>, DkgDealerMessageError>
where
E: BufferPooler + Clock + RuntimeStorage + Metrics,
SS: SecretStore,
D: Directory<C::PublicKey>,
{
if let Some(ack) = self.acks.get(&dealer) {
return Ok(ack.clone());
}
let ack = self
.player
.dealer_message::<N3f1>(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)
}
#[allow(clippy::type_complexity)]
pub fn finalize<M, B>(
self,
rng: &mut impl CryptoRng,
logs: Logs<V, C::PublicKey, M>,
strategy: &impl Strategy,
) -> Result<(Output<V, C::PublicKey>, group::Share), DkgFinalizeError<C::PublicKey>>
where
M: Faults,
B: BatchVerifier<PublicKey = C::PublicKey>,
{
self.player.finalize::<M, B>(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<E> = Store<E, MemorySecretStore, MinPk, PublicKey>;
fn summary(seed: u8) -> Summary {
let bytes = [seed; Summary::SIZE];
Summary::read(&mut bytes.as_ref()).expect("valid summary")
}
fn output(seed: u64) -> Output<MinPk, PublicKey> {
let (output, _) = deal::<MinPk, _, N3f1>(
TestRng::new(seed),
Mode::NonZeroCounter,
players(&signers()),
)
.expect("trusted deal");
output
}
fn epoch_info(epoch: Epoch, output: Output<MinPk, PublicKey>) -> EpochInfo<MinPk, PublicKey> {
EpochInfo {
outcome: EpochOutcome::Success,
epoch,
output,
players: Set::default(),
next_players: Set::default(),
directory: Unit,
}
}
fn signers() -> Vec<PrivateKey> {
(0..4).map(PrivateKey::from_seed).collect()
}
fn players(signers: &[PrivateKey]) -> Set<PublicKey> {
Set::from_iter_dedup(signers.iter().map(Signer::public_key))
}
async fn init_store<E>(
context: E,
partition: &str,
secret_store: MemorySecretStore,
) -> TestStore<E>
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::<N3f1>(
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::<PrivateKey, N3f1>(
Epoch::zero(),
signers[0].clone(),
info.clone(),
None,
summary(2),
)
.expect("dealer");
let mut player = store
.create_player::<PrivateKey, N3f1>(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);
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)
));
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::<N3f1>());
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;
assert!(store.current().is_none());
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::<PrivateKey, N3f1>(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::<N3f1>(
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::<PrivateKey, N3f1>(
Epoch::zero(),
signers[0].clone(),
info.clone(),
None,
summary(2),
)
.expect("dealer");
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::<N3f1>());
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 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::<PrivateKey, N3f1>(
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::<N3f1>(
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::<PrivateKey, N3f1>(
Epoch::zero(),
signers[0].clone(),
info.clone(),
None,
summary(2),
)
.expect("dealer");
let mut player = store
.create_player::<PrivateKey, N3f1>(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::<PrivateKey, N3f1>(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::<MinPk, _, N3f1>(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;
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::<MinPk, _, N3f1>(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;
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<Event<MinSig, ed25519::PublicKey>>,
}
}