use super::super::trace::{
ordered_send_error_label, queue_send_trace as qs_trace, QUEUE_SEND_TABLE,
};
use super::{
config::*, events::*, peer_registry::*, sequencer::*, sequencer_task::*, state::*, wire::*, *,
};
use crate::zakura::{
FrontierChange, FrontierUpdate, OrderedSendError, ServiceAdmissionDecision,
ServicePeerDirection, ServicePeerSnapshot, ZakuraBlockSyncCandidateState,
};
use iroh::NodeId;
const ACTION_SEND_TIMEOUT: Duration = Duration::from_secs(5);
const BS_ACTION_SPARE_POOL: usize = 128;
const ROUTINE_TO_REACTOR_DEPTH: usize = 1024;
const NEEDED_BLOCK_REFILL_LIMIT: u32 = 4_000;
#[derive(Copy, Clone, Debug, Eq, PartialEq)]
struct FloorGapDiagnostics {
height: block::Height,
state: &'static str,
servable_peers: usize,
available_peers: usize,
outstanding_peers: usize,
oldest_outstanding_ms: Option<u64>,
next_deadline_ms: Option<u64>,
}
#[derive(Copy, Clone, Debug)]
struct RangeResponseTrace {
start_height: block::Height,
requested_count: u32,
sent_count: u32,
sent_bytes: u64,
reason: &'static str,
prepare_elapsed: Option<Duration>,
send_elapsed: Duration,
total_elapsed: Option<Duration>,
}
pub fn spawn_block_sync_reactor(
startup: BlockSyncStartup,
) -> (
BlockSyncHandle,
mpsc::Receiver<BlockSyncAction>,
JoinHandle<()>,
) {
debug_assert!(
!startup.state_queries_enabled
|| (startup.header_tip.is_some() ^ startup.frontier_updates.is_some()),
"state-backed block sync must have exactly one frontier source",
);
let state = BlockSyncState::new(&startup);
let (events_tx, events_rx) =
mpsc::channel(startup.config.peer_limits.inbound_queue_depth.max(1));
let events_keepalive = events_tx.clone();
let (lifecycle_tx, lifecycle_rx) = mpsc::unbounded_channel();
let actions_capacity = startup
.config
.submitted_apply_limit()
.saturating_add(BS_ACTION_SPARE_POOL);
let (actions_tx, actions_rx) = mpsc::channel(actions_capacity);
let (peers_tx, peers_rx) = watch::channel(state.peer_snapshot(startup.config.peer_limits));
let (status_tx, status_rx) = watch::channel(state.last_advertised_status);
let (candidates_tx, candidates_rx) = watch::channel(ZakuraBlockSyncCandidateState::default());
let sequencer = Sequencer::new(
startup.frontiers.verified_block_tip,
startup.config.submitted_apply_limit(),
);
let committed_throughput = ThroughputMeter::new(Instant::now());
let (sequencer_input_tx, sequencer_body_input_rx) =
mpsc::channel(startup.config.submitted_apply_limit().max(1));
let (sequencer_control_tx, sequencer_control_rx) = mpsc::unbounded_channel();
let sequencer_input_bytes = Arc::new(std::sync::atomic::AtomicU64::new(0));
let sequencer_input_decoded_attributed_memory_bytes =
Arc::new(std::sync::atomic::AtomicU64::new(0));
let (sequencer_view_tx, sequencer_view_rx) = watch::channel(initial_view(startup.frontiers));
let sequencer_task = SequencerTask::new(
sequencer,
state.budget.clone(),
state.work_queue.clone(),
actions_tx.clone(),
committed_throughput,
startup.frontiers,
sequencer_body_input_rx,
sequencer_control_rx,
sequencer_input_bytes.clone(),
sequencer_input_decoded_attributed_memory_bytes.clone(),
sequencer_view_tx,
ACTION_SEND_TIMEOUT,
startup.trace.clone(),
);
tokio::spawn(sequencer_task.run());
let registry = Arc::new(PeerRegistry::new());
let (routine_to_reactor_tx, routine_to_reactor_rx) = mpsc::channel(ROUTINE_TO_REACTOR_DEPTH);
let routine_to_reactor_keepalive = routine_to_reactor_tx.clone();
let routine_wiring = RoutineWiring {
config: startup.config.clone(),
budget: state.budget.clone(),
work: state.work_queue.clone(),
registry: registry.clone(),
received_throughput: state.received_throughput.clone(),
sequencer_input: sequencer_input_tx.clone(),
sequencer_input_bytes: sequencer_input_bytes.clone(),
sequencer_input_decoded_attributed_memory_bytes:
sequencer_input_decoded_attributed_memory_bytes.clone(),
actions: actions_tx.clone(),
routine_to_reactor: routine_to_reactor_tx,
view: sequencer_view_rx.clone(),
trace: startup.trace.clone(),
};
let handle = BlockSyncHandle {
events: events_tx,
lifecycle: lifecycle_tx,
peers: peers_rx,
status: status_rx,
candidates: candidates_rx,
routine_wiring: Some(routine_wiring),
};
let reactor = BlockSyncReactor {
verified_block_tip: startup.frontiers.verified_block_tip,
request_floor: startup.frontiers.verified_block_tip,
pending_needed_query: None,
last_reset_epoch: 0,
last_reaction_epoch: 0,
last_view: initial_view(startup.frontiers),
startup,
state,
registry,
events: events_rx,
_events_keepalive: events_keepalive,
lifecycle: lifecycle_rx,
actions: actions_tx,
routine_to_reactor: routine_to_reactor_rx,
_routine_to_reactor_keepalive: routine_to_reactor_keepalive,
peers: peers_tx,
status: status_tx,
candidates: candidates_tx,
sequencer_input: sequencer_input_tx,
sequencer_input_bytes,
sequencer_input_decoded_attributed_memory_bytes,
sequencer_control: sequencer_control_tx,
sequencer_view: sequencer_view_rx,
};
let task = tokio::spawn(reactor.run());
(handle, actions_rx, task)
}
#[derive(Debug)]
pub(super) struct BlockSyncReactor {
startup: BlockSyncStartup,
state: BlockSyncState,
registry: Arc<PeerRegistry>,
events: mpsc::Receiver<BlockSyncEvent>,
_events_keepalive: mpsc::Sender<BlockSyncEvent>,
lifecycle: mpsc::UnboundedReceiver<BlockSyncEvent>,
actions: mpsc::Sender<BlockSyncAction>,
routine_to_reactor: mpsc::Receiver<RoutineToReactor>,
_routine_to_reactor_keepalive: mpsc::Sender<RoutineToReactor>,
peers: watch::Sender<ServicePeerSnapshot>,
status: watch::Sender<BlockSyncStatus>,
candidates: watch::Sender<ZakuraBlockSyncCandidateState>,
sequencer_input: mpsc::Sender<SequencedBody>,
sequencer_input_bytes: Arc<std::sync::atomic::AtomicU64>,
sequencer_input_decoded_attributed_memory_bytes: Arc<std::sync::atomic::AtomicU64>,
sequencer_control: mpsc::UnboundedSender<SequencerControlInput>,
sequencer_view: watch::Receiver<SequencerView>,
verified_block_tip: block::Height,
request_floor: block::Height,
pending_needed_query: Option<(block::Height, u32, block::Height, block::Hash)>,
last_reset_epoch: u64,
last_reaction_epoch: u64,
last_view: SequencerView,
}
impl BlockSyncReactor {
async fn run(mut self) {
let mut header_tip = self.startup.header_tip.clone();
let mut header_tip_open = header_tip.is_some();
let mut frontier_updates = self.startup.frontier_updates.clone();
let mut frontier_updates_open = frontier_updates.is_some();
set_block_reactor_active_connection_gauge(self.state.peers.len());
let mut metrics_ticks = time::interval(self.startup.config.request_timeout);
let mut status_ticks = time::interval(
self.startup
.config
.status_refresh_interval
.max(Duration::from_millis(1)),
);
self.query_needed_blocks().await;
self.publish_metrics();
self.refresh_throughput();
self.trace_sync_state(true);
loop {
let floor_watchdog = self.earliest_floor_deadline_sleep();
tokio::pin!(floor_watchdog);
tokio::select! {
_ = self.startup.shutdown.cancelled() => break,
event = self.lifecycle.recv() => {
let Some(event) = event else { break };
self.handle_event(event).await;
}
event = self.events.recv() => {
let Some(event) = event else { break };
self.handle_event(event).await;
}
changed = async {
match header_tip.as_mut() {
Some(header_tip) => header_tip.changed().await,
None => std::future::pending().await,
}
}, if header_tip_open => {
match changed {
Ok(()) => {
let header_tip = header_tip
.as_mut()
.expect("header tip receiver exists while header_tip_open is true");
let (height, hash) = *header_tip.borrow_and_update();
self.handle_header_tip_changed(height, hash).await;
self.publish_metrics();
}
Err(_) => header_tip_open = false,
}
}
changed = async {
match frontier_updates.as_mut() {
Some(frontier_updates) => frontier_updates.changed().await,
None => std::future::pending().await,
}
}, if frontier_updates_open => {
match changed {
Ok(()) => {
let frontier_updates = frontier_updates
.as_mut()
.expect("frontier update receiver exists while frontier_updates_open is true");
let update = *frontier_updates.borrow_and_update();
self.handle_frontier_update(update).await;
self.publish_metrics();
}
Err(_) => frontier_updates_open = false,
}
}
changed = self.sequencer_view.changed() => {
match changed {
Ok(()) => {
let view = *self.sequencer_view.borrow_and_update();
self.on_sequencer_view_changed(view).await;
self.publish_metrics();
self.trace_sync_state(false);
}
Err(_) => break,
}
}
message = self.routine_to_reactor.recv() => {
match message {
Some(message) => self.handle_routine_message(message).await,
None => break,
}
}
_ = metrics_ticks.tick() => {
self.publish_metrics();
self.refresh_throughput();
self.trace_sync_state(true);
}
_ = status_ticks.tick() => self.flush_status_refresh().await,
_ = &mut floor_watchdog => {
self.run_floor_watchdog(Instant::now());
self.publish_metrics();
}
}
}
}
fn earliest_floor_deadline_sleep(&self) -> time::Sleep {
let earliest = next_height(self.request_floor)
.and_then(|height| self.registry.earliest_outstanding_deadline_at(height));
match earliest {
Some(deadline) => time::sleep(deadline.saturating_duration_since(Instant::now())),
None => time::sleep(Duration::from_secs(3600)),
}
}
fn run_floor_watchdog(&mut self, now: Instant) {
let Some(height) = next_height(self.request_floor) else {
return;
};
let (servable_peers, _) = self.registry.floor_gap_servable(height);
let claims = self.registry.outstanding_claims_at(height);
for claim in claims {
if claim.meta.deadline > now {
continue;
}
self.registry
.clear_outstanding_height(&claim.peer, claim.height);
if servable_peers > 2 {
self.registry.avoid_floor_height_until(
&claim.peer,
claim.height,
now + self.startup.config.effective_floor_peer_avoid_cooldown(),
);
}
let released = self
.state
.work_queue
.release_reserved_and_return_items_detailed([claim.height]);
self.state.budget.release(released.released_bytes);
self.emit_trace(bs_trace::BLOCK_FLOOR_WATCHDOG_CANCELLED, |row| {
bs_insert_peer(row, bs_trace::PEER, &claim.peer);
bs_insert_height(row, bs_trace::HEIGHT, claim.height);
bs_insert_u64(row, bs_trace::ESTIMATED_BYTES, claim.meta.estimated_bytes);
bs_insert_u64(row, "released_bytes", released.released_bytes);
bs_insert_u64(row, "returned_count", released.returned_count);
bs_insert_u64(row, "already_pending_count", released.already_pending_count);
bs_insert_u64(row, "released_count", released.released_count);
bs_insert_u64(row, "missing_count", released.missing_count);
bs_insert_u64(
row,
"pending_after",
u64::from(self.state.work_queue.pending_contains(claim.height)),
);
bs_insert_u64(
row,
"in_flight_after",
u64::from(self.state.work_queue.in_flight_contains(claim.height)),
);
});
metrics::counter!("sync.block.floor_watchdog.cancelled").increment(1);
tracing::debug!(
peer = ?claim.peer,
height = ?claim.height,
estimated_bytes = claim.meta.estimated_bytes,
released = released.released_bytes,
"force-cancelled expired Zakura block-sync floor request"
);
}
}
async fn handle_event(&mut self, event: BlockSyncEvent) {
self.trace_event_received(&event);
match event {
BlockSyncEvent::PeerConnected(session) => self.handle_peer_connected(session).await,
BlockSyncEvent::PeerDisconnected(peer) => self.handle_peer_disconnected(peer),
BlockSyncEvent::HeaderTipChanged { height, hash } => {
self.handle_header_tip_changed(height, hash).await
}
BlockSyncEvent::StateFrontiersChanged(frontiers) => {
self.handle_state_frontiers_changed(frontiers).await
}
BlockSyncEvent::ChainTipGrow(frontiers) => {
self.handle_state_frontiers_changed(frontiers).await
}
BlockSyncEvent::ChainTipReset(frontiers) => {
self.handle_chain_tip_reset(frontiers, true).await
}
BlockSyncEvent::NeededBlocks(blocks) => {
self.handle_needed_blocks(blocks).await;
}
BlockSyncEvent::BlockApplyFinished {
token,
height,
hash,
result,
local_frontier,
} => {
self.handle_block_apply_finished(token, height, hash, result, local_frontier)
.await
}
BlockSyncEvent::BlockRangeResponseReady {
peer,
start_height,
requested_count,
blocks,
} => {
self.handle_block_range_response_ready(peer, start_height, requested_count, blocks)
.await;
}
BlockSyncEvent::BlockRangeResponseFinished {
peer,
start_height,
requested_count,
returned_count,
} => {
self.handle_block_range_response_finished(
peer,
start_height,
requested_count,
returned_count,
)
.await;
}
}
self.publish_metrics();
}
fn admission_decision_for(
&self,
peer: &ZakuraPeerId,
direction: ServicePeerDirection,
) -> ServiceAdmissionDecision {
if self.state.peers.contains_key(peer) {
return ServiceAdmissionDecision::Admit;
}
let limits = self.startup.config.peer_limits;
let admitted = self.admitted_count(direction);
let cap = match direction {
ServicePeerDirection::Inbound => limits.max_inbound_peers,
ServicePeerDirection::Outbound => limits.max_outbound_peers,
};
if admitted >= cap {
ServiceAdmissionDecision::RejectFull
} else {
ServiceAdmissionDecision::Admit
}
}
fn admitted_count(&self, direction: ServicePeerDirection) -> usize {
self.state
.peers
.values()
.filter(|peer| peer.direction == direction)
.count()
}
fn publish_peer_snapshot(&self) {
let _ = self
.peers
.send(self.state.peer_snapshot(self.startup.config.peer_limits));
}
fn publish_candidate_state(&self) {
let has_body_gaps = !self.state.needed_heights.is_empty();
let needed = &self.state.needed_heights;
let mut admitted_node_ids: Vec<_> = self
.registry
.candidate_snapshot()
.into_iter()
.filter_map(|(peer_id, received_status, servable_low, servable_high)| {
if has_body_gaps {
let can_serve_any = received_status
&& needed
.iter()
.any(|height| servable_low <= *height && *height <= servable_high);
if !can_serve_any {
return None;
}
}
node_id_from_block_peer_id(&peer_id)
})
.collect();
admitted_node_ids.sort_by(|left, right| left.as_bytes().cmp(right.as_bytes()));
admitted_node_ids.dedup();
let _ = self.candidates.send(ZakuraBlockSyncCandidateState {
missing_block_bodies: self.state.needed_heights.clone(),
admitted_node_ids,
});
}
async fn handle_peer_connected(&mut self, session: BlockSyncPeerSession) {
let peer = session.peer_id().clone();
let direction = session.direction();
let decision = self.admission_decision_for(&peer, direction);
if decision != ServiceAdmissionDecision::Admit {
metrics::counter!("sync.block.peer.parked").increment(1);
tracing::info!(
?peer,
?direction,
?decision,
"locally parking Zakura block-sync service session"
);
self.state.parked_peers.insert(peer.clone());
session.cancel_token().cancel();
self.registry.remove(&peer);
self.publish_peer_snapshot();
self.publish_candidate_state();
return;
}
self.state.parked_peers.remove(&peer);
let peer_state = PeerBlockState::new(session, &self.startup.config);
self.state.peers.insert(peer.clone(), peer_state);
set_block_reactor_active_connection_gauge(self.state.peers.len());
self.trace_peer_connected(&peer, direction, self.state.peers.len());
self.publish_peer_snapshot();
self.publish_candidate_state();
self.send_status_and_mark_refresh(&peer, "peer_connected", Instant::now());
}
fn handle_peer_disconnected(&mut self, peer: ZakuraPeerId) {
if self.state.peers.remove(&peer).is_some() {
set_block_reactor_active_connection_gauge(self.state.peers.len());
self.trace_peer_disconnected(
&peer,
self.registry_received_status(&peer),
self.state.peers.len(),
);
}
self.registry.remove(&peer);
self.state.parked_peers.remove(&peer);
self.publish_peer_snapshot();
self.publish_candidate_state();
}
fn registry_received_status(&self, peer: &ZakuraPeerId) -> bool {
self.registry.has_received_status(peer)
}
async fn handle_header_tip_changed(&mut self, height: block::Height, hash: block::Hash) {
self.state.best_header_tip = height;
self.state.best_header_hash = hash;
self.query_needed_blocks().await;
}
async fn handle_frontier_update(&mut self, update: FrontierUpdate) {
let frontier = update.frontier;
let state_frontiers = BlockSyncFrontiers {
finalized_height: frontier.finalized.height,
verified_block_tip: frontier.verified_body.height,
verified_block_hash: frontier.verified_body.hash,
};
match update.change {
FrontierChange::Snapshot => {
self.handle_header_tip_changed(
frontier.best_header.height,
frontier.best_header.hash,
)
.await;
self.handle_state_frontiers_changed(state_frontiers).await;
}
FrontierChange::HeaderAdvanced => {
self.handle_header_tip_changed(
frontier.best_header.height,
frontier.best_header.hash,
)
.await;
if frontier.verified_body.height > self.verified_block_tip {
self.handle_state_frontiers_changed(state_frontiers).await;
}
}
FrontierChange::HeaderReanchored => {
self.state.best_header_tip = frontier.best_header.height;
self.state.best_header_hash = frontier.best_header.hash;
self.handle_chain_tip_reset(state_frontiers, false).await;
}
FrontierChange::VerifiedGrow => {
self.handle_state_frontiers_changed(state_frontiers).await;
if frontier.best_header.height > self.state.best_header_tip {
self.handle_header_tip_changed(
frontier.best_header.height,
frontier.best_header.hash,
)
.await;
}
}
FrontierChange::VerifiedReset => {
self.handle_chain_tip_reset(state_frontiers, true).await;
if frontier.best_header.height > self.state.best_header_tip {
self.handle_header_tip_changed(
frontier.best_header.height,
frontier.best_header.hash,
)
.await;
}
}
}
}
async fn handle_state_frontiers_changed(&mut self, frontiers: BlockSyncFrontiers) {
self.state.finalized_height = self.state.finalized_height.max(frontiers.finalized_height);
if frontiers.verified_block_tip < self.verified_block_tip {
tracing::debug!(
current = ?self.verified_block_tip,
stale = ?frontiers.verified_block_tip,
"ignoring stale Zakura block-sync frontier update"
);
return;
}
let tip = frontiers.verified_block_tip;
let capacity = self.sequencer_input.capacity();
let max_capacity = self.sequencer_input.max_capacity();
let started = Instant::now();
let send_result = self
.sequencer_control
.send(SequencerControlInput::FrontierAdvance {
frontiers,
release_applied: true,
});
self.trace_sequencer_control_send(
"frontier_advance",
if send_result.is_ok() {
"queued"
} else {
"closed"
},
started.elapsed(),
Some(tip),
None,
capacity,
max_capacity,
);
}
fn prune_needed_below_floor(&mut self) {
let floor = self.request_floor;
let before = self.state.needed_heights.len();
self.state.needed_heights.retain(|height| *height > floor);
if self.state.needed_heights.len() != before {
self.publish_candidate_state();
}
}
async fn handle_chain_tip_reset(
&mut self,
frontiers: BlockSyncFrontiers,
preserve_active_successors: bool,
) {
self.pending_needed_query = None;
let tip = frontiers.verified_block_tip;
let peer_has_successor_after = next_height(tip)
.map(|next| self.registry.any_outstanding_at_or_above(next))
.unwrap_or(false);
let peer_outstanding_conflicts_at_tip = self
.registry
.any_outstanding_conflicts_at(tip, frontiers.verified_block_hash);
let capacity = self.sequencer_input.capacity();
let max_capacity = self.sequencer_input.max_capacity();
let started = Instant::now();
let send_result = self
.sequencer_control
.send(SequencerControlInput::FrontierReset {
frontiers,
preserve_active_successors,
peer_has_successor_after,
peer_outstanding_conflicts_at_tip,
});
self.trace_sequencer_control_send(
"frontier_reset",
if send_result.is_ok() {
"queued"
} else {
"closed"
},
started.elapsed(),
Some(tip),
None,
capacity,
max_capacity,
);
}
async fn on_sequencer_view_changed(&mut self, view: SequencerView) {
let reset_advanced = view.reset_epoch != self.last_reset_epoch;
let reaction_advanced = view.reaction_epoch != self.last_reaction_epoch;
let old_serving_tip = (self.state.servable_high, self.state.servable_hash);
let tip_advanced = view.verified_tip > self.verified_block_tip;
self.last_view = view;
self.state.finalized_height = self.state.finalized_height.max(view.finalized);
self.verified_block_tip = view.verified_tip;
self.request_floor = view.download_floor;
self.state.verified_block_hash = view.verified_hash;
self.state.servable_high = view.verified_tip;
self.state.servable_hash = view.verified_hash;
if !reaction_advanced {
return;
}
self.last_reaction_epoch = view.reaction_epoch;
if reset_advanced {
self.last_reset_epoch = view.reset_epoch;
self.trace_chain_tip_reset(view.verified_tip);
self.pending_needed_query = None;
} else if tip_advanced {
self.trace_frontiers_changed(view.verified_tip);
}
self.prune_needed_below_floor();
self.queue_status_refresh_if_changed(old_serving_tip);
self.flush_status_refresh().await;
self.query_needed_blocks_with_options(reset_advanced).await;
}
async fn handle_needed_blocks(&mut self, blocks: Vec<BlockSyncBlockMeta>) {
self.pending_needed_query = None;
let blocks: Vec<_> = blocks
.into_iter()
.filter(|block| {
block.height > self.request_floor
&& !self.state.work_queue.in_flight_contains(block.height)
&& !self
.registry
.has_outstanding_request(block.height, block.hash)
})
.collect();
self.state.needed_heights = blocks.iter().map(|block| block.height).collect();
self.state.needed_heights.sort_unstable();
self.state.needed_heights.dedup();
let count = self.state.work_queue.extend(
blocks
.into_iter()
.map(|block| (block.height, block.hash, block.size)),
);
self.trace_work_extended(count);
self.publish_candidate_state();
}
fn body_lag(&self) -> u32 {
self.state
.best_header_tip
.0
.saturating_sub(self.verified_block_tip.0)
}
async fn handle_routine_message(&mut self, message: RoutineToReactor) {
match message {
RoutineToReactor::StatusReceived { peer, send_reply } => {
self.handle_status_received(peer, send_reply);
}
RoutineToReactor::ServeGetBlocks {
peer,
start_height,
count,
} => {
if self.state.parked_peers.contains(&peer) {
return;
}
self.handle_get_blocks(peer, start_height, count).await;
}
RoutineToReactor::RequeryNeeded => {
self.query_needed_blocks().await;
}
RoutineToReactor::Misbehavior { peer, reason } => {
self.report_misbehavior(peer, reason).await;
}
}
}
fn handle_status_received(&mut self, peer: ZakuraPeerId, send_reply: bool) {
if !self.state.peers.contains_key(&peer) {
return;
}
self.publish_candidate_state();
if send_reply {
self.send_status(&peer, "status_reply");
}
}
async fn handle_get_blocks(
&mut self,
peer: ZakuraPeerId,
start_height: block::Height,
count: u32,
) {
let local_inflight_cap = self.startup.config.advertised_max_inflight_requests();
if !self.state.peers.contains_key(&peer) {
self.report_misbehavior(peer, BlockSyncMisbehavior::GetBlocksSpam)
.await;
return;
}
if !self.registry_received_status(&peer) {
self.report_misbehavior(peer, BlockSyncMisbehavior::GetBlocksSpam)
.await;
return;
}
if count == 0 {
self.report_misbehavior(peer, BlockSyncMisbehavior::GetBlocksTooLong)
.await;
return;
}
let started_serving = self.state.peers.get_mut(&peer).is_some_and(|peer_state| {
peer_state.try_start_serving_blocks(local_inflight_cap, start_height)
});
if !started_serving {
let unavailable_count = count.min(inbound_get_blocks_count_limit(&self.startup.config));
self.send_range_unavailable(&peer, start_height, unavailable_count);
return;
}
let requested_count = self.clamp_served_block_count(start_height, count);
if requested_count == 0 {
let unavailable_count = count.min(inbound_get_blocks_count_limit(&self.startup.config));
self.send_range_unavailable(&peer, start_height, unavailable_count);
self.finish_serving_blocks(&peer, start_height);
return;
}
if !self.dispatch_action(BlockSyncAction::QueryBlocksByHeightRange {
peer: peer.clone(),
start: start_height,
count: requested_count,
}) {
self.finish_serving_blocks(&peer, start_height);
}
}
async fn handle_block_range_response_ready(
&mut self,
peer: ZakuraPeerId,
start_height: block::Height,
requested_count: u32,
blocks: Vec<(block::Height, Arc<block::Block>, usize)>,
) {
let prepare_elapsed = self.serving_blocks_elapsed(&peer, start_height);
let send_started = Instant::now();
let max_response_bytes = u64::from(self.startup.config.advertised_max_response_bytes());
let mut sent_blocks = 0u32;
let mut sent_bytes = 0u64;
let mut reason = "complete";
for (height, block, size) in blocks {
let Ok(size) = u64::try_from(size) else {
reason = "size_overflow";
break;
};
let Some(next_bytes) = sent_bytes.checked_add(size) else {
reason = "byte_overflow";
break;
};
if next_bytes > max_response_bytes {
reason = "byte_cap";
break;
}
if height_after_count(start_height, sent_blocks) != Some(height) {
reason = "non_contiguous";
break;
}
if !self.send_block(&peer, block) {
reason = "send_failed";
break;
}
sent_blocks = sent_blocks.saturating_add(1);
sent_bytes = next_bytes;
}
if sent_blocks == 0 {
self.send_range_unavailable(&peer, start_height, requested_count);
} else {
self.send_blocks_done(&peer, start_height, sent_blocks);
}
let total_elapsed = self.finish_serving_blocks(&peer, start_height);
self.trace_range_response_sent(
&peer,
RangeResponseTrace {
start_height,
requested_count,
sent_count: sent_blocks,
sent_bytes,
reason,
prepare_elapsed,
send_elapsed: send_started.elapsed(),
total_elapsed,
},
);
}
async fn handle_block_range_response_finished(
&mut self,
peer: ZakuraPeerId,
start_height: block::Height,
requested_count: u32,
returned_count: u32,
) {
if returned_count == 0 {
self.send_range_unavailable(&peer, start_height, requested_count);
}
let elapsed = self.finish_serving_blocks(&peer, start_height);
self.trace_range_response_sent(
&peer,
RangeResponseTrace {
start_height,
requested_count,
sent_count: returned_count,
sent_bytes: 0,
reason: "driver_finished",
prepare_elapsed: elapsed,
send_elapsed: Duration::ZERO,
total_elapsed: elapsed,
},
);
}
async fn handle_block_apply_finished(
&mut self,
token: BlockApplyToken,
height: block::Height,
hash: block::Hash,
result: BlockApplyResult,
local_frontier: Option<BlockSyncFrontiers>,
) {
self.trace_apply_finished(height, token, result, self.state.budget.reserved());
let capacity = self.sequencer_input.capacity();
let max_capacity = self.sequencer_input.max_capacity();
let started = Instant::now();
let send_result = self
.sequencer_control
.send(SequencerControlInput::ApplyFinished {
token,
height,
hash,
result,
local_frontier,
});
self.trace_sequencer_control_send(
"apply_finished",
if send_result.is_ok() {
"queued"
} else {
"closed"
},
started.elapsed(),
Some(height),
Some(token),
capacity,
max_capacity,
);
}
fn serving_blocks_elapsed(
&self,
peer: &ZakuraPeerId,
start_height: block::Height,
) -> Option<Duration> {
self.state
.peers
.get(peer)
.and_then(|peer_state| peer_state.serving_blocks_elapsed(start_height))
}
fn finish_serving_blocks(
&mut self,
peer: &ZakuraPeerId,
start_height: block::Height,
) -> Option<Duration> {
if let Some(peer_state) = self.state.peers.get_mut(peer) {
peer_state.finish_serving_blocks(start_height)
} else {
None
}
}
async fn query_needed_blocks(&mut self) -> bool {
self.query_needed_blocks_with_options(false).await
}
async fn query_needed_blocks_with_options(&mut self, force: bool) -> bool {
if !self.startup.state_queries_enabled {
return false;
}
if self.request_floor >= self.state.best_header_tip {
self.pending_needed_query = None;
return true;
}
let Some(from) = self.next_needed_block_query_start() else {
return true;
};
if !force && self.local_body_work_blocks() >= self.refill_low_water_blocks() {
return true;
}
let limit = self.refill_query_limit_blocks(from);
let query = (
from,
limit,
self.state.best_header_tip,
self.state.best_header_hash,
);
if self.pending_needed_query == Some(query) {
return true;
}
let dispatched = self.dispatch_action(BlockSyncAction::QueryNeededBlocks {
from,
limit,
best_header_tip: self.state.best_header_tip,
});
if dispatched {
self.pending_needed_query = Some(query);
}
dispatched
}
fn next_needed_block_query_start(&self) -> Option<block::Height> {
let last_claimed = self
.state
.work_queue
.max_claimed()
.unwrap_or(self.request_floor);
if last_claimed >= self.state.best_header_tip {
return None;
}
last_claimed.next().ok()
}
fn refill_query_limit_blocks(&self, from: block::Height) -> u32 {
let remaining = self
.state
.best_header_tip
.0
.saturating_sub(from.0)
.saturating_add(1);
let fanout_window = self.refill_low_water_blocks().saturating_mul(2).max(1);
let fanout_window = u32::try_from(fanout_window).unwrap_or(u32::MAX);
remaining.min(fanout_window).min(NEEDED_BLOCK_REFILL_LIMIT)
}
fn local_body_work_blocks(&self) -> usize {
let outstanding = self.registry.total_unreceived();
self.state
.work_queue
.pending_len()
.saturating_add(outstanding)
}
fn refill_low_water_blocks(&self) -> usize {
let status_peers = self.registry.peers_with_status().max(1);
let max_blocks_per_response =
usize::try_from(self.startup.config.advertised_max_blocks_per_response())
.expect("advertised block count fits usize");
let max_inflight_per_peer =
usize::try_from(self.startup.config.advertised_max_inflight_requests())
.expect("u32 max inflight requests fits in usize on supported targets")
.min(EFFECTIVE_BS_OUTBOUND_INFLIGHT_PER_PEER);
status_peers
.saturating_mul(max_inflight_per_peer)
.saturating_mul(max_blocks_per_response)
.max(max_blocks_per_response)
}
fn send_status(&self, peer: &ZakuraPeerId, reason: &'static str) -> bool {
let Some(peer_state) = self.state.peers.get(peer) else {
return false;
};
let status = self.local_status();
let msg = BlockSyncMessage::Status(status);
let started = Instant::now();
let session = peer_state.session.clone();
match session.try_send_status(status) {
Ok(()) => {
self.trace_message_sent(peer, &msg, "queued", started.elapsed());
self.trace_status_sent(peer, reason, status);
true
}
Err(OrderedSendError::Full) => {
tracing::debug!(?peer, "Zakura block-sync Status queue is full");
self.trace_message_sent(peer, &msg, "full", started.elapsed());
self.trace_queue_send_failed(
peer,
&msg,
&OrderedSendError::Full,
session.outbound_capacity(),
session.outbound_max_capacity(),
Some(reason),
);
false
}
Err(error) => {
tracing::debug!(?peer, ?error, "failed to queue Zakura block-sync Status");
self.trace_status_send_failed(peer, reason);
self.trace_message_sent(peer, &msg, "error", started.elapsed());
self.trace_queue_send_failed(
peer,
&msg,
&error,
session.outbound_capacity(),
session.outbound_max_capacity(),
Some(reason),
);
session.cancel_token().cancel();
false
}
}
}
fn send_status_and_mark_refresh(
&mut self,
peer: &ZakuraPeerId,
reason: &'static str,
now: Instant,
) -> bool {
if !self.send_status(peer, reason) {
return false;
}
if let Some(peer_state) = self.state.peers.get_mut(peer) {
peer_state.refresh_meter.mark_taken(now);
}
true
}
fn send_block(&self, peer: &ZakuraPeerId, block: Arc<block::Block>) -> bool {
let Some(session) = self
.state
.peers
.get(peer)
.map(|peer_state| peer_state.session.clone())
else {
return false;
};
let msg = BlockSyncMessage::Block(block.clone());
let started = Instant::now();
match session.try_send_block(block) {
Ok(()) => {
metrics::counter!("sync.block.body.served").increment(1);
self.trace_message_sent(peer, &msg, "queued", started.elapsed());
true
}
Err(OrderedSendError::Full) => {
metrics::counter!("sync.block.body.serve_queue_full").increment(1);
tracing::debug!(?peer, "Zakura block-sync Block queue is full");
self.trace_message_sent(peer, &msg, "full", started.elapsed());
self.trace_queue_send_failed(
peer,
&msg,
&OrderedSendError::Full,
session.outbound_capacity(),
session.outbound_max_capacity(),
None,
);
false
}
Err(error) => {
tracing::debug!(?peer, ?error, "failed to queue Zakura block-sync Block");
self.trace_message_sent(peer, &msg, "error", started.elapsed());
self.trace_queue_send_failed(
peer,
&msg,
&error,
session.outbound_capacity(),
session.outbound_max_capacity(),
None,
);
session.cancel_token().cancel();
false
}
}
}
fn send_blocks_done(&self, peer: &ZakuraPeerId, start_height: block::Height, returned: u32) {
if returned == 0 {
return;
}
let Some(session) = self
.state
.peers
.get(peer)
.map(|peer_state| peer_state.session.clone())
else {
return;
};
let msg = BlockSyncMessage::BlocksDone {
start_height,
returned,
};
let started = Instant::now();
match session.try_send_blocks_done(start_height, returned) {
Ok(()) => self.trace_message_sent(peer, &msg, "queued", started.elapsed()),
Err(OrderedSendError::Full) => {
metrics::counter!("sync.block.done.serve_queue_full").increment(1);
tracing::debug!(?peer, "Zakura block-sync BlocksDone queue is full");
self.trace_message_sent(peer, &msg, "full", started.elapsed());
self.trace_queue_send_failed(
peer,
&msg,
&OrderedSendError::Full,
session.outbound_capacity(),
session.outbound_max_capacity(),
None,
);
}
Err(error) => {
tracing::debug!(
?peer,
?error,
"failed to queue Zakura block-sync BlocksDone"
);
self.trace_message_sent(peer, &msg, "error", started.elapsed());
self.trace_queue_send_failed(
peer,
&msg,
&error,
session.outbound_capacity(),
session.outbound_max_capacity(),
None,
);
session.cancel_token().cancel();
}
}
}
fn send_range_unavailable(&self, peer: &ZakuraPeerId, start_height: block::Height, count: u32) {
let count = count.max(1);
let Some(peer_state) = self.state.peers.get(peer) else {
return;
};
let msg = BlockSyncMessage::RangeUnavailable {
start_height,
count,
};
let started = Instant::now();
match peer_state
.session
.try_send_range_unavailable(start_height, count)
{
Ok(()) => self.trace_message_sent(peer, &msg, "queued", started.elapsed()),
Err(OrderedSendError::Full) => {
metrics::counter!("sync.block.unavailable.serve_queue_full").increment(1);
tracing::debug!(?peer, "Zakura block-sync RangeUnavailable queue is full");
self.trace_message_sent(peer, &msg, "full", started.elapsed());
self.trace_queue_send_failed(
peer,
&msg,
&OrderedSendError::Full,
peer_state.session.outbound_capacity(),
peer_state.session.outbound_max_capacity(),
None,
);
}
Err(error) => {
tracing::debug!(
?peer,
?error,
"failed to queue Zakura block-sync RangeUnavailable"
);
self.trace_message_sent(peer, &msg, "error", started.elapsed());
self.trace_queue_send_failed(
peer,
&msg,
&error,
peer_state.session.outbound_capacity(),
peer_state.session.outbound_max_capacity(),
None,
);
peer_state.session.cancel_token().cancel();
}
}
}
async fn flush_status_refresh(&mut self) {
let unready: HashSet<ZakuraPeerId> = self
.registry
.candidate_snapshot()
.into_iter()
.filter_map(|(peer, received_status, _, _)| (!received_status).then_some(peer))
.collect();
let has_unready_peers = !unready.is_empty();
if !self.state.pending_status_refresh && !has_unready_peers {
return;
}
let now = Instant::now();
let status = self.local_status();
let status_changed = self.state.pending_status_refresh
&& status != self.state.last_advertised_status
&& self.state.status_refresh.try_take(now);
self.state.pending_status_refresh = false;
if status_changed {
self.state.last_advertised_status = status;
let _ = self.status.send(status);
}
let peer_ids: Vec<_> = self
.state
.peers
.iter()
.filter_map(|(peer_id, peer)| {
let should_send_status = status_changed
|| (unready.contains(peer_id) && peer.refresh_meter.is_ready(now));
if should_send_status {
Some(peer_id.clone())
} else {
None
}
})
.collect();
for peer in peer_ids {
self.send_status_and_mark_refresh(&peer, "refresh", now);
}
}
fn queue_status_refresh_if_changed(&mut self, old_serving_tip: (block::Height, block::Hash)) {
if old_serving_tip != (self.state.servable_high, self.state.servable_hash)
&& self.local_status() != self.state.last_advertised_status
{
self.state.pending_status_refresh = true;
}
}
fn emit_trace(
&self,
event: &'static str,
build: impl FnOnce(&mut serde_json::Map<String, serde_json::Value>),
) {
self.startup.trace.emit_with(BLOCK_SYNC_TABLE, |row| {
row.insert(
bs_trace::EVENT.to_string(),
serde_json::Value::String(event.to_string()),
);
build(row);
});
}
fn refresh_throughput(&mut self) {
let now = Instant::now();
if let Ok(mut meter) = self.state.received_throughput.lock() {
meter.sample(now);
}
}
fn trace_sync_state(&self, include_diagnostics: bool) {
if !self.startup.trace.is_enabled() {
return;
}
let floor_gap = include_diagnostics
.then(|| self.floor_gap_diagnostics(Instant::now()))
.flatten();
let slots = include_diagnostics.then(|| self.registry.slot_summary());
let counts = include_diagnostics.then(|| self.registry.direction_status_counts());
let peers_with_status = include_diagnostics.then(|| self.registry.peers_with_status());
let view = self.last_view;
let submitted_applies = view.in_flight_submission_count;
let (received_bytes_per_sec, received_blocks_per_sec) = self
.state
.received_throughput
.lock()
.map(|meter| (meter.bytes_per_sec(), meter.blocks_per_sec()))
.unwrap_or((0, 0));
self.emit_trace(bs_trace::BLOCK_SYNC_STATE, |row| {
bs_insert_height(row, bs_trace::REQUEST_FLOOR, self.request_floor);
bs_insert_height(row, bs_trace::BODY_DOWNLOAD_FLOOR, view.download_floor);
bs_insert_height(row, bs_trace::VERIFIED_BLOCK_TIP, view.verified_tip);
bs_insert_height(row, bs_trace::BEST_HEADER_TIP, self.state.best_header_tip);
bs_insert_u64(row, bs_trace::BODY_LAG, u64::from(self.body_lag()));
bs_insert_u64(row, bs_trace::APPLYING, view.applying_len);
bs_insert_u64(row, bs_trace::SUBMITTED_APPLIES, submitted_applies);
bs_insert_u64(row, bs_trace::REORDER, view.reorder_len);
if let Some(slots) = slots {
bs_insert_u64(
row,
bs_trace::OUTSTANDING,
slots.outstanding_requests as u64,
);
}
if let Some(floor_gap) = floor_gap {
bs_insert_height(row, bs_trace::FLOOR_GAP_HEIGHT, floor_gap.height);
bs_insert_str(row, bs_trace::FLOOR_GAP_STATE, floor_gap.state);
bs_insert_u64(
row,
bs_trace::FLOOR_GAP_SERVABLE_PEERS,
floor_gap.servable_peers as u64,
);
bs_insert_u64(
row,
bs_trace::FLOOR_GAP_AVAILABLE_PEERS,
floor_gap.available_peers as u64,
);
bs_insert_u64(
row,
bs_trace::FLOOR_GAP_OUTSTANDING_PEERS,
floor_gap.outstanding_peers as u64,
);
if let Some(age) = floor_gap.oldest_outstanding_ms {
bs_insert_u64(row, bs_trace::FLOOR_GAP_OLDEST_OUTSTANDING_MS, age);
}
if let Some(deadline) = floor_gap.next_deadline_ms {
bs_insert_u64(row, bs_trace::FLOOR_GAP_NEXT_DEADLINE_MS, deadline);
}
}
bs_insert_u64(
row,
bs_trace::BUDGET_AVAILABLE,
self.state.budget.available(),
);
bs_insert_u64(row, bs_trace::BUDGET_RESERVED, self.state.budget.reserved());
let sequencer_input_queued_bytes = self
.sequencer_input_bytes
.load(std::sync::atomic::Ordering::Relaxed);
let sequencer_input_decoded_attributed_memory_bytes = self
.sequencer_input_decoded_attributed_memory_bytes
.load(std::sync::atomic::Ordering::Relaxed);
let sequencer_input_max_capacity = self.sequencer_input.max_capacity();
let sequencer_input_capacity = self.sequencer_input.capacity();
let sequencer_input_queued_blocks =
sequencer_input_max_capacity.saturating_sub(sequencer_input_capacity);
bs_insert_u64(
row,
"sequencer_input_queued_bytes",
sequencer_input_queued_bytes,
);
bs_insert_u64(
row,
bs_trace::SEQUENCER_INPUT_DECODED_ATTRIBUTED_MEMORY_BYTES,
sequencer_input_decoded_attributed_memory_bytes,
);
bs_insert_u64(
row,
bs_trace::REORDER_DECODED_ATTRIBUTED_MEMORY_BYTES,
view.reorder_decoded_attributed_memory_bytes,
);
bs_insert_u64(
row,
bs_trace::APPLYING_DECODED_ATTRIBUTED_MEMORY_BYTES,
view.applying_decoded_attributed_memory_bytes,
);
bs_insert_u64(
row,
bs_trace::ACTIVE_PIPELINE_DECODED_ATTRIBUTED_MEMORY_BYTES,
sequencer_input_decoded_attributed_memory_bytes
.saturating_add(view.reorder_decoded_attributed_memory_bytes)
.saturating_add(view.applying_decoded_attributed_memory_bytes),
);
bs_insert_u64(
row,
"sequencer_input_queued_blocks",
sequencer_input_queued_blocks as u64,
);
bs_insert_u64(
row,
"sequencer_input_capacity",
sequencer_input_capacity as u64,
);
bs_insert_u64(
row,
"sequencer_input_max_capacity",
sequencer_input_max_capacity as u64,
);
bs_insert_u64(row, "reorder_buffered_bytes", view.reorder_buffered_bytes);
bs_insert_u64(row, "applying_buffered_bytes", view.applying_buffered_bytes);
bs_insert_u64(
row,
"unsubmitted_applying_count",
view.unsubmitted_applying_count,
);
bs_insert_u64(
row,
"in_flight_submission_bytes",
view.in_flight_submission_bytes,
);
bs_insert_u64(
row,
"retained_pipeline_wire_bytes",
super::admission::RetainedPipelineBytes {
reorder_buffered_bytes: view.reorder_buffered_bytes,
applying_buffered_bytes: view.applying_buffered_bytes,
sequencer_input_queued_bytes,
}
.wire_bytes(),
);
bs_insert_u64(row, bs_trace::PEERS, self.state.peers.len() as u64);
if let Some(peers_with_status) = peers_with_status {
bs_insert_u64(row, bs_trace::PEERS_WITH_STATUS, peers_with_status as u64);
}
if let (Some(slots), Some(peers_with_status)) = (slots, peers_with_status) {
let peers_wanting_slots = peers_with_status.saturating_sub(slots.saturated_peers);
let download_blocked_on_budget = u64::from(
peers_wanting_slots > 0
&& self.state.budget.available() < BS_PER_BLOCK_WORST_CASE_BYTES,
);
bs_insert_u64(
row,
bs_trace::PEERS_WANTING_SLOTS,
peers_wanting_slots as u64,
);
bs_insert_u64(
row,
bs_trace::DOWNLOAD_BLOCKED_ON_BUDGET,
download_blocked_on_budget,
);
}
bs_insert_u64(
row,
bs_trace::RECEIVED_BYTES_PER_SEC,
received_bytes_per_sec,
);
bs_insert_u64(
row,
bs_trace::RECEIVED_BLOCKS_PER_SEC,
received_blocks_per_sec,
);
bs_insert_u64(
row,
bs_trace::COMMITTED_BYTES_PER_SEC,
view.committed_bytes_per_sec,
);
bs_insert_u64(
row,
bs_trace::COMMITTED_BLOCKS_PER_SEC,
view.committed_blocks_per_sec,
);
if let Some(counts) = counts {
bs_insert_u64(row, "inbound_peers", counts.inbound as u64);
bs_insert_u64(row, "outbound_peers", counts.outbound as u64);
bs_insert_u64(
row,
"inbound_peers_with_status",
counts.inbound_with_status as u64,
);
bs_insert_u64(
row,
"outbound_peers_with_status",
counts.outbound_with_status as u64,
);
}
if let Some(slots) = slots {
bs_insert_u64(row, "request_slot_capacity", slots.capacity as u64);
bs_insert_u64(
row,
"request_slot_effective_window",
slots.effective_window as u64,
);
bs_insert_u64(row, "request_slot_available", slots.available as u64);
bs_insert_u64(
row,
"request_slot_saturated_peers",
slots.saturated_peers as u64,
);
}
if include_diagnostics {
if let Some(min) = self.state.needed_heights.first() {
bs_insert_height(row, bs_trace::NEEDED_MIN, *min);
}
bs_insert_u64(
row,
bs_trace::NEEDED_COUNT,
self.state.needed_heights.len() as u64,
);
bs_insert_u64(
row,
bs_trace::QUEUE_LEN,
self.state.work_queue.pending_run_count() as u64,
);
bs_insert_u64(
row,
bs_trace::QUEUE_BLOCKS,
self.state.work_queue.pending_len() as u64,
);
if let Some(start) = self.state.work_queue.min_pending() {
bs_insert_height(row, bs_trace::QUEUE_MIN_START, start);
}
bs_insert_u64(
row,
bs_trace::ASSIGNED_LEN,
self.state.work_queue.in_flight_len() as u64,
);
bs_insert_u64(
row,
bs_trace::LOCAL_BODY_WORK,
self.local_body_work_blocks() as u64,
);
bs_insert_u64(
row,
bs_trace::REFILL_LOW_WATER,
self.refill_low_water_blocks() as u64,
);
if let Some(end) = self.state.work_queue.max_in_flight() {
bs_insert_height(row, bs_trace::COVERED_MAX_END, end);
}
}
});
}
fn trace_status_sent(
&self,
peer: &ZakuraPeerId,
reason: &'static str,
status: BlockSyncStatus,
) {
self.emit_trace(bs_trace::BLOCK_STATUS_SENT, |row| {
bs_insert_peer(row, bs_trace::PEER, peer);
bs_insert_str(row, bs_trace::REASON, reason);
bs_insert_height(row, bs_trace::RANGE_START, status.servable_low);
bs_insert_height(row, bs_trace::HEIGHT, status.servable_high);
});
}
fn trace_status_send_failed(&self, peer: &ZakuraPeerId, reason: &'static str) {
self.emit_trace(bs_trace::BLOCK_STATUS_SEND_FAILED, |row| {
bs_insert_peer(row, bs_trace::PEER, peer);
bs_insert_str(row, bs_trace::REASON, reason);
});
}
fn trace_queue_send_failed(
&self,
peer: &ZakuraPeerId,
msg: &BlockSyncMessage,
error: &OrderedSendError,
queue_capacity: usize,
queue_max_capacity: usize,
reason: Option<&'static str>,
) {
self.startup.trace.emit_with(QUEUE_SEND_TABLE, |row| {
bs_insert_str(row, qs_trace::EVENT, qs_trace::QUEUE_SEND_FAILED);
bs_insert_str(row, qs_trace::SERVICE, "block_sync");
bs_insert_str(row, qs_trace::MESSAGE, block_sync_message_label(msg));
bs_insert_peer(row, qs_trace::PEER, peer);
bs_insert_str(row, qs_trace::ERROR, ordered_send_error_label(error));
if let Some(reason) = reason {
bs_insert_str(row, qs_trace::REASON, reason);
}
bs_insert_u64(
row,
qs_trace::QUEUE_CAPACITY,
u64::try_from(queue_capacity).unwrap_or(u64::MAX),
);
bs_insert_u64(
row,
qs_trace::QUEUE_MAX_CAPACITY,
u64::try_from(queue_max_capacity).unwrap_or(u64::MAX),
);
trace_block_sync_message_fields(row, msg);
});
}
fn trace_peer_connected(
&self,
peer: &ZakuraPeerId,
direction: ServicePeerDirection,
active_connections: usize,
) {
self.emit_trace(bs_trace::BLOCK_PEER_CONNECTED, |row| {
bs_insert_peer(row, bs_trace::PEER, peer);
bs_insert_str(row, "direction", direction.trace_label());
bs_insert_u64(
row,
bs_trace::ACTIVE_CONNECTIONS,
u64::try_from(active_connections).unwrap_or(u64::MAX),
);
});
}
fn trace_peer_disconnected(
&self,
peer: &ZakuraPeerId,
received_status: bool,
active_connections: usize,
) {
self.emit_trace(bs_trace::BLOCK_PEER_DISCONNECTED, |row| {
bs_insert_peer(row, bs_trace::PEER, peer);
row.insert(
"received_status".to_string(),
serde_json::Value::Bool(received_status),
);
bs_insert_u64(
row,
bs_trace::ACTIVE_CONNECTIONS,
u64::try_from(active_connections).unwrap_or(u64::MAX),
);
});
}
fn trace_message_sent(
&self,
peer: &ZakuraPeerId,
msg: &BlockSyncMessage,
result: &'static str,
elapsed: Duration,
) {
self.emit_trace(bs_trace::BLOCK_MESSAGE_SENT, |row| {
bs_insert_peer(row, bs_trace::PEER, peer);
bs_insert_str(row, bs_trace::KIND, block_sync_message_label(msg));
bs_insert_str(row, bs_trace::RESULT, result);
bs_insert_duration_ms(row, bs_trace::ELAPSED_MS, elapsed);
trace_block_sync_message_fields(row, msg);
});
}
fn trace_apply_finished(
&self,
height: block::Height,
token: BlockApplyToken,
result: BlockApplyResult,
budget_reserved_after: u64,
) {
self.emit_trace(bs_trace::BLOCK_APPLY_FINISHED, |row| {
bs_insert_height(row, bs_trace::HEIGHT, height);
bs_insert_u64(row, bs_trace::APPLY_TOKEN, token);
bs_insert_str(row, bs_trace::RESULT, block_apply_result_label(result));
bs_insert_u64(row, bs_trace::BUDGET_RESERVED_AFTER, budget_reserved_after);
});
}
#[allow(clippy::too_many_arguments)]
fn trace_sequencer_control_send(
&self,
kind: &'static str,
result: &'static str,
elapsed: Duration,
height: Option<block::Height>,
token: Option<BlockApplyToken>,
capacity: usize,
max_capacity: usize,
) {
self.emit_trace(bs_trace::BLOCK_SEQUENCER_CONTROL_SENT, |row| {
bs_insert_str(row, bs_trace::KIND, kind);
bs_insert_str(row, bs_trace::RESULT, result);
bs_insert_duration_ms(row, bs_trace::ELAPSED_MS, elapsed);
if let Some(height) = height {
bs_insert_height(row, bs_trace::HEIGHT, height);
}
if let Some(token) = token {
bs_insert_u64(row, bs_trace::APPLY_TOKEN, token);
}
bs_insert_u64(
row,
"sequencer_input_capacity",
u64::try_from(capacity).unwrap_or(u64::MAX),
);
bs_insert_u64(
row,
"sequencer_input_max_capacity",
u64::try_from(max_capacity).unwrap_or(u64::MAX),
);
});
}
fn trace_range_response_sent(&self, peer: &ZakuraPeerId, response: RangeResponseTrace) {
self.emit_trace(bs_trace::BLOCK_RANGE_RESPONSE_SENT, |row| {
bs_insert_peer(row, bs_trace::PEER, peer);
bs_insert_height(row, bs_trace::RANGE_START, response.start_height);
bs_insert_u64(row, bs_trace::RANGE_COUNT, u64::from(response.sent_count));
bs_insert_u64(
row,
bs_trace::EXPECTED_COUNT,
u64::from(response.requested_count),
);
bs_insert_u64(row, bs_trace::SERIALIZED_BYTES, response.sent_bytes);
bs_insert_str(row, bs_trace::REASON, response.reason);
if let Some(prepare_elapsed) = response.prepare_elapsed {
bs_insert_duration_ms(row, bs_trace::PREPARE_ELAPSED_MS, prepare_elapsed);
}
bs_insert_duration_ms(row, bs_trace::SEND_ELAPSED_MS, response.send_elapsed);
if let Some(total_elapsed) = response.total_elapsed {
bs_insert_duration_ms(row, bs_trace::ELAPSED_MS, total_elapsed);
}
});
}
fn trace_work_extended(&self, inserted: usize) {
if !self.startup.trace.is_enabled() {
return;
}
self.emit_trace(bs_trace::BLOCK_WORK_EXTENDED, |row| {
bs_insert_u64(row, bs_trace::RANGE_COUNT, inserted as u64);
bs_insert_u64(
row,
bs_trace::QUEUE_BLOCKS,
self.state.work_queue.pending_len() as u64,
);
});
}
fn trace_frontiers_changed(&self, verified_block_tip: block::Height) {
self.emit_trace(bs_trace::BLOCK_FRONTIERS_CHANGED, |row| {
bs_insert_height(row, bs_trace::VERIFIED_BLOCK_TIP, verified_block_tip);
bs_insert_height(row, bs_trace::BEST_HEADER_TIP, self.state.best_header_tip);
});
}
fn trace_chain_tip_reset(&self, verified_block_tip: block::Height) {
self.emit_trace(bs_trace::BLOCK_CHAIN_TIP_RESET, |row| {
bs_insert_height(row, bs_trace::VERIFIED_BLOCK_TIP, verified_block_tip);
});
}
fn floor_gap_diagnostics(&self, now: Instant) -> Option<FloorGapDiagnostics> {
let height = next_height(self.request_floor)?;
if height > self.state.best_header_tip {
return None;
}
let (servable_peers, outstanding_peers) = self.registry.floor_gap_servable(height);
let claims = self.registry.outstanding_claims_at(height);
let available_peers = 0usize;
let oldest_outstanding_ms = claims
.iter()
.map(|claim| elapsed_ms_u64(now.saturating_duration_since(claim.meta.queued_at)))
.max();
let next_deadline_ms = claims
.iter()
.map(|claim| elapsed_ms_u64(claim.meta.deadline.saturating_duration_since(now)))
.min();
let state = if outstanding_peers > 0 {
"outstanding"
} else if self.state.work_queue.pending_contains(height) {
"queued"
} else if self.state.work_queue.in_flight_contains(height) {
"in_flight_without_outstanding"
} else if self.state.needed_heights.binary_search(&height).is_ok() {
"needed_unscheduled"
} else {
"absent"
};
Some(FloorGapDiagnostics {
height,
state,
servable_peers,
available_peers,
outstanding_peers,
oldest_outstanding_ms,
next_deadline_ms,
})
}
fn publish_metrics(&self) {
let view = *self.sequencer_view.borrow();
let sequencer_input_decoded_attributed_memory_bytes = self
.sequencer_input_decoded_attributed_memory_bytes
.load(std::sync::atomic::Ordering::Relaxed);
metrics::gauge!("sync.block.best_header_tip.height")
.set(self.state.best_header_tip.0 as f64);
metrics::gauge!("sync.block.verified_tip.height").set(self.verified_block_tip.0 as f64);
metrics::gauge!("sync.block.missing_bodies").set(self.state.needed_heights.len() as f64);
metrics::gauge!("sync.block.budget.reserved_bytes")
.set(self.state.budget.reserved() as f64);
metrics::gauge!("sync.block.reorder.buffered_bytes")
.set(self.last_view.reorder_buffered_bytes as f64);
metrics::gauge!("sync.block.sequencer_input.decoded.attributed_memory_bytes")
.set(sequencer_input_decoded_attributed_memory_bytes as f64);
metrics::gauge!("sync.block.reorder.decoded.attributed_memory_bytes")
.set(view.reorder_decoded_attributed_memory_bytes as f64);
metrics::gauge!("sync.block.applying.decoded.attributed_memory_bytes")
.set(view.applying_decoded_attributed_memory_bytes as f64);
metrics::gauge!("sync.block.active_pipeline.decoded.attributed_memory_bytes").set(
sequencer_input_decoded_attributed_memory_bytes
.saturating_add(view.reorder_decoded_attributed_memory_bytes)
.saturating_add(view.applying_decoded_attributed_memory_bytes) as f64,
);
metrics::gauge!("sync.block.applying").set(self.last_view.applying_len as f64);
metrics::gauge!("sync.block.outstanding").set(self.registry.total_unreceived() as f64);
}
fn clamp_served_block_count(&self, start_height: block::Height, count: u32) -> u32 {
if start_height > self.state.servable_high {
return 0;
}
let available = self
.state
.servable_high
.0
.checked_sub(start_height.0)
.and_then(|diff| diff.checked_add(1))
.unwrap_or(0);
count
.min(inbound_get_blocks_count_limit(&self.startup.config))
.min(available)
}
fn local_status(&self) -> BlockSyncStatus {
BlockSyncStatus {
servable_low: block::Height::MIN,
servable_high: self.state.servable_high,
tip_hash: self.state.servable_hash,
max_blocks_per_response: self.startup.config.advertised_max_blocks_per_response(),
max_inflight_requests: self.startup.config.advertised_max_inflight_requests(),
max_response_bytes: self.startup.config.advertised_max_response_bytes(),
}
}
fn dispatch_action(&self, action: BlockSyncAction) -> bool {
let action_label = action.metric_label();
let queue_depth = self
.actions
.max_capacity()
.saturating_sub(self.actions.capacity());
metrics::histogram!(
"sync.block.action.queue.depth",
"action" => action_label
)
.record(queue_depth as f64);
self.trace_action_dispatched(&action);
match self.actions.try_send(action) {
Ok(()) => true,
Err(mpsc::error::TrySendError::Full(_)) => {
metrics::counter!(
"sync.block.action.send_queue_full",
"action" => action_label
)
.increment(1);
false
}
Err(mpsc::error::TrySendError::Closed(_)) => false,
}
}
async fn report_misbehavior(&mut self, peer: ZakuraPeerId, reason: BlockSyncMisbehavior) {
metrics::counter!("sync.block.peer.violation").increment(1);
let action = BlockSyncAction::Misbehavior { peer, reason };
if !self.dispatch_action(action) {
metrics::counter!("sync.block.peer.violation.action_dropped").increment(1);
}
}
fn trace_event_received(&self, event: &BlockSyncEvent) {
self.emit_trace(bs_trace::BLOCK_EVENT_RECEIVED, |row| match event {
BlockSyncEvent::PeerConnected(session) => {
bs_insert_str(row, bs_trace::KIND, "peer_connected");
bs_insert_peer(row, bs_trace::PEER, session.peer_id());
}
BlockSyncEvent::PeerDisconnected(peer) => {
bs_insert_str(row, bs_trace::KIND, "peer_disconnected");
bs_insert_peer(row, bs_trace::PEER, peer);
}
BlockSyncEvent::HeaderTipChanged { height, hash } => {
bs_insert_str(row, bs_trace::KIND, "header_tip_changed");
bs_insert_height(row, bs_trace::HEIGHT, *height);
bs_insert_hash(row, bs_trace::HASH, *hash);
}
BlockSyncEvent::StateFrontiersChanged(frontiers) => {
bs_insert_str(row, bs_trace::KIND, "state_frontiers_changed");
bs_insert_frontiers(row, frontiers);
}
BlockSyncEvent::ChainTipGrow(frontiers) => {
bs_insert_str(row, bs_trace::KIND, "chain_tip_grow");
bs_insert_frontiers(row, frontiers);
}
BlockSyncEvent::ChainTipReset(frontiers) => {
bs_insert_str(row, bs_trace::KIND, "chain_tip_reset");
bs_insert_frontiers(row, frontiers);
}
BlockSyncEvent::NeededBlocks(blocks) => {
bs_insert_str(row, bs_trace::KIND, "needed_blocks");
bs_insert_u64(row, bs_trace::RANGE_COUNT, blocks.len() as u64);
if let Some(first) = blocks.first() {
bs_insert_height(row, bs_trace::RANGE_START, first.height);
}
}
BlockSyncEvent::BlockApplyFinished {
token,
height,
hash,
result,
local_frontier,
} => {
bs_insert_str(row, bs_trace::KIND, "block_apply_finished");
bs_insert_u64(row, bs_trace::APPLY_TOKEN, *token);
bs_insert_height(row, bs_trace::HEIGHT, *height);
bs_insert_hash(row, bs_trace::HASH, *hash);
bs_insert_str(row, bs_trace::RESULT, block_apply_result_label(*result));
if let Some(frontiers) = local_frontier {
bs_insert_frontiers(row, frontiers);
}
}
BlockSyncEvent::BlockRangeResponseFinished {
peer,
start_height,
requested_count,
returned_count,
} => {
bs_insert_str(row, bs_trace::KIND, "block_range_response_finished");
bs_insert_peer(row, bs_trace::PEER, peer);
bs_insert_height(row, bs_trace::RANGE_START, *start_height);
bs_insert_u64(row, bs_trace::RANGE_COUNT, u64::from(*returned_count));
bs_insert_u64(row, bs_trace::EXPECTED_COUNT, u64::from(*requested_count));
}
BlockSyncEvent::BlockRangeResponseReady {
peer,
start_height,
requested_count,
blocks,
} => {
bs_insert_str(row, bs_trace::KIND, "block_range_response_ready");
bs_insert_peer(row, bs_trace::PEER, peer);
bs_insert_height(row, bs_trace::RANGE_START, *start_height);
bs_insert_u64(row, bs_trace::RANGE_COUNT, blocks.len() as u64);
bs_insert_u64(row, bs_trace::EXPECTED_COUNT, u64::from(*requested_count));
}
});
}
fn trace_action_dispatched(&self, action: &BlockSyncAction) {
self.emit_trace(bs_trace::BLOCK_ACTION_DISPATCHED, |row| match action {
BlockSyncAction::QueryNeededBlocks {
from,
limit,
best_header_tip,
} => {
bs_insert_str(row, bs_trace::KIND, "query_needed_blocks");
bs_insert_height(row, bs_trace::RANGE_START, *from);
bs_insert_u64(row, bs_trace::RANGE_COUNT, u64::from(*limit));
bs_insert_height(row, bs_trace::BEST_HEADER_TIP, *best_header_tip);
}
BlockSyncAction::QueryBlocksByHeightRange { peer, start, count } => {
bs_insert_str(row, bs_trace::KIND, "query_blocks_by_height_range");
bs_insert_peer(row, bs_trace::PEER, peer);
bs_insert_height(row, bs_trace::RANGE_START, *start);
bs_insert_u64(row, bs_trace::RANGE_COUNT, u64::from(*count));
}
BlockSyncAction::SubmitBlock { token, block } => {
bs_insert_str(row, bs_trace::KIND, "submit_block");
bs_insert_u64(row, bs_trace::APPLY_TOKEN, *token);
bs_insert_hash(row, bs_trace::HASH, block.hash());
if let Some(height) = block.coinbase_height() {
bs_insert_height(row, bs_trace::HEIGHT, height);
}
}
BlockSyncAction::Misbehavior { peer, reason } => {
bs_insert_str(row, bs_trace::KIND, "misbehavior");
bs_insert_peer(row, bs_trace::PEER, peer);
bs_insert_str(row, bs_trace::REASON, block_misbehavior_label(*reason));
}
});
}
}
pub(super) fn node_id_from_block_peer_id(peer_id: &ZakuraPeerId) -> Option<NodeId> {
let bytes: [u8; 32] = peer_id.as_bytes().try_into().ok()?;
NodeId::from_bytes(&bytes).ok()
}
fn block_apply_result_label(result: BlockApplyResult) -> &'static str {
match result {
BlockApplyResult::Committed => "committed",
BlockApplyResult::Duplicate => "duplicate",
BlockApplyResult::Rejected => "rejected",
BlockApplyResult::TimedOut => "timed_out",
}
}
pub(super) fn bs_insert_peer(
row: &mut serde_json::Map<String, serde_json::Value>,
key: &'static str,
peer: &ZakuraPeerId,
) {
row.insert(
key.to_string(),
serde_json::Value::String(trace_peer_label(peer)),
);
}
pub(super) fn bs_insert_height(
row: &mut serde_json::Map<String, serde_json::Value>,
key: &'static str,
height: block::Height,
) {
bs_insert_u64(row, key, u64::from(height.0));
}
fn bs_insert_hash(
row: &mut serde_json::Map<String, serde_json::Value>,
key: &'static str,
hash: block::Hash,
) {
row.insert(
key.to_string(),
serde_json::Value::String(format!("{hash}")),
);
}
pub(super) fn bs_insert_u64(
row: &mut serde_json::Map<String, serde_json::Value>,
key: &'static str,
value: u64,
) {
row.insert(key.to_string(), serde_json::Value::from(value));
}
fn bs_insert_duration_ms(
row: &mut serde_json::Map<String, serde_json::Value>,
key: &'static str,
duration: Duration,
) {
bs_insert_u64(
row,
key,
u64::try_from(duration.as_millis()).unwrap_or(u64::MAX),
);
}
fn bs_insert_frontiers(
row: &mut serde_json::Map<String, serde_json::Value>,
frontiers: &BlockSyncFrontiers,
) {
bs_insert_height(
row,
bs_trace::VERIFIED_BLOCK_TIP,
frontiers.verified_block_tip,
);
bs_insert_hash(row, bs_trace::HASH, frontiers.verified_block_hash);
}
fn trace_block_sync_message_fields(
row: &mut serde_json::Map<String, serde_json::Value>,
msg: &BlockSyncMessage,
) {
match msg {
BlockSyncMessage::Status(status) => {
bs_insert_height(row, bs_trace::RANGE_START, status.servable_low);
bs_insert_height(row, bs_trace::HEIGHT, status.servable_high);
}
BlockSyncMessage::Block(block) => {
bs_insert_hash(row, bs_trace::HASH, block.hash());
if let Some(height) = block.coinbase_height() {
bs_insert_height(row, bs_trace::HEIGHT, height);
}
}
BlockSyncMessage::BlocksDone {
start_height,
returned,
} => {
bs_insert_height(row, bs_trace::RANGE_START, *start_height);
bs_insert_u64(row, bs_trace::RANGE_COUNT, u64::from(*returned));
}
BlockSyncMessage::RangeUnavailable {
start_height,
count,
}
| BlockSyncMessage::GetBlocks {
start_height,
count,
} => {
bs_insert_height(row, bs_trace::RANGE_START, *start_height);
bs_insert_u64(row, bs_trace::RANGE_COUNT, u64::from(*count));
}
}
}
pub(super) fn block_sync_message_label(msg: &BlockSyncMessage) -> &'static str {
match msg {
BlockSyncMessage::Status(_) => "status",
BlockSyncMessage::Block(_) => "block",
BlockSyncMessage::BlocksDone { .. } => "blocks_done",
BlockSyncMessage::RangeUnavailable { .. } => "range_unavailable",
BlockSyncMessage::GetBlocks { .. } => "get_blocks",
}
}
fn block_misbehavior_label(reason: BlockSyncMisbehavior) -> &'static str {
match reason {
BlockSyncMisbehavior::MalformedMessage => "malformed_message",
BlockSyncMisbehavior::UnsolicitedBlock => "unsolicited_block",
BlockSyncMisbehavior::GetBlocksTooLong => "get_blocks_too_long",
BlockSyncMisbehavior::GetBlocksSpam => "get_blocks_spam",
BlockSyncMisbehavior::InvalidBlock => "invalid_block",
BlockSyncMisbehavior::SizeMismatch => "size_mismatch",
BlockSyncMisbehavior::InvalidStatus => "invalid_status",
BlockSyncMisbehavior::UnsolicitedDone => "unsolicited_done",
BlockSyncMisbehavior::RangeUnavailable => "range_unavailable",
BlockSyncMisbehavior::StatusSpam => "status_spam",
}
}
pub(super) fn bs_insert_str(
row: &mut serde_json::Map<String, serde_json::Value>,
key: &'static str,
value: &str,
) {
row.insert(key.to_string(), serde_json::Value::from(value.to_string()));
}
fn elapsed_ms_u64(duration: Duration) -> u64 {
u64::try_from(duration.as_millis()).unwrap_or(u64::MAX)
}
fn set_block_reactor_active_connection_gauge(active_connections: usize) {
metrics::gauge!("zakura.p2p.reactor.active_connections", "reactor" => "block_sync")
.set(active_connections as f64);
}
pub(super) fn tolerated_bytes(reserved_bytes: u64, tolerance_percent: u32) -> u64 {
reserved_bytes.saturating_mul(u64::from(tolerance_percent.max(100))) / 100
}