mod actor;
pub use actor::{Config, Probe};
mod mailbox;
pub use mailbox::Mailbox;
mod wire;
#[cfg(test)]
mod test {
use super::{wire, Config, Mailbox, Probe};
use bytes::{Buf, BufMut};
use commonware_actor::Feedback;
use commonware_codec::{Encode, EncodeSize, Error as CodecError, Read, ReadExt as _, Write};
use commonware_consensus::{
marshal::{
self,
core::{Actor as MarshalActor, Mailbox as MarshalMailbox},
resolver::p2p as marshal_resolver,
standard::Standard,
Start, Update,
},
simplex::{
mocks::scheme::{self as scheme_mocks, Scheme as MockScheme},
types::{Activity, Context as SimplexContext, Finalization, Finalize, Proposal},
},
types::{Epoch, FixedEpocher, Height, Round, View, ViewDelta},
Block as ConsensusBlock, CertifiableBlock, Heightable, Reporter,
};
use commonware_cryptography::{
certificate::{
self, mocks::Fixture, Attestation, ConstantProvider, Provider, Scoped, Verification,
Verifier,
},
ed25519,
sha256::Digest as Sha256Digest,
Digest as _, Digestible, Signer as _,
};
use commonware_macros::test_collect_traces;
use commonware_p2p::{
simulated::{Config as SimConfig, Link, Network, Oracle, Sender},
Recipients, Sender as _,
};
use commonware_parallel::{Sequential, Strategy as ParallelStrategy};
use commonware_runtime::{
buffer::paged::CacheRef, deterministic, telemetry::traces::collector::TraceStorage, Clock,
Handle, Metrics, Quota, Runner as _, Supervisor,
};
use commonware_storage::archive::immutable;
use commonware_utils::{
channel::oneshot, sync::Mutex, test_rng, Acknowledgement, NZDuration, NZUsize,
NonZeroDuration, NZU16, NZU64,
};
use std::{
collections::BTreeMap,
num::{NonZeroU32, NonZeroU64},
sync::Arc,
time::Duration,
};
const NAMESPACE: &[u8] = b"_COMMONWARE_GLUE_PROBE_TEST";
const EPOCH_LENGTH: NonZeroU64 = NZU64!(u64::MAX);
const TEST_QUOTA: Quota = Quota::per_second(NonZeroU32::MAX);
const LINK: Link = Link {
latency: Duration::from_millis(10),
jitter: Duration::from_millis(1),
success_rate: 1.0,
};
const PROBE_CHANNEL: u64 = 0;
const BACKFILL_CHANNEL: u64 = 1;
type Scheme = MockScheme<ed25519::PublicKey>;
type Variant = Standard<Block>;
#[derive(Clone, Debug, PartialEq, Eq)]
struct Block {
context: SimplexContext<Sha256Digest, ed25519::PublicKey>,
height: Height,
digest: Sha256Digest,
}
impl Block {
fn new(height: u64, digest_byte: u8) -> Self {
Self {
context: SimplexContext {
round: Round::new(Epoch::zero(), View::new(height)),
leader: ed25519::PrivateKey::from_seed(0).public_key(),
parent: (View::zero(), Sha256Digest::EMPTY),
},
height: Height::new(height),
digest: Sha256Digest::from([digest_byte; 32]),
}
}
}
impl Write for Block {
fn write(&self, buf: &mut impl BufMut) {
self.context.write(buf);
buf.put_u64(self.height.get());
buf.put_slice(self.digest.as_ref());
}
}
impl EncodeSize for Block {
fn encode_size(&self) -> usize {
self.context.encode_size() + 8 + 32
}
}
impl Read for Block {
type Cfg = ();
fn read_cfg(buf: &mut impl Buf, _: &Self::Cfg) -> Result<Self, CodecError> {
let context = SimplexContext::read(buf)?;
let height = Height::new(buf.get_u64());
let mut digest = [0u8; 32];
buf.copy_to_slice(&mut digest);
Ok(Self {
context,
height,
digest: Sha256Digest::from(digest),
})
}
}
impl Digestible for Block {
type Digest = Sha256Digest;
fn digest(&self) -> Self::Digest {
self.digest
}
}
impl Heightable for Block {
fn height(&self) -> Height {
self.height
}
}
impl ConsensusBlock for Block {
fn parent(&self) -> Self::Digest {
Sha256Digest::EMPTY
}
}
impl CertifiableBlock for Block {
type Context = SimplexContext<Sha256Digest, ed25519::PublicKey>;
fn context(&self) -> Self::Context {
self.context.clone()
}
}
#[derive(Clone)]
struct NoopReporter;
impl Reporter for NoopReporter {
type Activity = Update<Block>;
fn report(&mut self, activity: Self::Activity) -> Feedback {
if let Update::Block(_, ack) = activity {
ack.acknowledge();
}
Feedback::Ok
}
}
#[derive(Clone, Default)]
struct EpochProvider(Arc<Mutex<BTreeMap<Epoch, Arc<Scheme>>>>);
impl EpochProvider {
fn insert(&self, epoch: Epoch, scheme: Scheme) {
self.0.lock().insert(epoch, Arc::new(scheme));
}
fn forget(&self, epoch: Epoch) {
self.0.lock().remove(&epoch);
}
}
impl Provider for EpochProvider {
type Scope = Epoch;
type Scheme = Scheme;
fn scoped(&self, scope: Epoch) -> Option<Scoped<Scheme>> {
self.0.lock().get(&scope).cloned().map(Scoped::scheme)
}
}
#[derive(Clone, Debug)]
struct MaybeEnumerableScheme {
inner: Scheme,
enumerable: bool,
}
impl MaybeEnumerableScheme {
const fn new(inner: Scheme, enumerable: bool) -> Self {
Self { inner, enumerable }
}
fn wrap_attestation(attestation: Attestation<Scheme>) -> Attestation<Self> {
Attestation {
signer: attestation.signer,
signature: attestation.signature,
}
}
fn unwrap_attestation(attestation: Attestation<Self>) -> Attestation<Scheme> {
Attestation {
signer: attestation.signer,
signature: attestation.signature,
}
}
}
impl Verifier for MaybeEnumerableScheme {
type Subject<'a, D: commonware_cryptography::Digest> =
commonware_consensus::simplex::types::Subject<'a, D>;
type PublicKey = ed25519::PublicKey;
type Certificate = <Scheme as Verifier>::Certificate;
fn verify_certificate<R, D, M>(
&self,
rng: &mut R,
subject: Self::Subject<'_, D>,
certificate: &Self::Certificate,
strategy: &impl ParallelStrategy,
) -> bool
where
R: rand_core::CryptoRng,
D: commonware_cryptography::Digest,
M: commonware_utils::Faults,
{
self.inner
.verify_certificate::<_, D, M>(rng, subject, certificate, strategy)
}
fn verify_certificates<'a, R, D, I, M>(
&self,
rng: &mut R,
certificates: I,
strategy: &impl ParallelStrategy,
) -> bool
where
R: rand_core::CryptoRng,
D: commonware_cryptography::Digest,
I: Iterator<Item = (Self::Subject<'a, D>, &'a Self::Certificate)>,
M: commonware_utils::Faults,
{
self.inner
.verify_certificates::<_, D, _, M>(rng, certificates, strategy)
}
fn is_batchable() -> bool {
Scheme::is_batchable()
}
fn certificate_codec_config(&self) -> <Self::Certificate as commonware_codec::Read>::Cfg {
self.inner.certificate_codec_config()
}
fn certificate_codec_config_unbounded() -> <Self::Certificate as commonware_codec::Read>::Cfg
{
Scheme::certificate_codec_config_unbounded()
}
}
impl certificate::Scheme for MaybeEnumerableScheme {
type Signature = <Scheme as certificate::Scheme>::Signature;
fn me(&self) -> Option<commonware_utils::Participant> {
self.inner.me()
}
fn participants(&self) -> &commonware_utils::ordered::Set<Self::PublicKey> {
assert!(
self.enumerable,
"verify-only scheme has no participant metadata"
);
self.inner.participants()
}
fn sign<D: commonware_cryptography::Digest>(
&self,
subject: Self::Subject<'_, D>,
) -> Option<Attestation<Self>> {
self.inner.sign(subject).map(Self::wrap_attestation)
}
fn verify_attestation<R, D>(
&self,
rng: &mut R,
subject: Self::Subject<'_, D>,
attestation: &Attestation<Self>,
strategy: &impl ParallelStrategy,
) -> bool
where
R: rand_core::CryptoRng,
D: commonware_cryptography::Digest,
{
let attestation = Attestation::<Scheme> {
signer: attestation.signer,
signature: attestation.signature.clone(),
};
self.inner
.verify_attestation(rng, subject, &attestation, strategy)
}
fn verify_attestations<R, D, I>(
&self,
rng: &mut R,
subject: Self::Subject<'_, D>,
attestations: I,
strategy: &impl ParallelStrategy,
) -> Verification<Self>
where
R: rand_core::CryptoRng,
D: commonware_cryptography::Digest,
I: IntoIterator<Item = Attestation<Self>>,
I::IntoIter: Send,
{
let verification = self.inner.verify_attestations(
rng,
subject,
attestations.into_iter().map(Self::unwrap_attestation),
strategy,
);
Verification::new(
verification
.verified
.into_iter()
.map(Self::wrap_attestation)
.collect(),
verification.invalid,
)
}
fn assemble<I, M>(
&self,
attestations: I,
strategy: &impl ParallelStrategy,
) -> Option<Self::Certificate>
where
I: IntoIterator<Item = Attestation<Self>>,
I::IntoIter: Send,
M: commonware_utils::Faults,
{
self.inner.assemble::<_, M>(
attestations.into_iter().map(Self::unwrap_attestation),
strategy,
)
}
fn is_attributable() -> bool {
Scheme::is_attributable()
}
}
#[derive(Clone)]
struct ParticipantlessAllProvider {
verifier: Arc<MaybeEnumerableScheme>,
scheme: Arc<MaybeEnumerableScheme>,
}
impl Provider for ParticipantlessAllProvider {
type Scope = Epoch;
type Scheme = MaybeEnumerableScheme;
fn scoped(&self, _: Epoch) -> Option<Scoped<MaybeEnumerableScheme>> {
Some(Scoped::verifier(self.verifier.clone()))
}
fn scheme(&self, _: Epoch) -> Option<Arc<MaybeEnumerableScheme>> {
Some(self.scheme.clone())
}
}
struct Node {
probe: Mailbox<Scheme, Variant>,
marshal: MarshalMailbox<Scheme, Variant>,
probe_sender: Sender<ed25519::PublicKey, deterministic::Context>,
start: Option<Box<dyn FnOnce() -> Handle<()>>>,
_handles: Vec<Handle<()>>,
}
struct Harness {
participants: Vec<ed25519::PublicKey>,
schemes: Vec<Scheme>,
nodes: Vec<Node>,
oracle: Oracle<ed25519::PublicKey, deterministic::Context>,
_network: Handle<()>,
}
impl Harness {
async fn setup(
context: &deterministic::Context,
n: u32,
retry_timeout: NonZeroDuration,
) -> Self {
Self::setup_with(context, n, retry_timeout, Epoch::zero(), |scheme| {
ConstantProvider::new(scheme.clone())
})
.await
}
async fn setup_with<D, F>(
context: &deterministic::Context,
n: u32,
retry_timeout: NonZeroDuration,
minimum_epoch: Epoch,
make_provider: F,
) -> Self
where
D: Provider<Scope = Epoch, Scheme = Scheme>,
F: Fn(&Scheme) -> D,
{
let mut rng = test_rng();
let Fixture {
participants,
schemes,
..
} = scheme_mocks::fixture(&mut rng, NAMESPACE, n);
let (network, oracle) = Network::new_with_peers(
context.child("network"),
SimConfig {
max_size: 1024 * 1024,
disconnect_on_block: true,
tracked_peer_sets: NZUsize!(1),
},
participants.clone(),
)
.await;
let network = network.start();
for a in &participants {
for b in &participants {
if a != b {
oracle
.add_link(a.clone(), b.clone(), LINK)
.await
.expect("failed to add link");
}
}
}
let genesis = Block::new(0, 0);
let mut nodes = Vec::with_capacity(n as usize);
for (index, public_key) in participants.iter().enumerate() {
let scheme = schemes[index].clone();
let node_ctx = context.child("node").with_attribute("index", index);
let partition_prefix = format!("node-{index}");
let page_cache = CacheRef::from_pooler(&node_ctx, NZU16!(1024), NZUsize!(10));
let control = oracle.control(public_key.clone());
let backfill = control
.register(BACKFILL_CHANNEL, TEST_QUOTA)
.await
.expect("failed to register backfill channel");
let resolver = marshal_resolver::init(
node_ctx.child("marshal_resolver"),
marshal_resolver::Config {
public_key: public_key.clone(),
peer_provider: oracle.manager(),
blocker: oracle.control(public_key.clone()),
mailbox_size: NZUsize!(100),
initial: Duration::from_secs(1),
timeout: Duration::from_secs(2),
fetch_retry_timeout: Duration::from_millis(100),
priority_requests: false,
priority_responses: false,
},
backfill,
);
let finalizations_by_height = immutable::Archive::init(
node_ctx.child("finalizations_by_height"),
archive_config(&partition_prefix, "finalizations", page_cache.clone()),
)
.await
.expect("failed to init finalizations archive");
let finalized_blocks = immutable::Archive::init(
node_ctx.child("finalized_blocks"),
archive_config(&partition_prefix, "blocks", page_cache.clone()),
)
.await
.expect("failed to init blocks archive");
let marshal_config = marshal::Config {
provider: ConstantProvider::new(scheme.clone()),
epocher: FixedEpocher::new(EPOCH_LENGTH),
start: Start::Genesis(genesis.clone()),
partition_prefix: partition_prefix.clone(),
mailbox_size: NZUsize!(100),
view_retention_timeout: ViewDelta::new(10),
prunable_items_per_section: NZU64!(10),
page_cache,
replay_buffer: NZUsize!(2048),
key_write_buffer: NZUsize!(2048),
value_write_buffer: NZUsize!(2048),
block_codec_config: (),
max_repair: NZUsize!(10),
max_pending_acks: NZUsize!(1),
strategy: Sequential,
};
let (marshal_actor, marshal_mailbox, _) =
MarshalActor::<_, Variant, _, _, _, _, _>::init(
node_ctx.child("marshal"),
finalizations_by_height,
finalized_blocks,
marshal_config,
)
.await;
let marshal_handle = marshal_actor.start_unbuffered(NoopReporter, resolver);
let probe_network = control
.register(PROBE_CHANNEL, TEST_QUOTA)
.await
.expect("failed to register probe channel");
let probe_sender = probe_network.0.clone();
let (probe, probe_mailbox) = Probe::new(Config {
context: node_ctx.child("probe"),
provider: make_provider(&scheme),
strategy: Sequential,
capacity: NZUsize!(100),
blocker: oracle.control(public_key.clone()),
minimum_epoch,
retry_timeout,
});
if index != 0 {
probe_mailbox.attach(marshal_mailbox.clone());
}
let start: Box<dyn FnOnce() -> Handle<()>> =
Box::new(move || probe.start(probe_network));
nodes.push(Node {
probe: probe_mailbox,
marshal: marshal_mailbox,
probe_sender,
start: Some(start),
_handles: vec![marshal_handle],
});
}
Self {
participants,
schemes,
nodes,
oracle,
_network: network,
}
}
fn start_probes(&mut self) {
for node in &mut self.nodes {
if let Some(start) = node.start.take() {
node._handles.push(start());
}
}
}
fn finalization(
&self,
height: u64,
digest_byte: u8,
) -> (Block, Finalization<Scheme, Sha256Digest>) {
build_finalization(&self.schemes, height, digest_byte)
}
async fn inject(
&self,
index: usize,
block: Block,
finalization: Finalization<Scheme, Sha256Digest>,
) {
let mut marshal = self.nodes[index].marshal.clone();
let round = finalization.proposal.round;
assert!(marshal.verified(round, block).await);
let _ = marshal.report(Activity::Finalization(finalization));
}
fn send_raw(&self, from: usize, to: usize, bytes: Vec<u8>) {
let mut sender = self.nodes[from].probe_sender.clone();
sender.send(Recipients::One(self.participants[to].clone()), bytes, false);
}
}
fn build_finalization(
schemes: &[Scheme],
height: u64,
digest_byte: u8,
) -> (Block, Finalization<Scheme, Sha256Digest>) {
build_finalization_at(schemes, Epoch::zero(), height, digest_byte)
}
fn build_finalization_at<S>(
schemes: &[S],
epoch: Epoch,
height: u64,
digest_byte: u8,
) -> (Block, Finalization<S, Sha256Digest>)
where
S: commonware_consensus::simplex::scheme::Scheme<Sha256Digest>,
{
let block = Block::new(height, digest_byte);
let round = Round::new(epoch, View::new(height));
let proposal = Proposal {
round,
parent: View::new(height.saturating_sub(1)),
payload: block.digest(),
};
let finalizes: Vec<_> = schemes
.iter()
.map(|scheme| Finalize::sign(scheme, proposal.clone()).expect("sign finalize"))
.collect();
let finalization =
Finalization::from_finalizes(&schemes[0], &finalizes, &Sequential).expect("recover");
(block, finalization)
}
fn epoch_provider(entries: impl IntoIterator<Item = (Epoch, Scheme)>) -> EpochProvider {
let provider = EpochProvider::default();
for (epoch, scheme) in entries {
provider.insert(epoch, scheme);
}
provider
}
fn finalization_bytes<S>(finalization: Finalization<S, Sha256Digest>) -> Vec<u8>
where
S: commonware_consensus::simplex::scheme::Scheme<Sha256Digest>,
{
wire::Message::<S, Variant>::Response(finalization)
.encode()
.to_vec()
}
#[cfg(feature = "arbitrary")]
mod conformance {
use super::{wire, Scheme, Variant};
use commonware_codec::conformance::CodecConformance;
commonware_conformance::conformance_tests! {
CodecConformance<wire::Tag>,
CodecConformance<wire::Message<Scheme, Variant>>,
}
}
fn archive_config(prefix: &str, name: &str, page_cache: CacheRef) -> immutable::Config<()> {
immutable::Config {
metadata_partition: format!("{prefix}-{name}-metadata"),
freezer_table_partition: format!("{prefix}-{name}-freezer-table"),
freezer_table_initial_size: 64,
freezer_table_resize_frequency: 10,
freezer_table_resize_chunk_size: 10,
freezer_key_partition: format!("{prefix}-{name}-freezer-key"),
freezer_key_page_cache: page_cache,
freezer_value_partition: format!("{prefix}-{name}-freezer-value"),
freezer_value_target_size: 1024,
freezer_value_compression: None,
ordinal_partition: format!("{prefix}-{name}-ordinal"),
items_per_section: NZU64!(10),
codec_config: (),
replay_buffer: NZUsize!(2048),
freezer_key_write_buffer: NZUsize!(2048),
freezer_value_write_buffer: NZUsize!(2048),
ordinal_write_buffer: NZUsize!(2048),
}
}
#[test]
fn test_resolves_floor_from_sample_peers() {
let runner = deterministic::Runner::timed(Duration::from_secs(30));
runner.start(|context| async move {
let mut harness =
Harness::setup(&context, 4, NZDuration!(Duration::from_millis(500))).await;
assert_eq!(harness.participants.len(), 4);
let (block, finalization) = harness.finalization(1, 1);
for index in [1, 2] {
harness
.inject(index, block.clone(), finalization.clone())
.await;
}
harness.start_probes();
let floor = harness.nodes[0]
.probe
.subscribe()
.await
.expect("floor resolved");
assert_eq!(floor, finalization);
});
}
#[test]
fn test_resolves_highest_floor_from_sample_replies() {
let runner = deterministic::Runner::timed(Duration::from_secs(30));
runner.start(|context| async move {
let mut harness =
Harness::setup(&context, 7, NZDuration!(Duration::from_secs(3600))).await;
harness.start_probes();
let mut subscription = harness.nodes[0].probe.subscribe();
let mut expected = None;
for index in 1..=3u8 {
let (_, finalization) = harness.finalization(index.into(), index);
if index == 3 {
expected = Some(finalization.clone());
}
harness.send_raw(index as usize, 0, finalization_bytes(finalization));
}
let expected = expected.expect("highest finalization present");
context.sleep(Duration::from_millis(100)).await;
let floor = subscription
.try_recv()
.expect("floor should resolve once enough replies arrive");
assert_eq!(floor, expected);
});
}
#[test_collect_traces]
fn test_retries_until_enough_replies(traces: TraceStorage) {
let retry_timeout = NZDuration!(Duration::from_millis(500));
let runner = deterministic::Runner::timed(Duration::from_secs(30));
runner.start(move |context| async move {
let mut harness = Harness::setup(&context, 4, retry_timeout).await;
harness.start_probes();
let start = context.current();
let mut subscription = harness.nodes[0].probe.subscribe();
let (_, finalization_f) = harness.finalization(1, 0xF);
let (_, finalization_g) = harness.finalization(2, 0x6);
harness.send_raw(1, 0, finalization_bytes(finalization_f.clone()));
context.sleep(Duration::from_millis(100)).await;
assert!(
matches!(
subscription.try_recv(),
Err(oneshot::error::TryRecvError::Empty)
),
"floor resolved before enough replies arrived"
);
context.sleep(retry_timeout.get()).await;
harness.send_raw(1, 0, finalization_bytes(finalization_f));
harness.send_raw(2, 0, finalization_bytes(finalization_g.clone()));
context.sleep(Duration::from_millis(100)).await;
let floor = subscription.try_recv().expect("floor resolved");
assert_eq!(floor, finalization_g);
let elapsed = context.current().duration_since(start).unwrap();
assert!(
elapsed >= retry_timeout.get(),
"floor resolved before a retry could occur ({elapsed:?})"
);
});
let events = traces.get_all();
events
.expect_event(|event| {
event.metadata.content == "re-requesting finalizations"
&& event
.metadata
.expect_field_exact("reason", "deadline elapsed")
.is_ok()
})
.expect("a deadline-driven retry should have occurred");
}
#[test]
fn test_waits_for_sample_size_even_with_matching_replies() {
let runner = deterministic::Runner::timed(Duration::from_secs(30));
runner.start(|context| async move {
let mut harness =
Harness::setup(&context, 7, NZDuration!(Duration::from_secs(3600))).await;
harness.start_probes();
let mut subscription = harness.nodes[0].probe.subscribe();
let (_, finalization) = harness.finalization(1, 1);
for index in 1..=2 {
harness.send_raw(index, 0, finalization_bytes(finalization.clone()));
}
context.sleep(Duration::from_millis(100)).await;
assert!(
matches!(
subscription.try_recv(),
Err(oneshot::error::TryRecvError::Empty)
),
"matching replies below the sample size must not resolve the floor"
);
});
}
#[test]
fn test_does_not_request_without_subscriber() {
let runner = deterministic::Runner::timed(Duration::from_secs(30));
runner.start(|context| async move {
let mut harness =
Harness::setup(&context, 4, NZDuration!(Duration::from_millis(500))).await;
harness.start_probes();
context.sleep(Duration::from_millis(100)).await;
let metrics = context.encode();
assert!(
!metrics.contains("network_messages_sent_total"),
"unexpected network messages before subscription: {metrics}"
);
});
}
#[test]
fn test_blocks_peer_sending_malformed_finalization() {
let runner = deterministic::Runner::timed(Duration::from_secs(30));
runner.start(|context| async move {
let mut harness =
Harness::setup(&context, 4, NZDuration!(Duration::from_millis(500))).await;
harness.start_probes();
let mut junk = vec![1u8];
junk.extend_from_slice(&[0xAB; 32]);
harness.send_raw(1, 0, junk);
context.sleep(Duration::from_millis(100)).await;
let blocked = harness.oracle.blocked().await.unwrap();
assert!(
blocked.contains(&(
harness.participants[0].clone(),
harness.participants[1].clone(),
)),
"node 0 should have blocked node 1"
);
});
}
#[test]
fn test_blocks_peer_sending_invalid_finalization() {
let runner = deterministic::Runner::timed(Duration::from_secs(30));
runner.start(|context| async move {
let mut harness =
Harness::setup(&context, 4, NZDuration!(Duration::from_millis(500))).await;
harness.start_probes();
let mut rng = test_rng();
let Fixture {
schemes: foreign, ..
} = scheme_mocks::fixture(&mut rng, b"_COMMONWARE_GLUE_PROBE_FOREIGN", 4);
let (_, finalization) = build_finalization(&foreign, 1, 1);
harness.send_raw(1, 0, finalization_bytes(finalization));
context.sleep(Duration::from_millis(100)).await;
let blocked = harness.oracle.blocked().await.unwrap();
assert!(
blocked.contains(&(
harness.participants[0].clone(),
harness.participants[1].clone(),
)),
"node 0 should have blocked node 1"
);
});
}
#[test]
fn test_blocks_peer_sending_invalid_message() {
let runner = deterministic::Runner::timed(Duration::from_secs(30));
runner.start(|context| async move {
let mut harness =
Harness::setup(&context, 4, NZDuration!(Duration::from_millis(500))).await;
harness.start_probes();
harness.send_raw(1, 0, vec![0xFF]);
context.sleep(Duration::from_millis(100)).await;
let blocked = harness.oracle.blocked().await.unwrap();
assert!(
blocked.contains(&(
harness.participants[0].clone(),
harness.participants[1].clone(),
)),
"node 0 should have blocked node 1"
);
});
}
#[test]
fn test_late_subscriber_receives_resolved_floor() {
let runner = deterministic::Runner::timed(Duration::from_secs(30));
runner.start(|context| async move {
let mut harness =
Harness::setup(&context, 4, NZDuration!(Duration::from_millis(500))).await;
let (block, finalization) = harness.finalization(1, 1);
for index in [1, 2, 3] {
harness
.inject(index, block.clone(), finalization.clone())
.await;
}
harness.start_probes();
let floor = harness.nodes[0]
.probe
.subscribe()
.await
.expect("floor resolved");
assert_eq!(floor, finalization);
let late = harness.nodes[0]
.probe
.subscribe()
.await
.expect("late subscriber served");
assert_eq!(late, finalization);
});
}
#[test]
fn test_ignores_finalizations_after_floor_set() {
let runner = deterministic::Runner::timed(Duration::from_secs(30));
runner.start(|context| async move {
let mut harness =
Harness::setup(&context, 4, NZDuration!(Duration::from_millis(500))).await;
let (block, finalization) = harness.finalization(1, 1);
for index in [1, 2, 3] {
harness
.inject(index, block.clone(), finalization.clone())
.await;
}
harness.start_probes();
let floor = harness.nodes[0]
.probe
.subscribe()
.await
.expect("floor resolved");
assert_eq!(floor, finalization);
let mut rng = test_rng();
let Fixture {
schemes: foreign, ..
} = scheme_mocks::fixture(&mut rng, b"_COMMONWARE_GLUE_PROBE_FOREIGN", 4);
let (_, invalid) = build_finalization(&foreign, 2, 9);
harness.send_raw(3, 0, finalization_bytes(invalid));
context.sleep(Duration::from_millis(100)).await;
let blocked = harness.oracle.blocked().await.unwrap();
assert!(
!blocked.contains(&(
harness.participants[0].clone(),
harness.participants[3].clone(),
)),
"finalizations after the floor is set must be ignored, not verified"
);
});
}
#[test]
fn test_duplicate_finalization_from_peer_is_ignored() {
let runner = deterministic::Runner::timed(Duration::from_secs(30));
runner.start(|context| async move {
let mut harness =
Harness::setup(&context, 7, NZDuration!(Duration::from_secs(3600))).await;
harness.start_probes();
let mut subscription = harness.nodes[0].probe.subscribe();
let (_, first) = harness.finalization(1, 1);
let (_, second) = harness.finalization(2, 2);
harness.send_raw(1, 0, finalization_bytes(first));
harness.send_raw(1, 0, finalization_bytes(second.clone()));
harness.send_raw(2, 0, finalization_bytes(second));
context.sleep(Duration::from_millis(100)).await;
let blocked = harness.oracle.blocked().await.unwrap();
assert!(
!blocked.contains(&(
harness.participants[0].clone(),
harness.participants[1].clone(),
)),
"a duplicate finalization must be ignored, not treated as a fault"
);
assert!(
matches!(
subscription.try_recv(),
Err(oneshot::error::TryRecvError::Empty)
),
"a duplicate finalization must not satisfy the sample size"
);
});
}
#[test]
fn test_invalid_duplicate_finalization_is_ignored() {
let runner = deterministic::Runner::timed(Duration::from_secs(30));
runner.start(|context| async move {
let mut harness =
Harness::setup(&context, 4, NZDuration!(Duration::from_secs(3600))).await;
harness.start_probes();
let _subscription = harness.nodes[0].probe.subscribe();
let (_, valid) = harness.finalization(1, 1);
harness.send_raw(1, 0, finalization_bytes(valid));
let mut rng = test_rng();
let Fixture {
schemes: foreign, ..
} = scheme_mocks::fixture(&mut rng, b"_COMMONWARE_GLUE_PROBE_FOREIGN", 4);
let (_, invalid) = build_finalization(&foreign, 1, 2);
harness.send_raw(1, 0, finalization_bytes(invalid));
context.sleep(Duration::from_millis(100)).await;
let blocked = harness.oracle.blocked().await.unwrap();
assert!(
!blocked.contains(&(
harness.participants[0].clone(),
harness.participants[1].clone(),
)),
"an invalid duplicate finalization must be ignored, not blocked"
);
});
}
#[test]
fn test_sample_selects_highest_over_stale_agreement() {
let runner = deterministic::Runner::timed(Duration::from_secs(30));
runner.start(|context| async move {
let mut harness =
Harness::setup(&context, 7, NZDuration!(Duration::from_secs(3600))).await;
harness.start_probes();
let mut subscription = harness.nodes[0].probe.subscribe();
let (_, stale) = harness.finalization(1, 0x0F);
let (_, newest) = harness.finalization(2, 0xA2);
harness.send_raw(1, 0, finalization_bytes(stale.clone()));
harness.send_raw(2, 0, finalization_bytes(stale.clone()));
harness.send_raw(3, 0, finalization_bytes(newest.clone()));
context.sleep(Duration::from_millis(100)).await;
let floor = subscription.try_recv().expect("floor resolved");
assert_eq!(floor, newest);
});
}
#[test]
fn test_resolves_floor_at_non_zero_epoch() {
let runner = deterministic::Runner::timed(Duration::from_secs(30));
runner.start(|context| async move {
let mut rng = test_rng();
let Fixture {
schemes: epoch_one, ..
} = scheme_mocks::fixture(&mut rng, b"_COMMONWARE_GLUE_FD_EPOCH_ONE", 4);
let provider = epoch_provider([(Epoch::new(1), epoch_one[0].clone())]);
let mut harness = Harness::setup_with(
&context,
4,
NZDuration!(Duration::from_secs(3600)),
Epoch::zero(),
{
let provider = provider.clone();
move |_scheme| provider.clone()
},
)
.await;
harness.start_probes();
let mut subscription = harness.nodes[0].probe.subscribe();
let (_, finalization) = build_finalization_at(&epoch_one, Epoch::new(1), 1, 7);
harness.send_raw(1, 0, finalization_bytes(finalization.clone()));
harness.send_raw(2, 0, finalization_bytes(finalization.clone()));
context.sleep(Duration::from_millis(100)).await;
let floor = subscription.try_recv().expect("floor resolved");
assert_eq!(floor, finalization);
});
}
#[test]
fn test_ignores_finalization_below_minimum_epoch() {
let runner = deterministic::Runner::timed(Duration::from_secs(30));
runner.start(|context| async move {
let provider = EpochProvider::default();
let mut harness = Harness::setup_with(
&context,
4,
NZDuration!(Duration::from_secs(3600)),
Epoch::new(1),
{
let provider = provider.clone();
move |_scheme| provider.clone()
},
)
.await;
provider.insert(Epoch::zero(), harness.schemes[0].clone());
provider.insert(Epoch::new(1), harness.schemes[0].clone());
harness.start_probes();
let mut subscription = harness.nodes[0].probe.subscribe();
let (_, old) = build_finalization_at(&harness.schemes, Epoch::zero(), 1, 1);
harness.send_raw(1, 0, finalization_bytes(old.clone()));
harness.send_raw(2, 0, finalization_bytes(old));
context.sleep(Duration::from_millis(50)).await;
assert!(
matches!(
subscription.try_recv(),
Err(oneshot::error::TryRecvError::Empty)
),
"below-minimum finalizations must not resolve the floor"
);
let (_, accepted) = build_finalization_at(&harness.schemes, Epoch::new(1), 2, 2);
harness.send_raw(1, 0, finalization_bytes(accepted.clone()));
harness.send_raw(2, 0, finalization_bytes(accepted.clone()));
context.sleep(Duration::from_millis(100)).await;
let floor = subscription.try_recv().expect("floor resolved");
assert_eq!(floor, accepted);
let blocked = harness.oracle.blocked().await.unwrap();
assert!(blocked.is_empty(), "no peer should be blocked");
});
}
#[test]
fn test_sample_size_uses_scoped_scheme_with_participantless_all_verifier() {
let runner = deterministic::Runner::timed(Duration::from_secs(30));
runner.start(|context| async move {
let mut rng = test_rng();
let Fixture {
participants,
schemes,
verifier,
..
} = scheme_mocks::fixture(&mut rng, b"_COMMONWARE_GLUE_FD_ALL_VERIFIER", 4);
let schemes: Vec<_> = schemes
.into_iter()
.map(|scheme| MaybeEnumerableScheme::new(scheme, true))
.collect();
let provider = ParticipantlessAllProvider {
verifier: Arc::new(MaybeEnumerableScheme::new(verifier.clone(), false)),
scheme: Arc::new(MaybeEnumerableScheme::new(verifier, true)),
};
let (network, oracle) = Network::new_with_peers(
context.child("network"),
SimConfig {
max_size: 1024 * 1024,
disconnect_on_block: true,
tracked_peer_sets: NZUsize!(1),
},
participants.clone(),
)
.await;
let _network = network.start();
for a in &participants {
for b in &participants {
if a != b {
oracle
.add_link(a.clone(), b.clone(), LINK)
.await
.expect("failed to add link");
}
}
}
let probe_network = oracle
.control(participants[0].clone())
.register(PROBE_CHANNEL, TEST_QUOTA)
.await
.expect("failed to register probe channel");
let (probe, probe_mailbox) = Probe::<
_,
MaybeEnumerableScheme,
ParticipantlessAllProvider,
Variant,
_,
ed25519::PublicKey,
_,
>::new(Config {
context: context.child("probe"),
provider,
strategy: Sequential,
capacity: NZUsize!(100),
blocker: oracle.control(participants[0].clone()),
minimum_epoch: Epoch::zero(),
retry_timeout: NZDuration!(Duration::from_secs(3600)),
});
let _probe = probe.start(probe_network);
let mut subscription = probe_mailbox.subscribe();
let mut peer_channels = Vec::new();
for public_key in participants.iter().take(3).skip(1) {
peer_channels.push(
oracle
.control(public_key.clone())
.register(PROBE_CHANNEL, TEST_QUOTA)
.await
.expect("failed to register peer probe channel"),
);
}
let (_, finalization) = build_finalization_at(&schemes, Epoch::zero(), 1, 1);
for (sender, _) in &peer_channels {
let mut sender = sender.clone();
sender.send(
Recipients::One(participants[0].clone()),
finalization_bytes(finalization.clone()),
false,
);
}
context.sleep(Duration::from_millis(100)).await;
let floor = subscription.try_recv().expect("floor resolved");
assert_eq!(floor, finalization);
});
}
#[test]
fn test_ignores_unknown_epoch_finalization() {
let runner = deterministic::Runner::timed(Duration::from_secs(30));
runner.start(|context| async move {
let mut rng = test_rng();
let Fixture {
schemes: epoch_one, ..
} = scheme_mocks::fixture(&mut rng, b"_COMMONWARE_GLUE_FD_EPOCH_ONE", 4);
let provider = epoch_provider([(Epoch::new(1), epoch_one[0].clone())]);
let mut harness = Harness::setup_with(
&context,
4,
NZDuration!(Duration::from_secs(3600)),
Epoch::zero(),
{
let provider = provider.clone();
move |_scheme| provider.clone()
},
)
.await;
harness.start_probes();
let (_, unknown) = build_finalization_at(&epoch_one, Epoch::new(5), 1, 1);
harness.send_raw(1, 0, finalization_bytes(unknown));
context.sleep(Duration::from_millis(100)).await;
let blocked = harness.oracle.blocked().await.unwrap();
assert!(
!blocked.contains(&(
harness.participants[0].clone(),
harness.participants[1].clone(),
)),
"an unknown-epoch finalization must be ignored, not blocked"
);
let mut subscription = harness.nodes[0].probe.subscribe();
context.sleep(Duration::from_millis(50)).await;
assert!(
matches!(
subscription.try_recv(),
Err(oneshot::error::TryRecvError::Empty)
),
"an unknown-epoch finalization must not resolve the floor"
);
});
}
#[test]
fn test_forgotten_epoch_finalization_is_not_counted() {
let runner = deterministic::Runner::timed(Duration::from_secs(30));
runner.start(|context| async move {
let provider = EpochProvider::default();
let mut harness = Harness::setup_with(
&context,
4,
NZDuration!(Duration::from_secs(3600)),
Epoch::zero(),
{
let provider = provider.clone();
move |_scheme| provider.clone()
},
)
.await;
provider.insert(Epoch::new(1), harness.schemes[0].clone());
provider.insert(Epoch::new(2), harness.schemes[0].clone());
harness.start_probes();
let mut subscription = harness.nodes[0].probe.subscribe();
let (_, epoch_one_finalization) =
build_finalization_at(&harness.schemes, Epoch::new(1), 1, 1);
harness.send_raw(1, 0, finalization_bytes(epoch_one_finalization));
context.sleep(Duration::from_millis(50)).await;
provider.forget(Epoch::new(1));
let (_, epoch_two_finalization) =
build_finalization_at(&harness.schemes, Epoch::new(2), 1, 2);
harness.send_raw(2, 0, finalization_bytes(epoch_two_finalization));
context.sleep(Duration::from_millis(50)).await;
assert!(
matches!(
subscription.try_recv(),
Err(oneshot::error::TryRecvError::Empty)
),
"a single judgeable vote must not resolve the floor"
);
let blocked = harness.oracle.blocked().await.unwrap();
assert!(blocked.is_empty(), "no peer should be blocked");
});
}
}