Skip to content
This repository was archived by the owner on Feb 27, 2025. It is now read-only.
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
12 changes: 6 additions & 6 deletions crates/hotshot/src/traits/networking/push_cdn_network.rs
Original file line number Diff line number Diff line change
Expand Up @@ -536,18 +536,18 @@ impl<K: SignatureKey + 'static> ConnectedNetwork<K> for PushCdnNetwork<K> {
/// - If we fail to serialize the message
/// - If we fail to send the direct message
async fn direct_message(&self, message: Vec<u8>, recipient: K) -> Result<(), NetworkError> {
// If we're paused, don't send the message
#[cfg(feature = "hotshot-testing")]
if self.is_paused.load(Ordering::Relaxed) {
return Ok(());
}

// If the message is to ourselves, just add it to the internal queue
if recipient == self.public_key {
self.internal_queue.lock().push_back(message);
return Ok(());
}

// If we're paused, don't send the message
#[cfg(feature = "hotshot-testing")]
if self.is_paused.load(Ordering::Relaxed) {
return Ok(());
}

// Send the message
if let Err(e) = self
.client
Expand Down
15 changes: 15 additions & 0 deletions crates/task-impls/src/da.rs
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@ use async_broadcast::{Receiver, Sender};
use async_lock::RwLock;
use async_trait::async_trait;
use hotshot_task::task::TaskState;
use hotshot_types::simple_vote::HasEpoch;
use hotshot_types::{
consensus::{Consensus, OuterConsensus},
data::{DaProposal2, PackedBundle},
Expand Down Expand Up @@ -236,10 +237,24 @@ impl<TYPES: NodeType, I: NodeImplementation<TYPES>, V: Versions> DaTaskState<TYP
let public_key = self.public_key.clone();
let chan = event_stream.clone();
let upgrade_lock = self.upgrade_lock.clone();
let proposal_epoch = proposal.data.epoch();
let next_epoch = proposal_epoch.map(|epoch| epoch + 1);

let membership_reader = membership.read().await;
let target_epoch = if membership_reader.has_stake(&public_key, proposal_epoch) {
proposal_epoch
} else if membership_reader.has_stake(&public_key, next_epoch) {
next_epoch
} else {
bail!("Not calculating VID, the node doesn't belong to the current epoch or the next epoch.");
};
drop(membership_reader);

spawn(async move {
Consensus::calculate_and_update_vid::<V>(
OuterConsensus::new(Arc::clone(&consensus.inner_consensus)),
view_number,
target_epoch,
membership,
&pk,
&upgrade_lock,
Expand Down
51 changes: 31 additions & 20 deletions crates/task-impls/src/request.rs
Original file line number Diff line number Diff line change
Expand Up @@ -30,7 +30,7 @@ use hotshot_types::{
node_implementation::{NodeImplementation, NodeType},
signature_key::SignatureKey,
},
utils::option_epoch_from_block_number,
utils::is_last_block_in_epoch,
vote::HasViewNumber,
};
use rand::{seq::SliceRandom, thread_rng};
Expand Down Expand Up @@ -113,21 +113,32 @@ impl<TYPES: NodeType, I: NodeImplementation<TYPES>> TaskState for NetworkRequest
match event.as_ref() {
HotShotEvent::QuorumProposalValidated(proposal, _) => {
let prop_view = proposal.data.view_number();
let prop_epoch = option_epoch_from_block_number::<TYPES>(
proposal.data.epoch().is_some(),
proposal.data.block_header().block_number(),
self.epoch_height,
);
let prop_epoch = proposal.data.epoch();
let next_epoch = prop_epoch.map(|epoch| epoch + 1);

// Request VID share only if:
// 1. we are part of the current epoch or
// 2. we are part of the next epoch and this is a proposal for the last block.
let membership_reader = self.membership.read().await;
if !membership_reader.has_stake(&self.public_key, prop_epoch)
&& (!membership_reader.has_stake(&self.public_key, next_epoch)
|| !is_last_block_in_epoch(
proposal.data.block_header().block_number(),
self.epoch_height,
))
{
return Ok(());
}
drop(membership_reader);

let consensus_reader = self.consensus.read().await;
let maybe_vid_share = consensus_reader
.vid_shares()
.get(&prop_view)
.and_then(|shares| shares.get(&self.public_key));
// If we already have the VID shares for the next view, do nothing.
if prop_view >= self.view
&& !self
.consensus
.read()
.await
.vid_shares()
.contains_key(&prop_view)
{
if prop_view >= self.view && maybe_vid_share.is_none() {
drop(consensus_reader);
self.spawn_requests(prop_view, prop_epoch, sender, receiver)
.await;
}
Expand Down Expand Up @@ -361,15 +372,15 @@ impl<TYPES: NodeType, I: NodeImplementation<TYPES>> NetworkRequestState<TYPES, I
) -> bool {
let consensus_reader = consensus.read().await;

let maybe_vid_share = consensus_reader
.vid_shares()
.get(view)
.and_then(|shares| shares.get(public_key));
let cancel = shutdown_flag.load(Ordering::Relaxed)
|| consensus_reader.vid_shares().contains_key(view)
|| maybe_vid_share.is_some()
|| consensus_reader.cur_view() > *view;
if cancel {
if let Some(Some(vid_share)) = consensus_reader
.vid_shares()
.get(view)
.map(|shares| shares.get(public_key).cloned())
{
if let Some(vid_share) = maybe_vid_share {
broadcast_event(
Arc::new(HotShotEvent::VidShareRecv(
public_key.clone(),
Expand Down
21 changes: 16 additions & 5 deletions crates/task-impls/src/response.rs
Original file line number Diff line number Diff line change
Expand Up @@ -84,14 +84,22 @@ impl<TYPES: NodeType, V: Versions> NetworkResponseState<TYPES, V> {
match event.as_ref() {
HotShotEvent::VidRequestRecv(request, sender) => {
let cur_epoch = self.consensus.read().await.cur_epoch();
let next_epoch = cur_epoch.map(|epoch| epoch + 1);
let target_epoch = if self.valid_sender(sender, cur_epoch).await {
cur_epoch
} else if self.valid_sender(sender, next_epoch).await {
next_epoch
} else {
// The sender neither belongs to the current nor to the next epoch.
continue;
};
// Verify request is valid
if !self.valid_sender(sender, cur_epoch).await
|| !valid_signature::<TYPES>(request, sender)
{
if !valid_signature::<TYPES>(request, sender) {
continue;
}
if let Some(proposal) =
self.get_or_calc_vid_share(request.view, sender).await
if let Some(proposal) = self
.get_or_calc_vid_share(request.view, target_epoch, sender)
.await
{
broadcast_event(
HotShotEvent::VidResponseSend(
Expand Down Expand Up @@ -151,6 +159,7 @@ impl<TYPES: NodeType, V: Versions> NetworkResponseState<TYPES, V> {
async fn get_or_calc_vid_share(
&self,
view: TYPES::View,
target_epoch: Option<TYPES::Epoch>,
key: &TYPES::SignatureKey,
) -> Option<Proposal<TYPES, VidDisperseShare<TYPES>>> {
let consensus_reader = self.consensus.read().await;
Expand All @@ -165,6 +174,7 @@ impl<TYPES: NodeType, V: Versions> NetworkResponseState<TYPES, V> {
if Consensus::calculate_and_update_vid::<V>(
OuterConsensus::new(Arc::clone(&self.consensus)),
view,
target_epoch,
Arc::clone(&self.membership),
&self.private_key,
&self.upgrade_lock,
Expand All @@ -177,6 +187,7 @@ impl<TYPES: NodeType, V: Versions> NetworkResponseState<TYPES, V> {
Consensus::calculate_and_update_vid::<V>(
OuterConsensus::new(Arc::clone(&self.consensus)),
view,
target_epoch,
Arc::clone(&self.membership),
&self.private_key,
&self.upgrade_lock,
Expand Down
177 changes: 175 additions & 2 deletions crates/testing/tests/tests_6/test_epochs.rs
Original file line number Diff line number Diff line change
Expand Up @@ -4,8 +4,6 @@
// You should have received a copy of the MIT License
// along with the HotShot repository. If not, see <https://mit-license.org/>.

use std::{collections::HashMap, time::Duration};

use hotshot_example_types::{
node_types::{
CombinedImpl, EpochUpgradeTestVersions, EpochsTestVersions, Libp2pImpl, MemoryImpl,
Expand All @@ -24,6 +22,7 @@ use hotshot_testing::{
view_sync_task::ViewSyncTaskDescription,
};
use hotshot_types::{data::ViewNumber, traits::node_implementation::ConsensusTime};
use std::{collections::HashMap, time::Duration};

cross_tests!(
TestName: test_success_with_epochs,
Expand Down Expand Up @@ -557,3 +556,177 @@ cross_tests!(
metadata
},
);

cross_tests!(
TestName: test_combined_network_with_epochs,
Impls: [CombinedImpl],
Types: [TestTypes, TestTwoStakeTablesTypes],
Versions: [EpochsTestVersions],
Ignore: false,
Metadata: {
let timing_data = TimingData {
next_view_timeout: 10_000,
..Default::default()
};

let overall_safety_properties = OverallSafetyPropertiesDescription {
num_failed_views: 0,
num_successful_views: 25,
..Default::default()
};

let completion_task_description = CompletionTaskDescription::TimeBasedCompletionTaskBuilder(
TimeBasedCompletionTaskDescription {
duration: Duration::from_secs(120),
},
);

let mut metadata = TestDescription::default_multiple_rounds();
metadata.timing_data = timing_data;
metadata.overall_safety_properties = overall_safety_properties;
metadata.completion_task_description = completion_task_description;

metadata
},
);

// A run where the CDN crashes part-way through, epochs enabled.
cross_tests!(
TestName: test_combined_network_cdn_crash_with_epochs,
Impls: [CombinedImpl],
Types: [TestTypes, TestTwoStakeTablesTypes],
Versions: [EpochsTestVersions],
Ignore: false,
Metadata: {
let timing_data = TimingData {
next_view_timeout: 10_000,
..Default::default()
};

let overall_safety_properties = OverallSafetyPropertiesDescription {
num_failed_views: 0,
num_successful_views: 35,
..Default::default()
};

let completion_task_description = CompletionTaskDescription::TimeBasedCompletionTaskBuilder(
TimeBasedCompletionTaskDescription {
duration: Duration::from_secs(220),
},
);

let mut metadata = TestDescription::default_multiple_rounds();
metadata.timing_data = timing_data;
metadata.overall_safety_properties = overall_safety_properties;
metadata.completion_task_description = completion_task_description;

let mut all_nodes = vec![];
for node in 0..metadata.test_config.num_nodes_with_stake.into() {
all_nodes.push(ChangeNode {
idx: node,
updown: NodeAction::NetworkDown,
});
}

metadata.spinning_properties = SpinningTaskDescription {
node_changes: vec![(5, all_nodes)],
};

metadata
},
);

cross_tests!(
TestName: test_combined_network_reup_with_epochs,
Impls: [CombinedImpl],
Types: [TestTypes, TestTwoStakeTablesTypes],
Versions: [EpochsTestVersions],
Ignore: false,
Metadata: {
let timing_data = TimingData {
next_view_timeout: 10_000,
..Default::default()
};

let overall_safety_properties = OverallSafetyPropertiesDescription {
num_failed_views: 0,
num_successful_views: 35,
..Default::default()
};

let completion_task_description = CompletionTaskDescription::TimeBasedCompletionTaskBuilder(
TimeBasedCompletionTaskDescription {
duration: Duration::from_secs(220),
},
);

let mut metadata = TestDescription::default_multiple_rounds();
metadata.timing_data = timing_data;
metadata.overall_safety_properties = overall_safety_properties;
metadata.completion_task_description = completion_task_description;

let mut all_down = vec![];
let mut all_up = vec![];
for node in 0..metadata.test_config.num_nodes_with_stake.into() {
all_down.push(ChangeNode {
idx: node,
updown: NodeAction::NetworkDown,
});
all_up.push(ChangeNode {
idx: node,
updown: NodeAction::NetworkUp,
});
}

metadata.spinning_properties = SpinningTaskDescription {
node_changes: vec![(13, all_up), (5, all_down)],
};

metadata
},
);

cross_tests!(
TestName: test_combined_network_half_dc_with_epochs,
Impls: [CombinedImpl],
Types: [TestTypes, TestTwoStakeTablesTypes],
Versions: [EpochsTestVersions],
Ignore: false,
Metadata: {
let timing_data = TimingData {
next_view_timeout: 10_000,
..Default::default()
};

let overall_safety_properties = OverallSafetyPropertiesDescription {
num_failed_views: 0,
num_successful_views: 35,
..Default::default()
};

let completion_task_description = CompletionTaskDescription::TimeBasedCompletionTaskBuilder(
TimeBasedCompletionTaskDescription {
duration: Duration::from_secs(220),
},
);

let mut metadata = TestDescription::default_multiple_rounds();
metadata.timing_data = timing_data;
metadata.overall_safety_properties = overall_safety_properties;
metadata.completion_task_description = completion_task_description;

let mut half = vec![];
for node in 0..usize::from(metadata.test_config.num_nodes_with_stake) / 2 {
half.push(ChangeNode {
idx: node,
updown: NodeAction::NetworkDown,
});
}

metadata.spinning_properties = SpinningTaskDescription {
node_changes: vec![(5, half)],
};

metadata
},
);
3 changes: 2 additions & 1 deletion crates/types/src/consensus.rs
Original file line number Diff line number Diff line change
Expand Up @@ -955,6 +955,7 @@ impl<TYPES: NodeType> Consensus<TYPES> {
pub async fn calculate_and_update_vid<V: Versions>(
consensus: OuterConsensus<TYPES>,
view: <TYPES as NodeType>::View,
target_epoch: Option<<TYPES as NodeType>::Epoch>,
membership: Arc<RwLock<TYPES::Membership>>,
private_key: &<TYPES::SignatureKey as SignatureKey>::PrivateKey,
upgrade_lock: &UpgradeLock<TYPES, V>,
Expand All @@ -972,7 +973,7 @@ impl<TYPES: NodeType> Consensus<TYPES> {
payload.as_ref(),
&membership,
view,
epoch,
target_epoch,
epoch,
upgrade_lock,
)
Expand Down