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_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)>,
block_wait_started: Option<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,
}
#[derive(Copy, Clone, Debug, Eq, PartialEq)]
pub(super) enum VctWriteRetryWait {
HeaderChainInsert,
Delay(Duration),
}
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,
block_wait_started: 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.finish_block_wait();
self.clear_root_stall();
self.clear_committer_repair();
}
pub(super) fn on_commit_success(&mut self) {
self.finish_block_wait();
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"
);
}
self.clear_root_stall();
self.clear_committer_repair();
}
pub(super) fn on_retryable_error(
&mut self,
height: Height,
retry_cause: VctWriteRetryCause,
block: QueuedCheckpointVerified,
) -> VctWriteRetryWait {
metrics::counter!("state.vct.root.retry.count").increment(1);
if let VctWriteRetryCause::MissingRoot { trigger } = retry_cause {
self.request_committer_repair(height, trigger);
}
self.block_wait_started.get_or_insert_with(Instant::now);
let new_stall = match self.root_stall {
Some((stalled_height, _)) if stalled_height == height => false,
_ => {
self.root_stall = Some((height, Instant::now()));
if self.root_stall_reported {
metrics::gauge!("state.vct.root.stalled.height").set(f64::from(height.0));
}
true
}
};
self.report_stall_if_due(retry_cause);
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 { .. } => VctWriteRetryWait::HeaderChainInsert,
VctWriteRetryCause::MissingSuccessor => {
VctWriteRetryWait::Delay(VCT_AWAIT_SUCCESSOR_WAIT)
}
}
}
pub(super) fn stall_warning_remaining(&self) -> Option<Duration> {
if self.root_stall_reported {
return None;
}
self.block_wait_started
.map(|since| VCT_ROOT_STALL_WARN_AFTER.saturating_sub(since.elapsed()))
}
pub(super) fn report_stall_if_due(&mut self, retry_cause: VctWriteRetryCause) {
let Some((height, _)) = self.root_stall else {
return;
};
let Some(since) = self.block_wait_started else {
return;
};
if self.root_stall_reported || since.elapsed() < VCT_ROOT_STALL_WARN_AFTER {
return;
}
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;
}
fn finish_block_wait(&mut self) {
if let Some(since) = self.block_wait_started.take() {
metrics::histogram!("state.vct.root.wait.seconds")
.record(since.elapsed().as_secs_f64());
}
}
fn clear_root_stall(&mut self) {
if self.root_stall_reported {
metrics::gauge!("state.vct.root.stalled.height").set(0.0);
}
self.root_stall = None;
self.root_stall_reported = false;
}
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;