use std::{collections::HashSet, env, sync::Arc};
use proptest::prelude::*;
use chrono::Duration;
use tokio::{sync::oneshot, time};
use tower::ServiceExt;
use zakura_chain::{
block::{Block, Height},
serialization::ZcashDeserializeInto,
transaction::{Transaction, UnminedTx},
};
use zakura_node_services::mempool::{Gossip, Request, Response};
use zakura_state::{BoxError, ReadRequest, ReadResponse};
use zakura_test::mock_service::MockService;
use crate::queue::{Queue, Runner, CHANNEL_AND_QUEUE_CAPACITY};
const DEFAULT_BLOCK_VEC_PROPTEST_CASES: u32 = 2;
proptest! {
#![proptest_config(
proptest::test_runner::Config::with_cases(env::var("PROPTEST_CASES")
.ok()
.and_then(|v| v.parse().ok())
.unwrap_or(DEFAULT_BLOCK_VEC_PROPTEST_CASES))
)]
#[test]
fn insert_remove_to_from_queue(transaction in any::<UnminedTx>()) {
let (mut runner, _sender) = Queue::start();
runner.queue.insert(transaction.clone());
let queue_transactions = runner.queue.transactions();
prop_assert_eq!(1, queue_transactions.len());
runner.queue.remove(transaction.id);
prop_assert_eq!(runner.queue.transactions().len(), 0);
}
#[test]
fn queue_size_limit(transactions in any::<[UnminedTx; CHANNEL_AND_QUEUE_CAPACITY + 1]>()) {
let (mut runner, _sender) = Queue::start();
transactions.iter().for_each(|t| runner.queue.insert(t.clone()));
let queue_transactions = runner.queue.transactions();
prop_assert_eq!(CHANNEL_AND_QUEUE_CAPACITY, queue_transactions.len());
}
#[test]
fn queue_order(transactions in any::<[UnminedTx; 32]>()) {
let (mut runner, _sender) = Queue::start();
for i in 0..CHANNEL_AND_QUEUE_CAPACITY {
let transaction = transactions[i].clone();
runner.queue.insert(transaction.clone());
let queue_transactions = runner.queue.transactions();
prop_assert_eq!(i + 1, queue_transactions.len());
prop_assert_eq!(UnminedTx::from(queue_transactions[i].0.clone()), transaction);
}
let queue_transactions = runner.queue.transactions();
prop_assert_eq!(CHANNEL_AND_QUEUE_CAPACITY, queue_transactions.len());
for transaction in transactions.iter().skip(CHANNEL_AND_QUEUE_CAPACITY) {
runner.queue.insert(transaction.clone());
let queue_transactions = runner.queue.transactions();
prop_assert_eq!(CHANNEL_AND_QUEUE_CAPACITY, queue_transactions.len());
prop_assert_eq!(UnminedTx::from(queue_transactions.last().unwrap().1.0.clone()), transaction.clone());
}
let queue_transactions = runner.queue.transactions();
for i in 0..CHANNEL_AND_QUEUE_CAPACITY {
let transaction = transactions[(CHANNEL_AND_QUEUE_CAPACITY - 8) + i].clone();
prop_assert_eq!(UnminedTx::from(queue_transactions[i].0.clone()), transaction);
}
}
#[test]
fn remove_expired_transactions_from_queue(transaction in any::<UnminedTx>()) {
let (runtime, _init_guard) = zakura_test::init_async();
runtime.block_on(async move {
time::pause();
let (mut runner, _sender) = Queue::start();
runner.queue.insert(transaction);
prop_assert_eq!(runner.queue.transactions().len(), 1);
let spacing = Duration::seconds(150);
runner.remove_expired(spacing);
prop_assert_eq!(runner.queue.transactions().len(), 1);
time::advance(spacing.to_std().unwrap()).await;
runner.remove_expired(spacing);
prop_assert_eq!(runner.queue.transactions().len(), 1);
time::advance(spacing.to_std().unwrap()).await;
runner.remove_expired(spacing);
prop_assert_eq!(runner.queue.transactions().len(), 1);
time::advance(spacing.to_std().unwrap()).await;
runner.remove_expired(spacing);
prop_assert_eq!(runner.queue.transactions().len(), 1);
time::advance(spacing.to_std().unwrap()).await;
runner.remove_expired(spacing);
prop_assert_eq!(runner.queue.transactions().len(), 1);
time::advance(spacing.to_std().unwrap()).await;
runner.remove_expired(spacing);
prop_assert_eq!(runner.queue.transactions().len(), 1);
time::advance(chrono::Duration::seconds(6).to_std().unwrap()).await;
runner.remove_expired(spacing);
prop_assert_eq!(runner.queue.transactions().len(), 0);
Ok::<_, TestCaseError>(())
})?;
}
#[test]
fn queue_runner_mempool(transaction in any::<Transaction>()) {
let (runtime, _init_guard) = zakura_test::init_async();
runtime.block_on(async move {
let mut mempool = MockService::build().for_prop_tests();
let (mut runner, _sender) = Queue::start();
let unmined_transaction = UnminedTx::from(transaction);
runner.queue.insert(unmined_transaction.clone());
let transactions = runner.queue.transactions();
prop_assert_eq!(transactions.len(), 1);
let transactions_hash_set = runner.transactions_as_hash_set();
let send_task = tokio::spawn(Runner::check_mempool(mempool.clone(), transactions_hash_set.clone()));
let expected_request = Request::TransactionsById(transactions_hash_set.clone());
let response = Response::Transactions(vec![]);
mempool
.expect_request(expected_request)
.await?
.respond(response);
let result = send_task.await.expect("Requesting transactions should not panic");
prop_assert_eq!(result, HashSet::new());
let request = Request::Queue(vec![Gossip::Tx(unmined_transaction.clone())]);
let expected_request = Request::Queue(vec![Gossip::Tx(unmined_transaction.clone())]);
let send_task = tokio::spawn(mempool.clone().oneshot(request));
let (rsp_tx, rsp_rx) = oneshot::channel();
let _ = rsp_tx.send(Ok(()));
let response = Response::Queued(vec![Ok(rsp_rx)]);
mempool
.expect_request(expected_request)
.await?
.respond(response);
let _ = send_task.await.expect("Inserting to mempool should not panic");
let send_task = tokio::spawn(Runner::check_mempool(mempool.clone(), transactions_hash_set.clone()));
let expected_request = Request::TransactionsById(transactions_hash_set);
let response = Response::Transactions(vec![unmined_transaction]);
mempool
.expect_request(expected_request)
.await?
.respond(response);
let result = send_task.await.expect("Requesting transactions should not panic");
prop_assert_eq!(result.len(), 1);
prop_assert_eq!(runner.queue.transactions().len(), 1);
runner.remove_committed(result);
prop_assert_eq!(runner.queue.transactions().len(), 0);
mempool.expect_no_requests().await?;
Ok::<_, TestCaseError>(())
})?;
}
#[test]
fn queue_runner_state(transaction in any::<Transaction>()) {
let (runtime, _init_guard) = zakura_test::init_async();
runtime.block_on(async move {
let mut read_state: MockService<_, _, _, BoxError> = MockService::build().for_prop_tests();
let mut write_state: MockService<_, _, _, BoxError> = MockService::build().for_prop_tests();
let (mut runner, _sender) = Queue::start();
let unmined_transaction = UnminedTx::from(&transaction);
runner.queue.insert(unmined_transaction.clone());
prop_assert_eq!(runner.queue.transactions().len(), 1);
let transactions_hash_set = runner.transactions_as_hash_set();
let send_task = tokio::spawn(Runner::check_state(read_state.clone(), transactions_hash_set.clone()));
let expected_request = ReadRequest::Transaction(transaction.hash());
let response = ReadResponse::Transaction(None);
read_state
.expect_request(expected_request)
.await?
.respond(response);
let result = send_task.await.expect("Requesting transaction should not panic");
prop_assert_eq!(HashSet::new(), result);
let block =
zakura_test::vectors::BLOCK_MAINNET_1_BYTES.zcash_deserialize_into::<Arc<Block>>()?;
let mut block = Arc::try_unwrap(block).expect("block should unwrap");
block.transactions.push(Arc::new(transaction.clone()));
let request = zakura_state::Request::CommitCheckpointVerifiedBlock(zakura_state::CheckpointVerifiedBlock::from(Arc::new(block.clone())));
let send_task = tokio::spawn(write_state.clone().oneshot(request.clone()));
let response = zakura_state::Response::Committed(block.hash());
write_state
.expect_request(request)
.await?
.respond(response);
let _ = send_task.await.expect("Inserting block to state should not panic");
let send_task = tokio::spawn(Runner::check_state(read_state.clone(), transactions_hash_set));
let expected_request = ReadRequest::Transaction(transaction.hash());
let response = ReadResponse::Transaction(Some(zakura_state::MinedTx::new(Arc::new(transaction), Height(1), 1, block.header.time)));
read_state
.expect_request(expected_request)
.await?
.respond(response);
let result = send_task.await.expect("Requesting transaction should not panic");
prop_assert_eq!(result.len(), 1);
read_state.expect_no_requests().await?;
write_state.expect_no_requests().await?;
Ok::<_, TestCaseError>(())
})?;
}
#[test]
fn queue_mempool_retry(transaction in any::<Transaction>()) {
let (runtime, _init_guard) = zakura_test::init_async();
runtime.block_on(async move {
let mut mempool = MockService::build().for_prop_tests();
let (mut runner, _sender) = Queue::start();
let unmined_transaction = UnminedTx::from(transaction.clone());
runner.queue.insert(unmined_transaction.clone());
let transactions = runner.queue.transactions();
prop_assert_eq!(transactions.len(), 1);
let transactions_vec = runner.transactions_as_vec();
let send_task = tokio::spawn(Runner::retry(mempool.clone(), transactions_vec.clone()));
let gossip = Gossip::Tx(UnminedTx::from(transaction.clone()));
let expected_request = Request::Queue(vec![gossip]);
let (rsp_tx, rsp_rx) = oneshot::channel();
let _ = rsp_tx.send(Ok(()));
let response = Response::Queued(vec![Ok(rsp_rx)]);
mempool
.expect_request(expected_request)
.await?
.respond(response);
let result = send_task.await.expect("Requesting transactions should not panic");
prop_assert_eq!(result.len(), 1);
Ok::<_, TestCaseError>(())
})?;
}
}