use std::time::{Duration, Instant};
use tokio::sync::watch;
use tracing::info;
use zakura_chain::block::Height;
use zakura_header_chain::EvidenceId;
use crate::service::{
finalized_state::FinalizedState,
queued_blocks::QueuedCheckpointVerified,
write::{VctRootRepairState, VctRootRepairStatus},
};
const VCT_ROOT_RETRY_WAIT: Duration = Duration::from_millis(500);
const VCT_AWAIT_SUCCESSOR_WAIT: Duration = Duration::from_millis(20);
const VCT_ROOT_STALL_WARN_AFTER: Duration = Duration::from_secs(30);
pub(super) struct VctWriteRetryManager {
retryable_block: Option<QueuedCheckpointVerified>,
root_stall: Option<(Height, Instant)>,
root_stall_reported: bool,
root_repair_sender: watch::Sender<VctRootRepairStatus>,
root_repair_status: VctRootRepairStatus,
committer_repair_height: Option<Height>,
unrecorded_committer_rejection: Option<(Height, EvidenceId)>,
sweep_repair_height: Option<Height>,
}
#[derive(Copy, Clone, Debug, Eq, PartialEq)]
enum VctRepairRequester {
Committer,
Sweep,
}
#[derive(Copy, Clone, Debug, Eq, PartialEq)]
pub(super) enum VctRepairTrigger {
MissingRootObserved,
RejectedDelivery,
UnrecordedRejectedDelivery(EvidenceId),
}
impl VctRepairTrigger {
fn starts_new_episode(self) -> bool {
self != Self::MissingRootObserved
}
}
#[derive(Copy, Clone, Debug, Eq, PartialEq)]
pub(super) enum VctWriteRetryCause {
MissingRoot {
trigger: VctRepairTrigger,
},
MissingSuccessor,
}
impl Default for VctWriteRetryManager {
fn default() -> Self {
let (root_repair_sender, _root_repair_receiver) =
watch::channel(VctRootRepairStatus::default());
Self::new(root_repair_sender)
}
}
impl VctWriteRetryManager {
pub(super) fn new(root_repair_sender: watch::Sender<VctRootRepairStatus>) -> Self {
Self {
retryable_block: None,
root_stall: None,
root_stall_reported: false,
root_repair_sender,
root_repair_status: VctRootRepairStatus::default(),
committer_repair_height: None,
unrecorded_committer_rejection: None,
sweep_repair_height: None,
}
}
pub(super) fn request_sweep_repair(&mut self, height: Height, trigger: VctRepairTrigger) {
let starts_new_episode = trigger.starts_new_episode();
self.sweep_repair_height = Some(height);
self.publish_effective_repair_status(
starts_new_episode.then_some(VctRepairRequester::Sweep),
);
}
pub(super) fn clear_sweep_repair(&mut self) {
self.sweep_repair_height = None;
self.publish_effective_repair_status(None);
}
pub(super) fn sweep_repair_height(&self) -> Option<Height> {
self.sweep_repair_height
}
#[cfg(test)]
pub(super) fn request_committer_repair_for_test(&mut self, height: Height) {
self.request_committer_repair(height, VctRepairTrigger::RejectedDelivery);
}
pub(super) fn take_retryable_block(&mut self) -> Option<QueuedCheckpointVerified> {
self.retryable_block.take()
}
pub(super) fn reset(&mut self, finalized_state: &mut FinalizedState) {
finalized_state.clear_vct_prevalidated_next();
self.clear_committer_repair();
}
pub(super) fn on_commit_success(&mut self) {
if self.root_stall.is_some() {
if self.root_stall_reported {
info!(
stalled_height = ?self.root_stall.map(|(height, _)| height),
"VCT: checkpoint commit recovered; the stalled height now has a verifiable supplied root"
);
metrics::gauge!("state.vct.root.stalled.height").set(0.0);
}
self.root_stall = None;
self.root_stall_reported = false;
}
self.clear_committer_repair();
}
pub(super) fn on_retryable_error(
&mut self,
height: Height,
retry_cause: VctWriteRetryCause,
block: QueuedCheckpointVerified,
) -> Duration {
metrics::counter!("state.vct.root.retry.count").increment(1);
if let VctWriteRetryCause::MissingRoot { trigger } = retry_cause {
self.request_committer_repair(height, trigger);
}
let new_stall = match self.root_stall {
Some((stalled_height, _)) if stalled_height == height => false,
_ => {
self.root_stall = Some((height, Instant::now()));
self.root_stall_reported = false;
true
}
};
if !self.root_stall_reported
&& self
.root_stall
.is_some_and(|(_, since)| since.elapsed() >= VCT_ROOT_STALL_WARN_AFTER)
{
tracing::error!(
?height,
?retry_cause,
stalled_for = ?VCT_ROOT_STALL_WARN_AFTER,
"VCT: checkpoint commit stalled waiting for a verifiable supplied root \
or successor witness; the node will not recompute against the frozen frontier"
);
metrics::gauge!("state.vct.root.stalled.height").set(f64::from(height.0));
self.root_stall_reported = true;
} else if new_stall {
tracing::warn!(
?height,
block_height = ?block.0.height,
block_hash = ?block.0.hash,
?retry_cause,
"VCT: supplied root not yet verifiable; retrying checkpoint commit in place"
);
} else {
tracing::trace!(
?height,
block_height = ?block.0.height,
block_hash = ?block.0.hash,
?retry_cause,
"VCT: supplied root still not verifiable; retrying checkpoint commit in place"
);
}
self.retryable_block = Some(block);
match retry_cause {
VctWriteRetryCause::MissingRoot { .. } => VCT_ROOT_RETRY_WAIT,
VctWriteRetryCause::MissingSuccessor => VCT_AWAIT_SUCCESSOR_WAIT,
}
}
fn request_committer_repair(&mut self, height: Height, trigger: VctRepairTrigger) {
let starts_new_episode = match trigger {
VctRepairTrigger::MissingRootObserved => false,
VctRepairTrigger::RejectedDelivery => true,
VctRepairTrigger::UnrecordedRejectedDelivery(delivery_id) => {
let rejection = (height, delivery_id);
let changed = self.unrecorded_committer_rejection != Some(rejection);
self.unrecorded_committer_rejection = Some(rejection);
changed
}
};
self.committer_repair_height = Some(height);
self.publish_effective_repair_status(
starts_new_episode.then_some(VctRepairRequester::Committer),
);
}
fn clear_committer_repair(&mut self) {
self.committer_repair_height = None;
self.unrecorded_committer_rejection = None;
self.publish_effective_repair_status(None);
}
fn publish_effective_repair_status(
&mut self,
requester_with_new_episode: Option<VctRepairRequester>,
) {
let effective_repair = match (self.committer_repair_height, self.sweep_repair_height) {
(Some(committer), _) => Some((committer, VctRepairRequester::Committer)),
(None, Some(sweep)) => Some((sweep, VctRepairRequester::Sweep)),
(None, None) => None,
};
let repair_state = effective_repair.map_or(VctRootRepairState::Idle, |(height, _)| {
VctRootRepairState::Unavailable { height }
});
let effective_requester = effective_repair.map(|(_, requester)| requester);
let effective_episode_changed = requester_with_new_episode.is_some_and(|requester| {
effective_requester == Some(requester) && repair_state == self.root_repair_status.state
});
if repair_state == self.root_repair_status.state && !effective_episode_changed {
return;
}
let repair_requested = repair_state != VctRootRepairState::Idle;
self.root_repair_status = VctRootRepairStatus {
state: repair_state,
generation: if repair_requested {
self.root_repair_status.generation.saturating_add(1)
} else {
self.root_repair_status.generation
},
};
let _ = self.root_repair_sender.send(self.root_repair_status);
if repair_requested {
metrics::counter!("state.vct.root.repair.requested").increment(1);
}
}
}
#[cfg(test)]
mod tests;