use std::{sync::Arc, time::Duration, time::Instant};
use tokio::sync::{mpsc, oneshot};
use zakura_chain::{block::Height, serialization::ZcashDeserializeInto};
use super::{VctWriteManager, VCT_AWAIT_SUCCESSOR_WAIT, VCT_ROOT_RETRY_WAIT};
use crate::{
request::CheckpointVerifiedBlock,
service::{
queued_blocks::QueuedCheckpointVerified,
write::{VctRootRepairState, VctRootRepairStatus},
},
tests::FakeChainHelper,
};
fn queued_block(seed: u128) -> QueuedCheckpointVerified {
let genesis = zakura_test::vectors::BLOCK_MAINNET_GENESIS_BYTES
.zcash_deserialize_into::<Arc<zakura_chain::block::Block>>()
.expect("genesis block deserializes");
let block = genesis.make_fake_child().set_work(seed);
let (rsp_tx, _rsp_rx) = oneshot::channel();
(CheckpointVerifiedBlock::from(block), rsp_tx)
}
fn queued_child_of(parent: &QueuedCheckpointVerified, seed: u128) -> QueuedCheckpointVerified {
let block = parent.0.block.clone().make_fake_child().set_work(seed);
let (rsp_tx, _rsp_rx) = oneshot::channel();
(CheckpointVerifiedBlock::from(block), rsp_tx)
}
#[test]
fn take_ready_returns_none_when_empty() {
let mut manager = VctWriteManager::default();
assert!(manager.take_ready().is_none());
}
#[test]
fn take_ready_prefers_retry_over_lookahead() {
let mut manager = VctWriteManager::default();
let retry_block = queued_block(1);
let retry_hash = retry_block.0.hash;
let lookahead_block = queued_block(2);
let lookahead_hash = lookahead_block.0.hash;
manager.lookahead.push_back(lookahead_block);
manager.retry = Some(retry_block);
let first = manager.take_ready().expect("retry block is ready");
assert_eq!(
first.0.hash, retry_hash,
"retry must be taken before lookahead"
);
let second = manager.take_ready().expect("lookahead block is ready");
assert_eq!(second.0.hash, lookahead_hash);
assert!(manager.take_ready().is_none());
}
#[test]
fn fill_successor_only_buffers_one_linking_block_at_a_time() {
let mut manager = VctWriteManager::default();
let (tx, mut rx) = mpsc::unbounded_channel();
let current = queued_block(1);
let first = queued_child_of(¤t, 2);
let first_hash = first.0.hash;
let second = queued_child_of(&first, 3);
let second_hash = second.0.hash;
tx.send(first).expect("channel is open");
tx.send(second).expect("channel is open");
manager.fill_successor(&mut rx, ¤t);
assert_eq!(manager.lookahead.len(), 1);
assert_eq!(manager.lookahead.front().unwrap().0.hash, first_hash);
manager.fill_successor(&mut rx, ¤t);
assert_eq!(manager.lookahead.len(), 1);
assert_eq!(manager.lookahead.front().unwrap().0.hash, first_hash);
let committed = manager.take_ready().expect("look-ahead block is ready");
manager.fill_successor(&mut rx, &committed);
assert_eq!(manager.lookahead.front().unwrap().0.hash, second_hash);
}
#[test]
fn fill_successor_discards_non_successors_and_keeps_the_linking_block() {
let mut manager = VctWriteManager::default();
let (tx, mut rx) = mpsc::unbounded_channel();
let current = queued_block(1);
let non_successor = queued_block(2);
let successor = queued_child_of(¤t, 3);
let successor_hash = successor.0.hash;
manager.lookahead.push_back(non_successor);
tx.send(successor).expect("channel is open");
manager.fill_successor(&mut rx, ¤t);
assert_eq!(manager.lookahead.len(), 1);
assert_eq!(
manager.lookahead.front().unwrap().0.hash,
successor_hash,
"the non-linking block is dropped and the true successor takes its place"
);
assert_eq!(
manager
.lookahead
.front()
.expect("successor is buffered")
.0
.block
.header
.previous_block_hash,
current.0.hash,
);
}
#[test]
fn fill_successor_empties_the_lookahead_when_no_successor_exists() {
let mut manager = VctWriteManager::default();
let (_tx, mut rx) = mpsc::unbounded_channel();
let current = queued_block(1);
manager.lookahead.push_back(queued_block(2));
manager.fill_successor(&mut rx, ¤t);
assert!(manager.lookahead.is_empty());
}
#[test]
fn buffered_successor_body_remains_queued_for_its_own_commit() {
let mut manager = VctWriteManager::default();
let block = queued_block(1);
let expected_hash = block.0.hash;
manager.lookahead.push_back(block);
let ready = manager
.take_ready()
.expect("the buffered successor body is ready for its own commit");
assert_eq!(ready.0.hash, expected_hash);
}
#[test]
fn reset_clears_the_lookahead() {
let mut manager = VctWriteManager::default();
manager.lookahead.push_back(queued_block(1));
manager.lookahead.push_back(queued_block(2));
let network = zakura_chain::parameters::Network::Mainnet;
let config = crate::Config::ephemeral();
let mut finalized_state =
crate::service::finalized_state::FinalizedState::new(&config, &network)
.expect("opening an ephemeral database should succeed");
manager.reset(&mut finalized_state);
assert!(manager.lookahead.is_empty());
}
#[test]
fn on_commit_success_is_a_no_op_without_a_stall() {
let mut manager = VctWriteManager::default();
manager.on_commit_success();
assert!(manager.stall.is_none());
assert!(!manager.stall_logged);
}
#[test]
fn on_commit_success_clears_an_escalated_stall() {
let mut manager = VctWriteManager::default();
let height = Height(1);
manager.stall = Some((height, Instant::now() - Duration::from_secs(31)));
manager.on_retryable_error(height, true, false, queued_block(1));
assert!(manager.stall_logged, "the stall should have been escalated");
manager.on_commit_success();
assert!(manager.stall.is_none());
assert!(!manager.stall_logged);
}
#[test]
fn on_retryable_error_keeps_the_same_stall_start_for_a_repeated_height() {
let mut manager = VctWriteManager::default();
let height = Height(5);
manager.on_retryable_error(height, true, false, queued_block(1));
let first_seen = manager.stall.expect("a stall is now tracked").1;
manager.on_retryable_error(height, true, false, queued_block(2));
let still_first_seen = manager.stall.expect("the stall is still tracked").1;
assert_eq!(
first_seen, still_first_seen,
"retrying the same height must not reset the stall's start time"
);
}
#[test]
fn on_retryable_error_resets_the_stall_for_a_different_height() {
let mut manager = VctWriteManager::default();
manager.on_retryable_error(Height(1), true, false, queued_block(1));
manager.stall_logged = true;
manager.on_retryable_error(Height(2), true, false, queued_block(2));
assert_eq!(manager.stall.map(|(h, _)| h), Some(Height(2)));
assert!(
!manager.stall_logged,
"a new height starts a fresh, unescalated stall"
);
}
#[test]
fn on_retryable_error_escalates_past_the_warn_threshold() {
let mut manager = VctWriteManager::default();
let height = Height(7);
manager.on_retryable_error(height, true, false, queued_block(1));
assert!(!manager.stall_logged);
manager.stall = Some((height, Instant::now() - Duration::from_secs(31)));
manager.on_retryable_error(height, true, false, queued_block(2));
assert!(manager.stall_logged);
}
#[test]
fn on_retryable_error_parks_the_block_for_retry() {
let mut manager = VctWriteManager::default();
let block = queued_block(1);
let hash = block.0.hash;
manager.on_retryable_error(Height(1), true, false, block);
let ready = manager
.take_ready()
.expect("the block was parked for retry");
assert_eq!(ready.0.hash, hash);
}
#[test]
fn on_retryable_error_wait_depends_on_root_availability() {
let mut manager = VctWriteManager::default();
let missing_root_wait = manager.on_retryable_error(Height(1), true, false, queued_block(1));
assert_eq!(missing_root_wait, VCT_ROOT_RETRY_WAIT);
let successor_wait = manager.on_retryable_error(Height(2), false, false, queued_block(2));
assert_eq!(successor_wait, VCT_AWAIT_SUCCESSOR_WAIT);
}
#[test]
fn root_repair_signal_deduplicates_repeated_missing_root_polls() {
let (tx, mut rx) = tokio::sync::watch::channel(VctRootRepairStatus::default());
let mut manager = VctWriteManager::new(tx);
manager.on_retryable_error(Height(42), true, false, queued_block(1));
let first = *rx.borrow_and_update();
assert_eq!(
first.state,
VctRootRepairState::Unavailable { height: Height(42) }
);
assert_eq!(first.generation, 1);
manager.on_retryable_error(Height(42), true, false, queued_block(2));
assert!(
!rx.has_changed().expect("watch channel remains open"),
"polling the same absent root must not flood repair notifications"
);
}
#[test]
fn root_repair_signal_advances_generation_after_rejected_replacement() {
let (tx, mut rx) = tokio::sync::watch::channel(VctRootRepairStatus::default());
let mut manager = VctWriteManager::new(tx);
manager.on_retryable_error(Height(42), true, false, queued_block(1));
let first = *rx.borrow_and_update();
assert_eq!(first.generation, 1);
manager.on_retryable_error(Height(42), true, true, queued_block(2));
let second = *rx.borrow_and_update();
assert_eq!(second.generation, 2);
assert_eq!(second.state, first.state);
}
#[test]
fn root_repair_signal_ignores_await_successor_and_clears_on_commit() {
let (tx, mut rx) = tokio::sync::watch::channel(VctRootRepairStatus::default());
let mut manager = VctWriteManager::new(tx);
manager.on_retryable_error(Height(42), false, false, queued_block(1));
assert!(
!rx.has_changed().expect("watch channel remains open"),
"await-successor stalls do not need root repair"
);
manager.on_retryable_error(Height(42), true, false, queued_block(2));
assert!(matches!(
rx.borrow_and_update().state,
VctRootRepairState::Unavailable { .. }
));
manager.on_commit_success();
assert_eq!(rx.borrow_and_update().state, VctRootRepairState::Idle);
}
#[test]
fn reset_withdraws_published_root_repair_need() {
let _init_guard = zakura_test::init();
let (tx, mut rx) = tokio::sync::watch::channel(VctRootRepairStatus::default());
let mut manager = VctWriteManager::new(tx);
let mut finalized_state = crate::service::finalized_state::FinalizedState::new(
&crate::Config::ephemeral(),
&zakura_chain::parameters::Network::Mainnet,
)
.expect("opening an ephemeral finalized state succeeds");
manager.on_retryable_error(Height(42), true, false, queued_block(1));
let published = *rx.borrow_and_update();
assert_eq!(
published.state,
VctRootRepairState::Unavailable { height: Height(42) }
);
manager.reset(&mut finalized_state);
let withdrawn = *rx.borrow_and_update();
assert_eq!(
withdrawn.state,
VctRootRepairState::Idle,
"a queue reset must withdraw the published repair need"
);
assert_eq!(
withdrawn.generation, published.generation,
"withdrawing a repair must not consume a new generation"
);
manager.on_retryable_error(Height(42), true, false, queued_block(2));
let republished = *rx.borrow_and_update();
assert_eq!(
republished.state,
VctRootRepairState::Unavailable { height: Height(42) }
);
assert_eq!(republished.generation, published.generation + 1);
}