Skip to content
Merged
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
146 changes: 2 additions & 144 deletions crates/blockchain/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,7 @@ use crate::aggregation::{
AggregationSession, PRIOR_WORKER_JOIN_TIMEOUT, run_aggregation_worker,
};
use crate::key_manager::ValidatorKeyPair;
use crate::sync_status::SyncStatusTracker;
use spawned_concurrency::actor;
use spawned_concurrency::error::ActorError;
use spawned_concurrency::protocol;
Expand All @@ -35,6 +36,7 @@ pub mod key_manager;
pub mod metrics;
pub mod reaggregate;
pub mod store;
mod sync_status;

pub struct BlockChain {
handle: ActorRef<BlockChainServer>,
Expand All @@ -55,56 +57,6 @@ pub use ethlambda_types::block::MAX_ATTESTATIONS_DATA;
///
/// See: leanSpec PR #682.
pub const GOSSIP_DISPARITY_INTERVALS: u64 = 1;
/// Local head lag beyond which the node is considered to be syncing.
///
/// See: leanSpec PR #708.
const SYNC_LAG_THRESHOLD: u64 = 4;
/// Freshest-known block lag beyond which the network is considered stalled.
///
/// During a network-wide stall the node remains synced so validators can help
/// the chain recover.
const NETWORK_STALL_THRESHOLD: u64 = 8;
/// Recovery band that prevents the sync status from flapping near the threshold.
const SYNC_HYSTERESIS_BAND: u64 = 2;

#[derive(Default)]
struct SyncStatusTracker {
syncing: bool,
}

impl SyncStatusTracker {
fn update(
&mut self,
current_slot: u64,
head_slot: u64,
max_seen_slot: u64,
) -> metrics::SyncStatus {
let head_lag = current_slot.saturating_sub(head_slot);
let network_lag = current_slot.saturating_sub(max_seen_slot);

if network_lag > NETWORK_STALL_THRESHOLD {
self.syncing = false;
} else if self.syncing {
self.syncing = head_lag > SYNC_LAG_THRESHOLD.saturating_sub(SYNC_HYSTERESIS_BAND);
} else {
self.syncing = head_lag > SYNC_LAG_THRESHOLD;
}

if self.syncing {
metrics::SyncStatus::Syncing
} else {
metrics::SyncStatus::Synced
}
}

fn duties_allowed(&self) -> bool {
!self.syncing
}

fn gate_proposer(&self, proposer: Option<u64>) -> Option<u64> {
proposer.filter(|_| self.duties_allowed())
}
}

/// Milliseconds until the next interval boundary, measured relative to genesis.
fn ms_until_next_interval(now_ms: u64, genesis_time_ms: u64) -> u64 {
Expand Down Expand Up @@ -1057,97 +1009,3 @@ impl Handler<AggregationDeadline> for BlockChainServer {
}
}
}

#[cfg(test)]
mod tests {
use super::*;

#[test]
fn sync_status_allows_lag_through_threshold() {
let mut tracker = SyncStatusTracker::default();

for lag in 0..=SYNC_LAG_THRESHOLD {
assert_eq!(
tracker.update(10 + lag, 10, 10 + lag),
metrics::SyncStatus::Synced
);
}

let first_syncing_slot = 10 + SYNC_LAG_THRESHOLD + 1;
assert_eq!(
tracker.update(first_syncing_slot, 10, first_syncing_slot),
metrics::SyncStatus::Syncing
);
}

#[test]
fn sync_status_detects_local_lag_when_fresh_blocks_are_known() {
let mut tracker = SyncStatusTracker::default();
let current_slot = 10 + SYNC_LAG_THRESHOLD + 1;

assert_eq!(
tracker.update(current_slot, 10, current_slot),
metrics::SyncStatus::Syncing
);
}

#[test]
fn sync_status_treats_stale_known_blocks_as_network_stall() {
let mut tracker = SyncStatusTracker::default();

assert_eq!(tracker.update(100, 0, 0), metrics::SyncStatus::Synced);
}

#[test]
fn sync_status_hysteresis_prevents_flapping() {
let mut tracker = SyncStatusTracker::default();

assert_eq!(tracker.update(15, 10, 15), metrics::SyncStatus::Syncing);
assert_eq!(tracker.update(15, 11, 15), metrics::SyncStatus::Syncing);
assert_eq!(tracker.update(15, 10, 15), metrics::SyncStatus::Syncing);
assert_eq!(tracker.update(15, 13, 15), metrics::SyncStatus::Synced);
}

#[test]
fn network_stall_reopens_sync_status() {
let mut tracker = SyncStatusTracker::default();

assert_eq!(tracker.update(20, 0, 20), metrics::SyncStatus::Syncing);
assert_eq!(tracker.update(30, 0, 20), metrics::SyncStatus::Synced);
}

#[test]
fn future_head_saturates_lag_at_zero() {
let mut tracker = SyncStatusTracker::default();

assert_eq!(tracker.update(15, 20, 20), metrics::SyncStatus::Synced);
}

#[test]
fn syncing_gates_proposals_and_attestations() {
let mut tracker = SyncStatusTracker::default();
tracker.update(20, 0, 20);

assert!(!tracker.duties_allowed());
assert_eq!(tracker.gate_proposer(Some(3)), None);
}

#[test]
fn caught_up_node_allows_proposals_and_attestations() {
let mut tracker = SyncStatusTracker::default();
tracker.update(20, 0, 20);
tracker.update(20, 18, 20);

assert!(tracker.duties_allowed());
assert_eq!(tracker.gate_proposer(Some(3)), Some(3));
}

#[test]
fn network_stall_keeps_proposals_and_attestations_enabled() {
let mut tracker = SyncStatusTracker::default();
tracker.update(100, 0, 0);

assert!(tracker.duties_allowed());
assert_eq!(tracker.gate_proposer(Some(3)), Some(3));
}
}
143 changes: 143 additions & 0 deletions crates/blockchain/src/sync_status.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,143 @@
use crate::metrics::SyncStatus;

/// Local head lag beyond which the node is considered to be syncing.
///
/// See: leanSpec PR #708.
const SYNC_LAG_THRESHOLD: u64 = 4;
/// Freshest-known block lag beyond which the network is considered stalled.
///
/// During a network-wide stall the node remains synced so validators can help
/// the chain recover.
const NETWORK_STALL_THRESHOLD: u64 = 8;
/// Recovery band that prevents the sync status from flapping near the threshold.
const SYNC_HYSTERESIS_BAND: u64 = 2;

#[derive(Default)]
pub(crate) struct SyncStatusTracker {
syncing: bool,
}

impl SyncStatusTracker {
pub(crate) fn update(
&mut self,
current_slot: u64,
head_slot: u64,
max_seen_slot: u64,
) -> SyncStatus {
let head_lag = current_slot.saturating_sub(head_slot);
let network_lag = current_slot.saturating_sub(max_seen_slot);

if network_lag > NETWORK_STALL_THRESHOLD {
self.syncing = false;
} else if self.syncing {
self.syncing = head_lag > SYNC_LAG_THRESHOLD.saturating_sub(SYNC_HYSTERESIS_BAND);
} else {
self.syncing = head_lag > SYNC_LAG_THRESHOLD;
}

if self.syncing {
SyncStatus::Syncing
} else {
SyncStatus::Synced
}
}

pub(crate) fn duties_allowed(&self) -> bool {
!self.syncing
}

pub(crate) fn gate_proposer(&self, proposer: Option<u64>) -> Option<u64> {
proposer.filter(|_| self.duties_allowed())
}
}

#[cfg(test)]
mod tests {
use super::*;

#[test]
fn sync_status_allows_lag_through_threshold() {
let mut tracker = SyncStatusTracker::default();

for lag in 0..=SYNC_LAG_THRESHOLD {
assert_eq!(tracker.update(10 + lag, 10, 10 + lag), SyncStatus::Synced);
}

let first_syncing_slot = 10 + SYNC_LAG_THRESHOLD + 1;
assert_eq!(
tracker.update(first_syncing_slot, 10, first_syncing_slot),
SyncStatus::Syncing
);
}

#[test]
fn sync_status_detects_local_lag_when_fresh_blocks_are_known() {
let mut tracker = SyncStatusTracker::default();
let current_slot = 10 + SYNC_LAG_THRESHOLD + 1;

assert_eq!(
tracker.update(current_slot, 10, current_slot),
SyncStatus::Syncing
);
}

#[test]
fn sync_status_treats_stale_known_blocks_as_network_stall() {
let mut tracker = SyncStatusTracker::default();

assert_eq!(tracker.update(100, 0, 0), SyncStatus::Synced);
}

#[test]
fn sync_status_hysteresis_prevents_flapping() {
let mut tracker = SyncStatusTracker::default();

assert_eq!(tracker.update(15, 10, 15), SyncStatus::Syncing);
assert_eq!(tracker.update(15, 11, 15), SyncStatus::Syncing);
assert_eq!(tracker.update(15, 10, 15), SyncStatus::Syncing);
assert_eq!(tracker.update(15, 13, 15), SyncStatus::Synced);
}

#[test]
fn network_stall_reopens_sync_status() {
let mut tracker = SyncStatusTracker::default();

assert_eq!(tracker.update(20, 0, 20), SyncStatus::Syncing);
assert_eq!(tracker.update(30, 0, 20), SyncStatus::Synced);
}

#[test]
fn future_head_saturates_lag_at_zero() {
let mut tracker = SyncStatusTracker::default();

assert_eq!(tracker.update(15, 20, 20), SyncStatus::Synced);
}

#[test]
fn syncing_gates_proposals_and_attestations() {
let mut tracker = SyncStatusTracker::default();
tracker.update(20, 0, 20);

assert!(!tracker.duties_allowed());
assert_eq!(tracker.gate_proposer(Some(3)), None);
}

#[test]
fn caught_up_node_allows_proposals_and_attestations() {
let mut tracker = SyncStatusTracker::default();
tracker.update(20, 0, 20);
tracker.update(20, 18, 20);

assert!(tracker.duties_allowed());
assert_eq!(tracker.gate_proposer(Some(3)), Some(3));
}

#[test]
fn network_stall_keeps_proposals_and_attestations_enabled() {
let mut tracker = SyncStatusTracker::default();
tracker.update(100, 0, 0);

assert!(tracker.duties_allowed());
assert_eq!(tracker.gate_proposer(Some(3)), Some(3));
}
}