use std::{
collections::VecDeque,
path::{Path, PathBuf},
sync::Arc,
time::Duration,
};
use indexmap::IndexMap;
use tokio::sync::{
mpsc::{error::TryRecvError, UnboundedReceiver, UnboundedSender},
oneshot, watch,
};
use tracing::Span;
use zakura_chain::{
block::{self, Height},
parallel::{commitment_aux::BlockCommitmentRoots, tree::NoteCommitmentTrees},
};
use crate::{
constants::MAX_BLOCK_REORG_HEIGHT,
error::CommitHeaderRangeError,
service::{
check,
finalized_state::{
AuthenticateHeaderRootsError, AuthenticatedHeaderRoots, FinalizedState,
HeaderRootAuthFrontierError, HeaderRootAuthState, HighestCompletedCheckpoint,
HighestCompletedCheckpointTracker, ZakuraDb,
},
non_finalized_state::NonFinalizedState,
queued_blocks::{QueuedCheckpointVerified, QueuedSemanticallyVerified},
ChainTipBlock, ChainTipSender, InvalidateError, ReconsiderError,
},
SemanticallyVerifiedBlock, ValidateContextError,
};
#[allow(unused_imports)]
use crate::service::{
chain_tip::{ChainTipChange, LatestChainTip},
non_finalized_state::Chain,
};
mod vct_write;
use vct_write::VctWriteManager;
#[derive(Copy, Clone, Debug, Eq, PartialEq)]
pub struct VctRootRepairStatus {
pub state: VctRootRepairState,
pub generation: u64,
}
impl Default for VctRootRepairStatus {
fn default() -> Self {
Self {
state: VctRootRepairState::Idle,
generation: 0,
}
}
}
#[derive(Copy, Clone, Debug, Eq, PartialEq)]
pub enum VctRootRepairState {
Idle,
Unavailable {
height: block::Height,
},
}
const PARENT_ERROR_MAP_LIMIT: usize = MAX_BLOCK_REORG_HEIGHT as usize * 2;
#[tracing::instrument(
level = "debug",
skip(finalized_state, non_finalized_state, prepared),
fields(
height = ?prepared.height,
hash = %prepared.hash,
chains = non_finalized_state.chain_count()
)
)]
pub(crate) fn validate_and_commit_non_finalized(
finalized_state: &ZakuraDb,
non_finalized_state: &mut NonFinalizedState,
prepared: SemanticallyVerifiedBlock,
) -> Result<(), ValidateContextError> {
check::initial_contextual_validity(finalized_state, non_finalized_state, &prepared)?;
let parent_hash = prepared.block.header.previous_block_hash;
if finalized_state.finalized_tip_hash() == parent_hash {
non_finalized_state.commit_new_chain(prepared, finalized_state)?;
} else {
non_finalized_state.commit_block(prepared, finalized_state)?;
}
Ok(())
}
#[instrument(
level = "debug",
skip(
non_finalized_state,
chain_tip_sender,
non_finalized_state_sender,
backup_dir_path,
),
fields(chains = non_finalized_state.chain_count())
)]
fn update_latest_chain_channels(
non_finalized_state: &NonFinalizedState,
chain_tip_sender: &mut ChainTipSender,
non_finalized_state_sender: &watch::Sender<NonFinalizedState>,
backup_dir_path: Option<&Path>,
) -> block::Height {
let best_chain = non_finalized_state.best_chain().expect("unexpected empty non-finalized state: must commit at least one block before updating channels");
let tip_block = best_chain
.tip_block()
.expect("unexpected empty chain: must commit at least one block before updating channels")
.clone();
let tip_block = ChainTipBlock::from(tip_block);
let tip_block_height = tip_block.height;
if let Some(backup_dir_path) = backup_dir_path {
non_finalized_state.write_to_backup(backup_dir_path);
}
let _ = non_finalized_state_sender.send(non_finalized_state.clone());
chain_tip_sender.set_best_non_finalized_tip(tip_block);
tip_block_height
}
fn commit_header_range(
finalized_state: &FinalizedState,
completed_checkpoint: &mut HighestCompletedCheckpointTracker,
anchor: block::Hash,
headers: Vec<Arc<block::Header>>,
body_sizes: Vec<u32>,
tree_aux_roots: Vec<BlockCommitmentRoots>,
) -> Result<block::Hash, CommitHeaderRangeError> {
if let Err(height) =
completed_checkpoint.check_immutable_conflicts(&finalized_state.db, anchor, &headers)
{
return Err(CommitHeaderRangeError::ImmutableConflict { height });
}
let mut batch = crate::service::finalized_state::DiskWriteBatch::new();
batch
.prepare_header_range_batch_with_roots(
&finalized_state.db,
anchor,
&headers,
&body_sizes,
&tree_aux_roots,
)
.and_then(|hash| {
let proposed = completed_checkpoint.propose_after_headers(
&finalized_state.db,
anchor,
&headers,
)?;
finalized_state
.db
.write_batch(batch)
.map(|()| {
completed_checkpoint.commit_success(proposed);
hash
})
.map_err(|error| {
tracing::error!(?error, "failed to write validated header range");
CommitHeaderRangeError::StorageWriteError {
error: error.to_string(),
}
})
})
}
fn completed_checkpoint_for_auth_frontier(
completed_checkpoint: &HighestCompletedCheckpointTracker,
) -> Result<HighestCompletedCheckpoint, HeaderRootAuthFrontierError> {
completed_checkpoint
.current()
.ok_or(HeaderRootAuthFrontierError::MissingCompletedCheckpoint)
}
#[allow(clippy::too_many_arguments)]
fn authenticate_header_roots(
finalized_state: &FinalizedState,
completed_checkpoint: &HighestCompletedCheckpointTracker,
header_root_auth_sender: &watch::Sender<Option<HeaderRootAuthState>>,
expected_state: HeaderRootAuthState,
anchor: block::Hash,
start: Height,
headers: Vec<Arc<block::Header>>,
roots: Vec<BlockCommitmentRoots>,
rsp_tx: oneshot::Sender<Result<AuthenticatedHeaderRoots, AuthenticateHeaderRootsError>>,
) {
respond_if_requested(rsp_tx, || {
let completed_checkpoint = completed_checkpoint_for_auth_frontier(completed_checkpoint)?;
let result = finalized_state.db.authenticate_header_roots(
completed_checkpoint,
expected_state,
anchor,
start,
&headers,
&roots,
);
if let Ok(success) = &result {
let _ = header_root_auth_sender.send(Some(success.state));
}
result
});
}
fn respond_if_requested<T, E>(
rsp_tx: oneshot::Sender<Result<T, E>>,
work: impl FnOnce() -> Result<T, E>,
) {
if rsp_tx.is_closed() {
metrics::counter!("state.write.cancelled_before_start").increment(1);
return;
}
let _ = rsp_tx.send(work());
}
fn publish_header_root_auth_state(
db: &ZakuraDb,
completed_checkpoint: &HighestCompletedCheckpointTracker,
sender: &watch::Sender<Option<HeaderRootAuthState>>,
) {
match db.load_header_root_auth_frontier() {
Ok(Some(frontier)) => match completed_checkpoint_for_auth_frontier(completed_checkpoint) {
Ok(completed_checkpoint) => {
let _ = sender.send(Some(frontier.state(completed_checkpoint)));
}
Err(error) => {
tracing::warn!(
?error,
"skipping header-root auth state publish: durable frontier without a completed checkpoint"
);
}
},
Ok(None) => {
let _ = sender.send(None);
}
Err(error) => {
tracing::error!(
?error,
"durable header-root authentication state failed validation after write"
);
}
}
}
struct WriteBlockWorkerTask {
finalized_block_write_receiver: UnboundedReceiver<QueuedCheckpointVerified>,
non_finalized_block_write_receiver: UnboundedReceiver<NonFinalizedWriteMessage>,
finalized_state: FinalizedState,
non_finalized_state: NonFinalizedState,
seed_zakura_header_from_best_chain_commits: bool,
invalid_block_reset_sender: UnboundedSender<block::Hash>,
non_finalized_rejected_sender: UnboundedSender<block::Hash>,
chain_tip_sender: ChainTipSender,
non_finalized_state_sender: watch::Sender<NonFinalizedState>,
highest_completed_checkpoint: HighestCompletedCheckpointTracker,
vct_root_repair_sender: watch::Sender<VctRootRepairStatus>,
header_root_auth_sender: watch::Sender<Option<HeaderRootAuthState>>,
backup_dir_path: Option<PathBuf>,
}
pub enum NonFinalizedWriteMessage {
Commit(QueuedSemanticallyVerified),
CommitHeaderRange {
anchor: block::Hash,
headers: Vec<Arc<block::Header>>,
body_sizes: Vec<u32>,
tree_aux_roots: Vec<BlockCommitmentRoots>,
rsp_tx: oneshot::Sender<Result<block::Hash, CommitHeaderRangeError>>,
},
AuthenticateHeaderRoots {
expected_state: HeaderRootAuthState,
anchor: block::Hash,
start: Height,
headers: Vec<Arc<block::Header>>,
roots: Vec<BlockCommitmentRoots>,
rsp_tx: oneshot::Sender<Result<AuthenticatedHeaderRoots, AuthenticateHeaderRootsError>>,
},
Invalidate {
hash: block::Hash,
rsp_tx: oneshot::Sender<Result<block::Hash, InvalidateError>>,
},
Reconsider {
hash: block::Hash,
rsp_tx: oneshot::Sender<Result<Vec<block::Hash>, ReconsiderError>>,
},
}
impl From<QueuedSemanticallyVerified> for NonFinalizedWriteMessage {
fn from(block: QueuedSemanticallyVerified) -> Self {
NonFinalizedWriteMessage::Commit(block)
}
}
#[derive(Clone, Debug)]
pub struct BlockWriteSender {
pub non_finalized: Option<tokio::sync::mpsc::UnboundedSender<NonFinalizedWriteMessage>>,
pub finalized: Option<tokio::sync::mpsc::UnboundedSender<QueuedCheckpointVerified>>,
}
impl BlockWriteSender {
#[instrument(
level = "debug",
skip_all,
fields(
network = %non_finalized_state.network
)
)]
pub fn spawn(
finalized_state: FinalizedState,
non_finalized_state: NonFinalizedState,
chain_tip_sender: ChainTipSender,
non_finalized_state_sender: watch::Sender<NonFinalizedState>,
should_use_finalized_block_write_sender: bool,
backup_dir_path: Option<PathBuf>,
) -> (
Self,
tokio::sync::mpsc::UnboundedReceiver<block::Hash>,
tokio::sync::mpsc::UnboundedReceiver<block::Hash>,
watch::Receiver<Option<HighestCompletedCheckpoint>>,
watch::Receiver<VctRootRepairStatus>,
watch::Receiver<Option<HeaderRootAuthState>>,
Option<Arc<std::thread::JoinHandle<()>>>,
) {
let (non_finalized_block_write_sender, non_finalized_block_write_receiver) =
tokio::sync::mpsc::unbounded_channel();
let (finalized_block_write_sender, finalized_block_write_receiver) =
tokio::sync::mpsc::unbounded_channel();
let (invalid_block_reset_sender, invalid_block_write_reset_receiver) =
tokio::sync::mpsc::unbounded_channel();
let (non_finalized_rejected_sender, non_finalized_rejected_receiver) =
tokio::sync::mpsc::unbounded_channel();
let (vct_root_repair_sender, vct_root_repair_receiver) =
watch::channel(VctRootRepairStatus::default());
let (highest_completed_checkpoint, highest_completed_checkpoint_receiver) =
HighestCompletedCheckpointTracker::open(&finalized_state.db);
let initial_header_root_auth_state = finalized_state
.db
.validate_header_root_auth_state()
.expect("authenticated header-root state was validated during database startup")
.and_then(|frontier| {
match completed_checkpoint_for_auth_frontier(&highest_completed_checkpoint) {
Ok(completed_checkpoint) => Some(frontier.state(completed_checkpoint)),
Err(error) => {
tracing::warn!(
?error,
"durable header-root authentication frontier exists without a completed checkpoint; publishing no auth state"
);
None
}
}
});
let (header_root_auth_sender, header_root_auth_receiver) =
watch::channel(initial_header_root_auth_state);
let seed_zakura_header_from_best_chain_commits = finalized_state
.db
.config()
.enable_zakura_header_seed_from_committed_blocks;
let span = Span::current();
let task = std::thread::spawn(move || {
span.in_scope(|| {
WriteBlockWorkerTask {
finalized_block_write_receiver,
non_finalized_block_write_receiver,
finalized_state,
non_finalized_state,
seed_zakura_header_from_best_chain_commits,
invalid_block_reset_sender,
non_finalized_rejected_sender,
chain_tip_sender,
non_finalized_state_sender,
highest_completed_checkpoint,
vct_root_repair_sender,
header_root_auth_sender,
backup_dir_path,
}
.run()
})
});
(
Self {
non_finalized: Some(non_finalized_block_write_sender),
finalized: should_use_finalized_block_write_sender
.then_some(finalized_block_write_sender),
},
invalid_block_write_reset_receiver,
non_finalized_rejected_receiver,
highest_completed_checkpoint_receiver,
vct_root_repair_receiver,
header_root_auth_receiver,
Some(Arc::new(task)),
)
}
}
impl WriteBlockWorkerTask {
#[instrument(
level = "debug",
skip(self),
fields(
network = %self.non_finalized_state.network
)
)]
pub fn run(mut self) {
let Self {
finalized_block_write_receiver,
non_finalized_block_write_receiver,
finalized_state,
non_finalized_state,
invalid_block_reset_sender,
non_finalized_rejected_sender,
chain_tip_sender,
non_finalized_state_sender,
highest_completed_checkpoint,
vct_root_repair_sender,
header_root_auth_sender,
seed_zakura_header_from_best_chain_commits,
backup_dir_path,
} = &mut self;
let mut prev_finalized_note_commitment_trees: Option<NoteCommitmentTrees> = None;
let mut deferred_non_finalized_messages = VecDeque::new();
let mut vct_write_manager = VctWriteManager::new(vct_root_repair_sender.clone());
loop {
match non_finalized_block_write_receiver.try_recv() {
Ok(NonFinalizedWriteMessage::CommitHeaderRange {
anchor,
headers,
body_sizes,
tree_aux_roots,
rsp_tx,
}) => {
let result = commit_header_range(
finalized_state,
highest_completed_checkpoint,
anchor,
headers,
body_sizes,
tree_aux_roots,
);
if result.is_ok() {
publish_header_root_auth_state(
&finalized_state.db,
highest_completed_checkpoint,
header_root_auth_sender,
);
}
let _ = rsp_tx.send(result);
continue;
}
Ok(NonFinalizedWriteMessage::AuthenticateHeaderRoots {
expected_state,
anchor,
start,
headers,
roots,
rsp_tx,
}) => {
authenticate_header_roots(
finalized_state,
highest_completed_checkpoint,
header_root_auth_sender,
expected_state,
anchor,
start,
headers,
roots,
rsp_tx,
);
continue;
}
Ok(msg) => deferred_non_finalized_messages.push_back(msg),
Err(TryRecvError::Empty) => {}
Err(TryRecvError::Disconnected) => {}
}
let ordered_block = match vct_write_manager.take_ready() {
Some(block) => block,
None => match finalized_block_write_receiver.try_recv() {
Ok(block) => block,
Err(TryRecvError::Empty) => {
std::thread::park_timeout(Duration::from_millis(10));
continue;
}
Err(TryRecvError::Disconnected) => break,
},
};
if invalid_block_reset_sender.is_closed() {
info!("StateService closed the block reset channel. Is Zakura shutting down?");
return;
}
let next_valid_height = finalized_state
.db
.finalized_tip_height()
.map(|height| (height + 1).expect("committed heights are valid"))
.unwrap_or(Height(0));
if ordered_block.0.height != next_valid_height {
debug!(
?next_valid_height,
invalid_height = ?ordered_block.0.height,
invalid_hash = ?ordered_block.0.hash,
"got a block that was the wrong height. \
Assuming a parent block failed, and dropping this block",
);
vct_write_manager.reset(finalized_state);
std::mem::drop(ordered_block);
continue;
}
vct_write_manager.fill_successor(finalized_block_write_receiver, &ordered_block);
let needs_vct_successor =
finalized_state.vct_fast_needs_successor(ordered_block.0.height);
let next_vct_block = if needs_vct_successor {
finalized_state
.vct_successor_from_header_store(ordered_block.0.height, ordered_block.0.hash)
} else {
None
};
if needs_vct_successor && next_vct_block.is_none() {
let height = ordered_block.0.height;
let wait =
vct_write_manager.on_retryable_error(height, false, false, ordered_block);
std::thread::park_timeout(wait);
continue;
}
let prev_note_commitment_trees = prev_finalized_note_commitment_trees.take();
let prev_note_commitment_trees_for_retry = prev_note_commitment_trees.clone();
let next_block_took_vct_path =
finalized_state.vct_fast_will_apply(ordered_block.0.height);
match finalized_state.commit_finalized(
ordered_block,
prev_note_commitment_trees,
next_vct_block,
) {
Ok((finalized, note_commitment_trees, rsp_tx)) => {
if next_block_took_vct_path {
metrics::counter!("state.vct.fast_path.hit").increment(1);
} else {
metrics::counter!("state.vct.fast_path.miss").increment(1);
}
vct_write_manager.on_commit_success();
let tip_hash = finalized.hash;
let tip_block = ChainTipBlock::from(finalized);
prev_finalized_note_commitment_trees = Some(note_commitment_trees);
match highest_completed_checkpoint.rebind_from_db(&finalized_state.db) {
Ok(()) => {
publish_header_root_auth_state(
&finalized_state.db,
highest_completed_checkpoint,
header_root_auth_sender,
);
}
Err(error) => {
tracing::warn!(
?error,
"failed to refresh highest completed checkpoint after finalized block commit"
);
}
}
chain_tip_sender.set_finalized_tip(tip_block);
let _ = rsp_tx.send(Ok(tip_hash));
}
Err((ordered_block, error)) => {
if let Some(height) = error.vct_retryable_height() {
let root_unavailable = error.vct_supplied_root_unavailable_height();
prev_finalized_note_commitment_trees = prev_note_commitment_trees_for_retry;
let wait = vct_write_manager.on_retryable_error(
height,
root_unavailable.is_some(),
next_block_took_vct_path,
ordered_block,
);
std::thread::park_timeout(wait);
continue;
}
let finalized_tip = finalized_state.db.tip();
let _ = ordered_block.1.send(Err(error.clone()));
vct_write_manager.reset(finalized_state);
info!(
?error,
last_valid_height = ?finalized_tip.map(|tip| tip.0),
last_valid_hash = ?finalized_tip.map(|tip| tip.1),
"committing a block to the finalized state failed, resetting state queue",
);
let send_result =
invalid_block_reset_sender.send(finalized_state.db.finalized_tip_hash());
if send_result.is_err() {
info!(
"StateService closed the block reset channel. Is Zakura shutting down?"
);
return;
}
}
}
}
if invalid_block_reset_sender.is_closed() {
info!("StateService closed the block reset channel. Is Zakura shutting down?");
return;
}
let mut parent_error_map: IndexMap<block::Hash, ValidateContextError> = IndexMap::new();
while let Some(msg) = deferred_non_finalized_messages
.pop_front()
.or_else(|| non_finalized_block_write_receiver.blocking_recv())
{
let queued_child_and_rsp_tx = match msg {
NonFinalizedWriteMessage::Commit(queued_child) => Some(queued_child),
NonFinalizedWriteMessage::CommitHeaderRange {
anchor,
headers,
body_sizes,
tree_aux_roots,
rsp_tx,
} => {
let result = commit_header_range(
finalized_state,
highest_completed_checkpoint,
anchor,
headers,
body_sizes,
tree_aux_roots,
);
if result.is_ok() {
publish_header_root_auth_state(
&finalized_state.db,
highest_completed_checkpoint,
header_root_auth_sender,
);
}
let _ = rsp_tx.send(result);
continue;
}
NonFinalizedWriteMessage::AuthenticateHeaderRoots {
expected_state,
anchor,
start,
headers,
roots,
rsp_tx,
} => {
authenticate_header_roots(
finalized_state,
highest_completed_checkpoint,
header_root_auth_sender,
expected_state,
anchor,
start,
headers,
roots,
rsp_tx,
);
continue;
}
NonFinalizedWriteMessage::Invalidate { hash, rsp_tx } => {
tracing::info!(?hash, "invalidating a block in the non-finalized state");
let _ = rsp_tx.send(non_finalized_state.invalidate_block(hash));
None
}
NonFinalizedWriteMessage::Reconsider { hash, rsp_tx } => {
tracing::info!(?hash, "reconsidering a block in the non-finalized state");
let _ = rsp_tx
.send(non_finalized_state.reconsider_block(hash, &finalized_state.db));
None
}
};
let Some((queued_child, rsp_tx)) = queued_child_and_rsp_tx else {
update_latest_chain_channels(
non_finalized_state,
chain_tip_sender,
non_finalized_state_sender,
backup_dir_path.as_deref(),
);
continue;
};
let child_hash = queued_child.hash;
let parent_hash = queued_child.block.header.previous_block_hash;
let child_height = queued_child.height;
let child_block = queued_child.block.clone();
let parent_error = parent_error_map.get(&parent_hash);
let result = if let Some(parent_error) = parent_error {
Err(parent_error.clone())
} else {
tracing::trace!(?child_hash, "validating queued child");
validate_and_commit_non_finalized(
&finalized_state.db,
non_finalized_state,
queued_child,
)
};
if let Err(ref error) = result {
parent_error_map.insert(child_hash, error.clone());
if parent_error_map.len() > PARENT_ERROR_MAP_LIMIT {
parent_error_map.shift_remove_index(0);
}
let _ = non_finalized_rejected_sender.send(child_hash);
let _ = rsp_tx.send(result.map(|()| child_hash).map_err(Into::into));
continue;
}
parent_error_map.shift_remove(&child_hash);
if should_seed_zakura_header_from_non_finalized_commit(
*seed_zakura_header_from_best_chain_commits,
non_finalized_state,
child_height,
child_hash,
) && seed_zakura_header_from_committed_block(
&finalized_state.db,
highest_completed_checkpoint,
child_height,
&child_block,
) {
publish_header_root_auth_state(
&finalized_state.db,
highest_completed_checkpoint,
header_root_auth_sender,
);
}
let tip_block_height = update_latest_chain_channels(
non_finalized_state,
chain_tip_sender,
non_finalized_state_sender,
backup_dir_path.as_deref(),
);
let _ = rsp_tx.send(result.map(|()| child_hash).map_err(Into::into));
while non_finalized_state
.best_chain_len()
.expect("just successfully inserted a non-finalized block above")
> MAX_BLOCK_REORG_HEIGHT
{
tracing::trace!("finalizing block past the reorg limit");
let contextually_verified_with_trees = non_finalized_state.finalize();
prev_finalized_note_commitment_trees = finalized_state
.commit_finalized_direct(
contextually_verified_with_trees,
prev_finalized_note_commitment_trees.take(),
None,
"commit contextually-verified request",
)
.expect(
"unexpected finalized block commit error: note commitment and history trees were already checked by the non-finalized state",
)
.1
.into();
match highest_completed_checkpoint.rebind_from_db(&finalized_state.db) {
Ok(()) => publish_header_root_auth_state(
&finalized_state.db,
highest_completed_checkpoint,
header_root_auth_sender,
),
Err(error) => {
tracing::warn!(
?error,
"failed to refresh highest completed checkpoint after finalized block commit"
);
}
}
}
metrics::counter!("state.full_verifier.committed.block.count").increment(1);
metrics::counter!("zcash.chain.verified.block.total").increment(1);
metrics::gauge!("state.full_verifier.committed.block.height")
.set(tip_block_height.0 as f64);
metrics::gauge!("zcash.chain.verified.block.height").set(tip_block_height.0 as f64);
tracing::trace!("finished processing queued block");
}
finalized_state.db.shutdown(true);
std::mem::drop(self.finalized_state);
}
}
fn seed_zakura_header_from_committed_block(
finalized_state: &ZakuraDb,
highest_completed_checkpoint: &mut HighestCompletedCheckpointTracker,
height: block::Height,
block: &Arc<block::Block>,
) -> bool {
match finalized_state.seed_zakura_header_from_committed_block(height, block) {
Ok(()) => {
if let Err(error) = highest_completed_checkpoint.rebind_from_db(finalized_state) {
tracing::warn!(
?error,
"failed to refresh highest completed checkpoint after seeding a header"
);
return false;
}
tracing::trace!(?height, hash = ?block.hash(), "seeded Zakura header from committed block");
true
}
Err(error) => {
tracing::warn!(
?height,
hash = ?block.hash(),
?error,
"failed to seed Zakura header from committed block"
);
false
}
}
}
fn should_seed_zakura_header_from_non_finalized_commit(
enabled: bool,
non_finalized_state: &NonFinalizedState,
height: block::Height,
hash: block::Hash,
) -> bool {
enabled && non_finalized_state.best_tip() == Some((height, hash))
}
#[cfg(test)]
mod tests {
use std::sync::{
atomic::{AtomicBool, Ordering},
Arc,
};
use zakura_chain::{
block::Height, history_tree::HistoryTree, parameters::Network,
serialization::ZcashDeserializeInto, value_balance::ValueBalance,
};
use crate::{
arbitrary::Prepare,
service::{
finalized_state::{
AuthenticateHeaderRootsError, AuthenticateHeaderRootsOutcome, DiskWriteBatch,
FinalizedState, HeaderRootAuthFrontierError, HeaderRootAuthState,
HighestCompletedCheckpointTracker, WriteDisk,
},
non_finalized_state::NonFinalizedState,
write::{
authenticate_header_roots, completed_checkpoint_for_auth_frontier,
publish_header_root_auth_state, respond_if_requested,
seed_zakura_header_from_committed_block,
should_seed_zakura_header_from_non_finalized_commit,
},
},
tests::FakeChainHelper,
Config,
};
#[test]
fn cancelled_response_skips_serialized_write_work() {
let (rsp_tx, rsp_rx) = tokio::sync::oneshot::channel();
drop(rsp_rx);
let ran = AtomicBool::new(false);
respond_if_requested(rsp_tx, || {
ran.store(true, Ordering::SeqCst);
Ok::<_, ()>(())
});
assert!(!ran.load(Ordering::SeqCst));
}
#[test]
fn missing_completed_checkpoint_is_a_local_frontier_error() {
let _init_guard = zakura_test::init();
let finalized_state = FinalizedState::new(&Config::ephemeral(), &Network::Mainnet)
.expect("opening an ephemeral database should succeed");
let (completed_checkpoint, _receiver) =
HighestCompletedCheckpointTracker::open(&finalized_state.db);
assert!(completed_checkpoint.current().is_none());
assert!(matches!(
completed_checkpoint_for_auth_frontier(&completed_checkpoint),
Err(HeaderRootAuthFrontierError::MissingCompletedCheckpoint)
));
let expected_state = HeaderRootAuthState {
authenticated_height: Height::MIN,
authenticated_hash: Network::Mainnet.genesis_hash(),
completed_checkpoint_height: Height::MIN,
completed_checkpoint_hash: Network::Mainnet.genesis_hash(),
header_witness: None,
};
let (header_root_auth_sender, mut header_root_auth_receiver) =
tokio::sync::watch::channel(Some(expected_state));
let _ = header_root_auth_receiver.borrow_and_update();
let (rsp_tx, rsp_rx) = tokio::sync::oneshot::channel();
authenticate_header_roots(
&finalized_state,
&completed_checkpoint,
&header_root_auth_sender,
expected_state,
Network::Mainnet.genesis_hash(),
Height::MIN,
Vec::new(),
Vec::new(),
rsp_tx,
);
let error = rsp_rx
.blocking_recv()
.expect("write path answers the oneshot")
.expect_err("missing completed checkpoint is an authentication error");
assert!(matches!(
error,
AuthenticateHeaderRootsError::Frontier(
HeaderRootAuthFrontierError::MissingCompletedCheckpoint
)
));
assert_eq!(error.outcome(), AuthenticateHeaderRootsOutcome::Local);
assert!(!header_root_auth_receiver
.has_changed()
.expect("sender open"));
}
#[test]
fn publish_skips_when_durable_frontier_lacks_completed_checkpoint() {
let _init_guard = zakura_test::init();
let finalized_state = FinalizedState::new(&Config::ephemeral(), &Network::Mainnet)
.expect("opening an ephemeral database should succeed");
let genesis = zakura_test::vectors::BLOCK_MAINNET_GENESIS_BYTES
.zcash_deserialize_into::<Arc<zakura_chain::block::Block>>()
.expect("mainnet genesis block deserializes");
let hash_by_height = finalized_state
.db
.db()
.cf_handle("hash_by_height")
.expect("hash_by_height column family exists");
let height_by_hash = finalized_state
.db
.db()
.cf_handle("height_by_hash")
.expect("height_by_hash column family exists");
let block_header_by_height = finalized_state
.db
.db()
.cf_handle("block_header_by_height")
.expect("block_header_by_height column family exists");
let mut batch = DiskWriteBatch::new();
batch.zs_insert(&hash_by_height, Height::MIN, genesis.hash());
batch.zs_insert(&height_by_hash, genesis.hash(), Height::MIN);
batch.zs_insert(&block_header_by_height, Height::MIN, &genesis.header);
finalized_state
.db
.write_batch(batch)
.expect("genesis tip rows write");
let mut batch = DiskWriteBatch::new();
batch
.rebase_header_root_auth_frontier(
&finalized_state.db,
Height::MIN,
genesis.hash(),
&HistoryTree::default(),
)
.expect("genesis frontier is coherent");
finalized_state
.db
.write_batch(batch)
.expect("genesis frontier writes");
let (mut completed_checkpoint, _receiver) =
HighestCompletedCheckpointTracker::open(&finalized_state.db);
assert!(
completed_checkpoint.current().is_some(),
"body tip at genesis completes the genesis checkpoint"
);
completed_checkpoint.clear_published_for_test();
assert!(completed_checkpoint.current().is_none());
assert!(finalized_state
.db
.load_header_root_auth_frontier()
.expect("frontier loads after clear")
.is_some());
let prior = HeaderRootAuthState {
authenticated_height: Height::MIN,
authenticated_hash: genesis.hash(),
completed_checkpoint_height: Height::MIN,
completed_checkpoint_hash: genesis.hash(),
header_witness: None,
};
let (sender, mut receiver) = tokio::sync::watch::channel(Some(prior));
let _ = receiver.borrow_and_update();
publish_header_root_auth_state(&finalized_state.db, &completed_checkpoint, &sender);
assert!(!receiver.has_changed().expect("sender open"));
assert_eq!(*receiver.borrow(), Some(prior));
}
#[test]
fn side_chain_commit_does_not_seed_zakura_headers() {
let _init_guard = zakura_test::init();
let network = Network::Mainnet;
let mut config = Config::ephemeral();
config.enable_zakura_header_seed_from_committed_blocks = true;
let finalized_state = FinalizedState::new(&config, &network)
.expect("opening an ephemeral database should succeed");
finalized_state.set_finalized_value_pool(ValueBalance::fake_populated_pool());
let parent = zakura_test::vectors::BLOCK_MAINNET_434873_BYTES
.zcash_deserialize_into::<Arc<zakura_chain::block::Block>>()
.expect("block deserializes");
let best_block = parent.make_fake_child().set_work(10);
let side_block = parent.make_fake_child().set_work(1);
let best_height = best_block
.coinbase_height()
.expect("fake child block has a coinbase height");
let mut non_finalized_state = NonFinalizedState::new(&network);
let (mut completed_checkpoint, _receiver) =
HighestCompletedCheckpointTracker::open(&finalized_state.db);
let parent_height = parent
.coinbase_height()
.expect("test vector block has a coinbase height");
let zakura_hash_by_height = finalized_state
.db
.db()
.cf_handle("zakura_header_hash_by_height")
.unwrap();
let mut batch = DiskWriteBatch::new();
batch.zs_insert(&zakura_hash_by_height, parent_height, parent.hash());
finalized_state
.db
.db()
.write(batch)
.expect("parent hash row writes");
non_finalized_state
.commit_new_chain(best_block.clone().prepare(), &finalized_state)
.expect("best block commits to a new chain");
assert!(should_seed_zakura_header_from_non_finalized_commit(
true,
&non_finalized_state,
best_height,
best_block.hash(),
));
seed_zakura_header_from_committed_block(
&finalized_state.db,
&mut completed_checkpoint,
best_height,
&best_block,
);
non_finalized_state
.commit_new_chain(side_block.clone().prepare(), &finalized_state)
.expect("side block commits to a losing fork");
assert!(!should_seed_zakura_header_from_non_finalized_commit(
true,
&non_finalized_state,
best_height,
side_block.hash(),
));
assert_eq!(
finalized_state.db.best_header_tip(),
Some((best_height, best_block.hash()))
);
assert_eq!(
finalized_state.db.headers_by_height_range(best_height, 1),
vec![(best_height, best_block.hash(), best_block.header.clone())],
);
}
}