diff --git a/Cargo.lock b/Cargo.lock index 6a0b1c9989..4256d61d79 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -7795,7 +7795,7 @@ dependencies = [ [[package]] name = "zakura-consensus" -version = "7.0.2-rc0" +version = "8.0.0-rc0" dependencies = [ "blake2b_simd", "chrono", @@ -7895,7 +7895,7 @@ dependencies = [ [[package]] name = "zakura-header-chain" -version = "2.0.1-rc0" +version = "2.0.2-rc0" dependencies = [ "chrono", "criterion", @@ -7961,7 +7961,7 @@ dependencies = [ [[package]] name = "zakura-network" -version = "7.1.0-rc0" +version = "7.1.1-rc0" dependencies = [ "bitflags", "blake2b_simd", @@ -8007,7 +8007,7 @@ dependencies = [ [[package]] name = "zakura-node-services" -version = "3.2.3-rc0" +version = "3.2.4-rc0" dependencies = [ "jsonrpsee-types", "reqwest 0.12.28", @@ -8166,7 +8166,7 @@ dependencies = [ [[package]] name = "zakura-rpc" -version = "9.0.1-rc0" +version = "10.0.0-rc0" dependencies = [ "anyhow", "base64 0.22.1", @@ -8269,7 +8269,7 @@ dependencies = [ [[package]] name = "zakura-script" -version = "3.2.2-rc0" +version = "3.2.3-rc0" dependencies = [ "hex", "lazy_static", @@ -8299,7 +8299,7 @@ dependencies = [ [[package]] name = "zakura-state" -version = "7.2.0-rc0" +version = "8.0.0-rc0" dependencies = [ "bincode", "chrono", @@ -8410,7 +8410,7 @@ dependencies = [ [[package]] name = "zakura-utils" -version = "2.2.3-rc0" +version = "2.2.4-rc0" dependencies = [ "clap", "color-eyre", diff --git a/crates/zakura-consensus/Cargo.toml b/crates/zakura-consensus/Cargo.toml index da4a29447a..ec68d7c9aa 100644 --- a/crates/zakura-consensus/Cargo.toml +++ b/crates/zakura-consensus/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "zakura-consensus" -version = "7.0.2-rc0" +version = "8.0.0-rc0" authors.workspace = true description = "Implementation of Zcash consensus checks for the Zakura node. Internal crate, published to support cargo install zakura" license.workspace = true @@ -67,11 +67,11 @@ zcash_primitives = { workspace = true } tower-fallback = { package = "zakura-tower-fallback", path = "../tower-fallback/", version = "1.2.0" } tower-batch-control = { package = "zakura-tower-batch-control", path = "../tower-batch-control/", version = "1.3.0" } -zakura-script = { path = "../zakura-script", version = "3.2.2-rc0" } -zakura-state = { path = "../zakura-state", version = "7.2.0-rc0" } -zakura-node-services = { path = "../zakura-node-services", version = "3.2.3-rc0" } +zakura-script = { path = "../zakura-script", version = "3.2.3-rc0" } +zakura-state = { path = "../zakura-state", version = "8.0.0-rc0" } +zakura-node-services = { path = "../zakura-node-services", version = "3.2.4-rc0" } zakura-chain = { path = "../zakura-chain", version = "7.0.0-rc0" } -zakura-header-chain = { path = "../zakura-header-chain", version = "2.0.1-rc0" } +zakura-header-chain = { path = "../zakura-header-chain", version = "2.0.2-rc0" } zcash_protocol.workspace = true @@ -93,7 +93,7 @@ toml = { workspace = true } tokio = { workspace = true, features = ["full", "tracing", "test-util"] } -zakura-state = { path = "../zakura-state", version = "7.2.0-rc0", features = ["proptest-impl"] } +zakura-state = { path = "../zakura-state", version = "8.0.0-rc0", features = ["proptest-impl"] } zakura-chain = { path = "../zakura-chain", version = "7.0.0-rc0", features = ["proptest-impl"] } zakura-test = { path = "../zakura-test/", version = "2.1.0" } diff --git a/crates/zakura-consensus/benches/worst_case_tx_verification.rs b/crates/zakura-consensus/benches/worst_case_tx_verification.rs index c93b452af5..c2f7e87248 100644 --- a/crates/zakura-consensus/benches/worst_case_tx_verification.rs +++ b/crates/zakura-consensus/benches/worst_case_tx_verification.rs @@ -466,6 +466,7 @@ fn build_workload(case: &BenchmarkCase, candidates: &[CandidateTx]) -> Option> = Arc::new(known_utxos.keys().map(|outpoint| outpoint.hash).collect()); + // Reject cross-transaction double spends before cloning historical outputs. + let mut spent_outpoints = HashSet::new(); + for outpoint in block + .transactions + .iter() + .flat_map(|tx| tx.spent_outpoints()) + { + if !spent_outpoints.insert(outpoint) { + return Err(crate::error::TransactionError::DuplicateTransparentSpend( + outpoint, + ) + .into()); + } + } + let utxo_resolver = + tx::BlockUtxos::for_block(&block, &known_utxos, state_service.clone()); // Keep this guard after `known_outpoint_hashes` so its `Drop` removes the // pointer-keyed registration before the `Arc` address can be reused. let _block_batch_flush = primitives::register_block_verifier_batch_flush( @@ -367,6 +383,7 @@ where transaction: transaction.clone(), known_outpoint_hashes: known_outpoint_hashes.clone(), known_utxos: known_utxos.clone(), + utxo_resolver: utxo_resolver.clone(), height, time: block.header.time, }); diff --git a/crates/zakura-consensus/src/block/tests.rs b/crates/zakura-consensus/src/block/tests.rs index e5c6737d58..fb6a124a37 100644 --- a/crates/zakura-consensus/src/block/tests.rs +++ b/crates/zakura-consensus/src/block/tests.rs @@ -739,6 +739,122 @@ fn librustzcash_conversion_test_network(network_upgrade: NetworkUpgrade) -> Netw .expect("failed to build configured network") } +#[tokio::test] +async fn block_verification_uses_batched_external_outputs_without_committing_them() { + let _init_guard = zakura_test::init(); + let network = librustzcash_conversion_test_network(NetworkUpgrade::Nu5); + let outpoint = transparent::OutPoint { + hash: [42; 32].into(), + index: 0, + }; + for (script_succeeds, duplicate_spend) in [(true, false), (false, false), (true, true)] { + let mut block: Block = zakura_test::vectors::BLOCK_MAINNET_GENESIS_BYTES + .zcash_deserialize_into() + .unwrap(); + block.transactions = vec![ + Arc::new(v5_coinbase_transaction( + NetworkUpgrade::Nu5, + Height(1), + &network, + )), + Arc::new(Transaction::V5 { + network_upgrade: NetworkUpgrade::Nu5, + inputs: vec![transparent::Input::PrevOut { + outpoint, + unlock_script: transparent::Script::new(&[]), + sequence: u32::MAX, + }], + outputs: vec![transparent::Output { + value: Amount::try_from(1).unwrap(), + lock_script: transparent::Script::new(&[0x51]), + }], + lock_time: LockTime::unlocked(), + expiry_height: Height(2), + sapling_shielded_data: None, + orchard_shielded_data: None, + }), + ]; + if duplicate_spend { + let mut duplicate = (*block.transactions[1]).clone(); + let Transaction::V5 { expiry_height, .. } = &mut duplicate else { + unreachable!() + }; + *expiry_height = Height(3); + block.transactions.push(Arc::new(duplicate)); + } + Arc::make_mut(&mut block.header).merkle_root = block.transactions.iter().collect(); + let expected_hash = block.hash(); + let state = + tower::service_fn(move |request| async move { + Ok::<_, BoxError>(match request { + zs::Request::KnownBlock(_) => zs::Response::KnownBlock(None), + zs::Request::AwaitUtxos(outpoints) => { + assert!( + !duplicate_spend, + "reject duplicate spends before reading outputs" + ); + assert_eq!(outpoints, vec![outpoint]); + zs::Response::Utxos( + [( + outpoint, + transparent::Utxo::new( + transparent::Output { + value: Amount::try_from(1).unwrap(), + lock_script: transparent::Script::new(&[ + if script_succeeds { 0x51 } else { 0x00 }, + ]), + }, + Height(0), + false, + ), + )] + .into_iter() + .collect(), + ) + } + zs::Request::CheckBlockProposalValidity(prepared) => { + assert!(!prepared.new_outputs.contains_key(&outpoint)); + assert_eq!( + prepared.new_outputs, + transparent::new_ordered_outputs( + &prepared.block, + &prepared.transaction_hashes + ) + ); + zs::Response::ValidBlockProposal + } + _ => panic!("semantic block verification must use the block resolver"), + }) + }); + let transaction = transaction::Verifier::new_for_tests(&network, state); + let transaction = Buffer::new(BoxService::new(transaction), 1); + let result = tokio::time::timeout( + std::time::Duration::from_secs(30), + SemanticBlockVerifier::new(&network, state, transaction) + .oneshot(Request::CheckProposal(Arc::new(block))), + ) + .await + .unwrap(); + if duplicate_spend { + assert!( + matches!(result, Err(VerifyBlockError::Transaction( + TransactionError::DuplicateTransparentSpend(found))) if found == outpoint), + "{result:?}" + ); + } else if script_succeeds { + assert_eq!(result.unwrap(), expected_hash); + } else { + assert!( + matches!( + result, + Err(VerifyBlockError::Transaction(TransactionError::Script(_))) + ), + "{result:?}" + ); + } + } +} + fn block_with_librustzcash_conversion_failure( case: LibrustzcashConversionFailure, network: &Network, diff --git a/crates/zakura-consensus/src/transaction.rs b/crates/zakura-consensus/src/transaction.rs index 14bc9f4db7..fe928e8c8d 100644 --- a/crates/zakura-consensus/src/transaction.rs +++ b/crates/zakura-consensus/src/transaction.rs @@ -44,6 +44,8 @@ use zakura_state as zs; use crate::{error::TransactionError, primitives, script, BoxError}; pub mod check; +mod utxo_resolver; +pub use utxo_resolver::BlockUtxos; #[cfg(test)] mod tests; @@ -176,6 +178,9 @@ pub enum Request { known_outpoint_hashes: Arc>, /// Additional UTXOs which are known at the time of verification. known_utxos: Arc>, + /// Shared external-output resolver supplied by the semantic block verifier. + /// Standalone transaction verification can use `None` for legacy state lookups. + utxo_resolver: Option, /// The height of the block containing this transaction. height: block::Height, /// The time that the block was mined. @@ -792,6 +797,7 @@ where let mut spent_outputs: Vec> = vec![None; inputs.len()]; // Stores (input_idx, outpoint) for UTXOs not found in the best chain (fetched from mempool later). let mut spent_mempool_outpoints: Vec<(usize, transparent::OutPoint)> = Vec::new(); + let mut resolved_utxos = None; for (input_idx, input) in inputs.iter().enumerate() { if let transparent::Input::PrevOut { outpoint, .. } = input { @@ -817,6 +823,19 @@ where }; utxo + } else if let Request::Block { + utxo_resolver: Some(resolver), + .. + } = &req + { + if resolved_utxos.is_none() { + resolved_utxos = Some(resolver.resolve().await?); + } + resolved_utxos + .as_ref() + .expect("the resolver completed above") + .get(*outpoint, state.clone()) + .await? } else { let response = state .clone() diff --git a/crates/zakura-consensus/src/transaction/tests.rs b/crates/zakura-consensus/src/transaction/tests.rs index 415f71a058..d8d17a635c 100644 --- a/crates/zakura-consensus/src/transaction/tests.rs +++ b/crates/zakura-consensus/src/transaction/tests.rs @@ -54,6 +54,7 @@ use super::{check, Request, Verifier}; #[cfg(test)] mod prop; +mod utxo_resolver; /// Returns the timeout duration for tests, extended when running under coverage /// instrumentation to account for the performance overhead. @@ -1283,6 +1284,7 @@ async fn dont_skip_verification_of_block_transactions_in_mempool() { ); let make_request = |known_outpoint_hashes| Request::Block { + utxo_resolver: None, transaction_hash: tx_hash, transaction: Arc::new(tx), known_outpoint_hashes, @@ -1696,6 +1698,7 @@ async fn v5_transaction_is_rejected_before_nu5_activation() { assert_eq!( verifier .oneshot(Request::Block { + utxo_resolver: None, transaction_hash: tx.hash(), transaction: Arc::new(tx), known_utxos: Arc::new(HashMap::new()), @@ -1723,6 +1726,7 @@ async fn v5_transaction_is_accepted_after_nu5_activation() { let verif_res = Verifier::new_for_tests(&net, state) .oneshot(Request::Block { + utxo_resolver: None, transaction_hash: tx.hash(), transaction: Arc::new(tx), known_utxos: Arc::new(HashMap::new()), @@ -1777,6 +1781,7 @@ async fn v4_transaction_with_transparent_transfer_is_accepted() { let result = verifier .oneshot(Request::Block { + utxo_resolver: None, transaction_hash: transaction.hash(), transaction: Arc::new(transaction), known_utxos: Arc::new(known_utxos), @@ -1823,6 +1828,7 @@ async fn v4_transaction_with_last_valid_expiry_height() { let result = verifier .oneshot(Request::Block { + utxo_resolver: None, transaction_hash: transaction.hash(), transaction: Arc::new(transaction.clone()), known_utxos: Arc::new(known_utxos), @@ -1870,6 +1876,7 @@ async fn v4_coinbase_transaction_with_low_expiry_height() { let result = verifier .oneshot(Request::Block { + utxo_resolver: None, transaction_hash: transaction.hash(), transaction: Arc::new(transaction.clone()), known_utxos: Arc::new(HashMap::new()), @@ -1919,6 +1926,7 @@ async fn v4_transaction_with_too_low_expiry_height() { let result = verifier .oneshot(Request::Block { + utxo_resolver: None, transaction_hash: transaction.hash(), transaction: Arc::new(transaction.clone()), known_utxos: Arc::new(known_utxos), @@ -1971,6 +1979,7 @@ async fn v4_transaction_with_exceeding_expiry_height() { let result = verifier .oneshot(Request::Block { + utxo_resolver: None, transaction_hash: transaction.hash(), transaction: Arc::new(transaction.clone()), known_utxos: Arc::new(known_utxos), @@ -2026,6 +2035,7 @@ async fn v4_coinbase_transaction_with_exceeding_expiry_height() { let result = verifier .oneshot(Request::Block { + utxo_resolver: None, transaction_hash: transaction.hash(), transaction: Arc::new(transaction.clone()), known_utxos: Arc::new(HashMap::new()), @@ -2079,6 +2089,7 @@ async fn v4_coinbase_transaction_is_accepted() { let result = verifier .oneshot(Request::Block { + utxo_resolver: None, transaction_hash: transaction.hash(), transaction: Arc::new(transaction), known_utxos: Arc::new(HashMap::new()), @@ -2136,6 +2147,7 @@ async fn v4_transaction_with_transparent_transfer_is_rejected_by_the_script() { let result = verifier .oneshot(Request::Block { + utxo_resolver: None, transaction_hash: transaction.hash(), transaction: Arc::new(transaction), known_utxos: Arc::new(known_utxos), @@ -2192,6 +2204,7 @@ async fn v4_transaction_with_conflicting_transparent_spend_is_rejected() { let result = verifier .oneshot(Request::Block { + utxo_resolver: None, transaction_hash: transaction.hash(), transaction: Arc::new(transaction), known_utxos: Arc::new(known_utxos), @@ -2262,6 +2275,7 @@ fn v4_transaction_with_conflicting_sprout_nullifier_inside_joinsplit_is_rejected let result = verifier .oneshot(Request::Block { + utxo_resolver: None, transaction_hash: transaction.hash(), transaction: Arc::new(transaction), known_utxos: Arc::new(HashMap::new()), @@ -2337,6 +2351,7 @@ fn v4_transaction_with_conflicting_sprout_nullifier_across_joinsplits_is_rejecte let result = verifier .oneshot(Request::Block { + utxo_resolver: None, transaction_hash: transaction.hash(), transaction: Arc::new(transaction), known_utxos: Arc::new(HashMap::new()), @@ -2398,6 +2413,7 @@ async fn v5_transaction_with_transparent_transfer_is_accepted() { let result = verifier .oneshot(Request::Block { + utxo_resolver: None, transaction_hash: transaction.hash(), transaction: Arc::new(transaction), known_utxos: Arc::new(known_utxos), @@ -2446,6 +2462,7 @@ async fn v5_transaction_with_last_valid_expiry_height() { let result = verifier .oneshot(Request::Block { + utxo_resolver: None, transaction_hash: transaction.hash(), transaction: Arc::new(transaction.clone()), known_utxos: Arc::new(known_utxos), @@ -2493,6 +2510,7 @@ async fn v5_coinbase_transaction_expiry_height() { let result = verifier .clone() .oneshot(Request::Block { + utxo_resolver: None, transaction_hash: transaction.hash(), transaction: Arc::new(transaction.clone()), known_utxos: Arc::new(HashMap::new()), @@ -2516,6 +2534,7 @@ async fn v5_coinbase_transaction_expiry_height() { let result = verifier .clone() .oneshot(Request::Block { + utxo_resolver: None, transaction_hash: transaction.hash(), transaction: Arc::new(new_transaction.clone()), known_utxos: Arc::new(HashMap::new()), @@ -2547,6 +2566,7 @@ async fn v5_coinbase_transaction_expiry_height() { let result = verifier .clone() .oneshot(Request::Block { + utxo_resolver: None, transaction_hash: transaction.hash(), transaction: Arc::new(new_transaction.clone()), known_utxos: Arc::new(HashMap::new()), @@ -2587,6 +2607,7 @@ async fn v5_coinbase_transaction_expiry_height() { let verification_result = verifier .clone() .oneshot(Request::Block { + utxo_resolver: None, transaction_hash: transaction.hash(), transaction: Arc::new(new_transaction.clone()), known_utxos: Arc::new(HashMap::new()), @@ -2640,6 +2661,7 @@ async fn v5_transaction_with_too_low_expiry_height() { let result = verifier .oneshot(Request::Block { + utxo_resolver: None, transaction_hash: transaction.hash(), transaction: Arc::new(transaction.clone()), known_utxos: Arc::new(known_utxos), @@ -2691,6 +2713,7 @@ async fn v5_transaction_with_exceeding_expiry_height() { let verification_result = Verifier::new_for_tests(&Network::Mainnet, state) .oneshot(Request::Block { + utxo_resolver: None, transaction_hash: transaction.hash(), transaction: Arc::new(transaction.clone()), known_utxos: Arc::new(known_utxos), @@ -2747,6 +2770,7 @@ async fn v5_coinbase_transaction_is_accepted() { let result = verifier .oneshot(Request::Block { + utxo_resolver: None, transaction_hash: transaction.hash(), transaction: Arc::new(transaction), known_utxos: Arc::new(known_utxos), @@ -2806,6 +2830,7 @@ async fn v5_transaction_with_transparent_transfer_is_rejected_by_the_script() { let result = verifier .oneshot(Request::Block { + utxo_resolver: None, transaction_hash: transaction.hash(), transaction: Arc::new(transaction), known_utxos: Arc::new(known_utxos), @@ -2855,6 +2880,7 @@ async fn v5_transaction_with_conflicting_transparent_spend_is_rejected() { let verification_result = Verifier::new_for_tests(&network, state) .oneshot(Request::Block { + utxo_resolver: None, transaction_hash: transaction.hash(), transaction: Arc::new(transaction), known_utxos: Arc::new(known_utxos), @@ -2900,6 +2926,7 @@ fn v4_with_signed_sprout_transfer_is_accepted() { // Test the transaction verifier let result = verifier .oneshot(Request::Block { + utxo_resolver: None, transaction_hash: transaction.hash(), transaction, known_utxos: Arc::new(HashMap::new()), @@ -2988,6 +3015,7 @@ async fn v4_with_joinsplit_is_rejected_for_modification( let result = verifier .clone() .oneshot(Request::Block { + utxo_resolver: None, transaction_hash: transaction.hash(), transaction: transaction.clone(), known_utxos: Arc::new(HashMap::new()), @@ -3040,6 +3068,7 @@ fn v4_and_v5_with_sapling_spends() { let result = timeout( test_timeout(), verifier.oneshot(Request::Block { + utxo_resolver: None, transaction_hash: transaction.hash(), transaction, known_utxos: Arc::new(HashMap::new()), @@ -3080,6 +3109,7 @@ fn v4_and_v5_with_sapling_spends() { timeout( test_timeout(), verifier.oneshot(Request::Block { + utxo_resolver: None, transaction_hash: tx.hash(), transaction: Arc::new(tx), known_utxos: Arc::new(HashMap::new()), @@ -3277,6 +3307,7 @@ fn v4_with_invalid_sapling_proof_returns_typed_error() { let block_transaction_result = timeout( test_timeout(), transaction_verifier.oneshot(Request::Block { + utxo_resolver: None, transaction_hash, transaction, known_utxos: Arc::new(HashMap::new()), @@ -3383,6 +3414,7 @@ fn v4_with_duplicate_sapling_spends() { // Test the transaction verifier let result = verifier .oneshot(Request::Block { + utxo_resolver: None, transaction_hash: transaction.hash(), transaction, known_utxos: Arc::new(HashMap::new()), @@ -3429,6 +3461,7 @@ fn v4_with_sapling_outputs_and_no_spends() { // Test the transaction verifier let result = verifier .oneshot(Request::Block { + utxo_resolver: None, transaction_hash: transaction.hash(), transaction, known_utxos: Arc::new(HashMap::new()), @@ -3484,6 +3517,7 @@ fn sapling_output_with_invalid_ephemeral_key_is_rejected() { let result = verifier .oneshot(Request::Block { + utxo_resolver: None, transaction_hash: transaction.hash(), transaction, known_utxos: Arc::new(HashMap::new()), @@ -3554,6 +3588,7 @@ fn sapling_v4_output_with_invalid_value_commitment_is_rejected_after_roundtrip() let result = verifier .oneshot(Request::Block { + utxo_resolver: None, transaction_hash: transaction.hash(), transaction, known_utxos: Arc::new(HashMap::new()), @@ -3639,6 +3674,7 @@ fn sapling_spends_with_invalid_value_commitments_are_rejected_after_roundtrip() let result = verifier .oneshot(Request::Block { + utxo_resolver: None, transaction_hash: transaction.hash(), transaction, known_utxos: Arc::new(HashMap::new()), @@ -3854,6 +3890,7 @@ async fn v5_with_duplicate_sapling_spends() { assert_eq!( verifier .oneshot(Request::Block { + utxo_resolver: None, transaction_hash: tx.hash(), transaction: Arc::new(tx), known_utxos: Arc::new(HashMap::new()), @@ -3916,6 +3953,7 @@ async fn v5_with_duplicate_orchard_action() { assert_eq!( verifier .oneshot(Request::Block { + utxo_resolver: None, transaction_hash: tx.hash(), transaction: Arc::new(tx), known_utxos: Arc::new(HashMap::new()), @@ -4037,6 +4075,7 @@ async fn orchard_disabling_soft_fork_rejects_orchard_actions_in_blocks_and_mempo service_fn(|_| async { unreachable!("state service should not be called") }), ) .oneshot(Request::Block { + utxo_resolver: None, transaction_hash: tx.hash(), transaction: Arc::new(tx.clone()), known_utxos: Arc::new(HashMap::new()), @@ -4226,6 +4265,7 @@ fn orchard_disabling_soft_fork_accepts_orchard_actions_below_activation_height() let accept_response = accept_verifier .oneshot(Request::Block { + utxo_resolver: None, transaction_hash: tx.hash(), transaction: Arc::new(tx.clone()), known_utxos: Arc::new(HashMap::new()), @@ -4252,6 +4292,7 @@ fn orchard_disabling_soft_fork_accepts_orchard_actions_below_activation_height() service_fn(|_| async { unreachable!("state service should not be called") }), ) .oneshot(Request::Block { + utxo_resolver: None, transaction_hash: tx.hash(), transaction: Arc::new(tx), known_utxos: Arc::new(HashMap::new()), @@ -4388,6 +4429,7 @@ async fn v5_consensus_branch_ids() { let block_req = verifier .clone() .oneshot(Request::Block { + utxo_resolver: None, transaction_hash: tx.hash(), transaction: Arc::new(tx.clone()), known_utxos: known_utxos.clone(), @@ -4456,6 +4498,7 @@ async fn v5_consensus_branch_ids() { let block_req = verifier .clone() .oneshot(Request::Block { + utxo_resolver: None, transaction_hash: tx.hash(), transaction: Arc::new(tx.clone()), known_utxos: known_utxos.clone(), @@ -4515,6 +4558,7 @@ async fn v5_consensus_branch_ids() { let block_req = verifier .clone() .oneshot(Request::Block { + utxo_resolver: None, transaction_hash: tx.hash(), transaction: Arc::new(tx.clone()), known_utxos: known_utxos.clone(), @@ -5483,6 +5527,7 @@ async fn block_with_garbage_orchard_proofs_is_rejected() { let resp = verifier .clone() .oneshot(Request::Block { + utxo_resolver: None, transaction_hash: tx_hash, transaction: Arc::new(garbage_tx), known_outpoint_hashes: Arc::new([input_outpoint.hash].into()), @@ -5559,6 +5604,7 @@ async fn block_transaction_past_expiry_height_is_rejected() { let result = timeout( test_timeout(), verifier.clone().oneshot(Request::Block { + utxo_resolver: None, transaction_hash: tx_hash, transaction: Arc::new(tx.clone()), known_outpoint_hashes: Arc::new([input_outpoint.hash].into()), @@ -6014,6 +6060,7 @@ async fn verify( /// Returns a block request that mines `tx` at `height`. fn block_request(tx: &Transaction, height: block::Height) -> Request { Request::Block { + utxo_resolver: None, transaction_hash: tx.hash(), transaction: Arc::new(tx.clone()), known_outpoint_hashes: Arc::new(HashSet::new()), diff --git a/crates/zakura-consensus/src/transaction/tests/prop.rs b/crates/zakura-consensus/src/transaction/tests/prop.rs index a2dc56bed9..64698ffcc5 100644 --- a/crates/zakura-consensus/src/transaction/tests/prop.rs +++ b/crates/zakura-consensus/src/transaction/tests/prop.rs @@ -509,6 +509,7 @@ fn validate( verifier .clone() .oneshot(transaction::Request::Block { + utxo_resolver: None, transaction_hash, transaction: Arc::new(transaction), known_utxos: Arc::new(known_utxos), diff --git a/crates/zakura-consensus/src/transaction/tests/utxo_resolver.rs b/crates/zakura-consensus/src/transaction/tests/utxo_resolver.rs new file mode 100644 index 0000000000..22fbeb0604 --- /dev/null +++ b/crates/zakura-consensus/src/transaction/tests/utxo_resolver.rs @@ -0,0 +1,478 @@ +//! Regression checks for the block-owned resolver. + +use super::*; +use crate::transaction::BlockUtxos; + +fn fixture( + input_count: u32, + known_every: u32, +) -> ( + Block, + HashMap, + HashMap, +) { + let mut inputs = Vec::new(); + let mut known = HashMap::new(); + let mut expected = HashMap::new(); + for index in 0..input_count { + let (input, _, utxos) = mock_transparent_transfer( + Height(1), + true, + index, + Amount::try_from(u64::from(index) + 1).unwrap(), + ); + inputs.push(input); + for (outpoint, utxo) in utxos { + expected.insert(outpoint, utxo.utxo.clone()); + if known_every != 0 && index % known_every == 0 { + known.insert(outpoint, utxo); + } + } + } + let mut block: Block = zakura_test::vectors::BLOCK_MAINNET_1_BYTES + .zcash_deserialize_into() + .unwrap(); + let second_half = inputs.split_off(inputs.len() / 2); + block.transactions = [inputs, second_half] + .into_iter() + .map(|inputs| { + Arc::new(Transaction::V5 { + network_upgrade: NetworkUpgrade::Nu5, + inputs, + outputs: vec![transparent::Output { + value: Amount::try_from(1).unwrap(), + lock_script: transparent::Script::new(&[0x51]), + }], + lock_time: LockTime::unlocked(), + expiry_height: Height(2_000_000), + sapling_shielded_data: None, + orchard_shielded_data: None, + }) + }) + .collect(); + (block, known, expected) +} + +fn request( + tx: Arc, + known: Arc>, + resolver: Option, +) -> Request { + Request::Block { + transaction_hash: tx.hash(), + transaction: tx, + known_utxos: known, + known_outpoint_hashes: Arc::new(HashSet::new()), + utxo_resolver: resolver, + height: NetworkUpgrade::Nu5 + .activation_height(&Network::Mainnet) + .unwrap(), + time: Utc::now(), + } +} + +#[tokio::test] +async fn block_resolver_shares_bounded_batches_and_preserves_v5_sighashes() { + for input_count in [0, 1, 2, 63, 64, 65, 257, 513] { + for known_every in [0, 1, 3] { + let (mut block, known, expected) = fixture(input_count, known_every); + // A second consumer of the same transaction must not repeat its lookups. + block.transactions.push(block.transactions[1].clone()); + let (requests, mut received) = tokio::sync::mpsc::unbounded_channel(); + let state = tower::service_fn(move |request| { + let zakura_state::Request::AwaitUtxos(outpoints) = request else { + panic!("the shared resolver must use batched state requests") + }; + let (send, receive) = tokio::sync::oneshot::channel(); + requests.send((outpoints, send)).unwrap(); + async move { receive.await.unwrap() } + }); + let resolver = BlockUtxos::for_block(&block, &known, state.clone()); + let mut remaining = expected.len() - known.len(); + assert_eq!(resolver.is_none(), remaining == 0); + let known = Arc::new(known); + let lookups = futures::future::join_all(block.transactions.iter().map(|tx| { + Verifier::< + _, + tower::util::BoxCloneService, + >::spent_utxos( + tx.clone(), + request(tx.clone(), known.clone(), resolver.clone()), + tower::timeout::Timeout::new(state.clone(), super::super::UTXO_LOOKUP_TIMEOUT), + None, + ) + })); + futures::pin_mut!(lookups); + let mut seen = HashSet::new(); + while remaining > 0 { + tokio::task::yield_now().await; + assert!(futures::poll!(&mut lookups).is_pending()); + let mut batches = Vec::new(); + while let Ok(batch) = received.try_recv() { + batches.push(batch); + } + assert_eq!(batches.len(), remaining.div_ceil(64).min(4)); + for (outpoints, send) in batches.into_iter().rev() { + assert!(!outpoints.is_empty() && outpoints.len() <= 64); + remaining -= outpoints.len(); + let response = outpoints + .into_iter() + .rev() + .map(|outpoint| { + assert!( + seen.insert(outpoint), + "resolve each external outpoint once per block" + ); + assert!(!known.contains_key(&outpoint)); + (outpoint, expected[&outpoint].clone()) + }) + .collect(); + send.send(Ok(zakura_state::Response::Utxos(response))) + .unwrap(); + } + } + let results = timeout(test_timeout(), lookups).await.unwrap(); + for (tx, result) in block.transactions.iter().zip(results) { + let (utxos, outputs, mempool_outpoints) = result.unwrap(); + let expected_outputs: Vec<_> = tx + .spent_outpoints() + .map(|outpoint| expected[&outpoint].output.clone()) + .collect(); + assert_eq!(outputs, expected_outputs); + assert!(mempool_outpoints.is_empty()); + assert_eq!(utxos.len(), tx.inputs().len()); + let sighash = |outputs| { + zakura_chain::transaction::SigHasher::new( + tx, + NetworkUpgrade::Nu5, + Arc::new(outputs), + ) + .unwrap() + .sighash(HashType::ALL, None) + }; + assert_eq!(sighash(outputs), sighash(expected_outputs.clone())); + if expected_outputs.len() > 1 { + let mut reversed = expected_outputs.clone(); + reversed.reverse(); + assert_ne!(sighash(reversed), sighash(expected_outputs)); + } + } + assert!(received.try_recv().is_err()); + } + } +} + +#[tokio::test(start_paused = true)] +async fn block_resolver_errors_timeout_and_drop_cancel_batches() { + for failure in ["error", "timeout", "drop", "incomplete"] { + let (block, known, _) = fixture(513, 0); + let (requests, mut received) = tokio::sync::mpsc::unbounded_channel(); + let state = tower::service_fn(move |_| { + let (send, receive) = tokio::sync::oneshot::channel(); + requests.send(send).unwrap(); + async move { receive.await.unwrap() } + }); + let resolver = BlockUtxos::for_block(&block, &known, state).unwrap(); + let consumer = resolver.clone(); + let mut lookup = Box::pin(async move { consumer.resolve().await }); + assert!(futures::poll!(&mut lookup).is_pending()); + let mut pending = Vec::new(); + while let Ok(send) = received.try_recv() { + pending.push(send); + } + assert_eq!(pending.len(), 4); + match failure { + "error" => pending + .pop() + .unwrap() + .send(Err("batch failed".into())) + .unwrap(), + "incomplete" => pending + .pop() + .unwrap() + .send(Ok(zakura_state::Response::Utxos(HashMap::new()))) + .unwrap(), + "timeout" => tokio::time::advance(super::super::UTXO_LOOKUP_TIMEOUT).await, + "drop" => {} + _ => unreachable!(), + } + if failure != "drop" { + let error = lookup.as_mut().await.unwrap_err(); + if failure == "error" { + assert!(error.to_string().contains("batch failed")); + } else { + assert_eq!(error, TransactionError::TransparentInputNotFound); + } + } + drop(lookup); + drop(resolver); + assert!(pending.iter().all(tokio::sync::oneshot::Sender::is_closed)); + assert!(received.try_recv().is_err()); + } +} + +#[tokio::test] +async fn block_resolver_keeps_quick_rejection_and_transaction_results() { + let (block, known, expected) = fixture(4, 0); + let expected = Arc::new(expected); + let state = tower::service_fn(move |request| { + let response = match request { + zakura_state::Request::AwaitUtxo(outpoint) => { + zakura_state::Response::Utxo(expected[&outpoint].clone()) + } + zakura_state::Request::AwaitUtxos(outpoints) => zakura_state::Response::Utxos( + outpoints + .into_iter() + .map(|outpoint| (outpoint, expected[&outpoint].clone())) + .collect(), + ), + _ => panic!("transparent block verification only needs UTXO lookups"), + }; + async move { Ok::<_, BoxError>(response) } + }); + let resolver = BlockUtxos::for_block(&block, &known, state.clone()); + for tx in &block.transactions { + let baseline = Verifier::new_for_tests(&Network::Mainnet, state.clone()) + .oneshot(request(tx.clone(), Arc::new(known.clone()), None)) + .await; + let batched = Verifier::new_for_tests(&Network::Mainnet, state.clone()) + .oneshot(request( + tx.clone(), + Arc::new(known.clone()), + resolver.clone(), + )) + .await; + assert!(baseline.is_ok(), "{baseline:?}"); + assert_eq!(batched, baseline); + } + + let mut invalid = block.clone(); + let Transaction::V5 { inputs, .. } = Arc::make_mut(&mut invalid.transactions[0]) else { + unreachable!() + }; + inputs.push(inputs[0].clone()); + let state = tower::service_fn(|_| async { + Err::("quick rejection queried state".into()) + }); + let resolver = BlockUtxos::for_block(&invalid, &known, state); + let error = Verifier::new_for_tests(&Network::Mainnet, state) + .oneshot(request( + invalid.transactions[0].clone(), + Arc::new(known), + resolver, + )) + .await + .unwrap_err(); + assert!(matches!( + error, + TransactionError::DuplicateTransparentSpend(_) + )); +} + +#[tokio::test] +async fn block_resolver_refills_batches_while_earlier_batches_wait() { + let (block, known, expected) = fixture(513, 0); + let (requests, mut received) = tokio::sync::mpsc::unbounded_channel(); + let state = tower::service_fn(move |request| { + let zakura_state::Request::AwaitUtxos(outpoints) = request else { + unreachable!() + }; + let (send, receive) = tokio::sync::oneshot::channel(); + requests.send((outpoints, send)).unwrap(); + async move { receive.await.unwrap() } + }); + let resolver = BlockUtxos::for_block(&block, &known, state).unwrap(); + let consumer = resolver.clone(); + let mut lookup = Box::pin(async move { consumer.resolve().await }); + assert!(futures::poll!(&mut lookup).is_pending()); + let mut pending = Vec::new(); + while let Ok(batch) = received.try_recv() { + pending.push(batch); + } + assert_eq!(pending.len(), 4); + for _ in 0..5 { + let (outpoints, send) = pending.pop().unwrap(); + send.send(Ok(zakura_state::Response::Utxos( + outpoints + .into_iter() + .map(|outpoint| (outpoint, expected[&outpoint].clone())) + .collect(), + ))) + .unwrap(); + tokio::task::yield_now().await; + assert!(futures::poll!(&mut lookup).is_pending()); + pending.push( + received + .try_recv() + .expect("one completed batch must free one slot"), + ); + assert!(received.try_recv().is_err()); + } + drop(lookup); + assert!( + pending.iter().all(|(_, send)| !send.is_closed()), + "the block still owns the resolver" + ); + drop(resolver); + assert!(pending.iter().all(|(_, send)| send.is_closed())); +} + +#[tokio::test(start_paused = true)] +async fn block_resolver_preserves_inner_state_timeout_errors() { + let (block, known, _) = fixture(1, 0); + let state = tower::service_fn(|_| { + futures::future::pending::>() + }); + let state = tower::timeout::Timeout::new(state, std::time::Duration::from_secs(1)); + let resolver = BlockUtxos::for_block(&block, &known, state).unwrap(); + let lookup = resolver.resolve(); + futures::pin_mut!(lookup); + assert!(futures::poll!(&mut lookup).is_pending()); + tokio::time::advance(std::time::Duration::from_secs(1)).await; + assert_eq!( + lookup.await.unwrap_err(), + TransactionError::TransparentInputNotFound + ); +} + +#[tokio::test(start_paused = true)] +async fn block_resolver_falls_back_without_losing_large_outputs() { + use super::super::utxo_resolver::MAX_CACHED_UTXO_SCRIPT_BYTES; + use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering}; + + for exceeds_limit in [false, true] { + let (block, known, mut expected) = fixture(1024, 0); + let script = vec![0x51; MAX_CACHED_UTXO_SCRIPT_BYTES / 1024 + usize::from(exceeds_limit)]; + for utxo in expected.values_mut() { + utxo.output.lock_script = transparent::Script::new(&script); + } + let expected = Arc::new(expected); + let single_requests = Arc::new(AtomicUsize::new(0)); + let stall = Arc::new(AtomicBool::new(false)); + let state = tower::service_fn({ + let expected = expected.clone(); + let stall = stall.clone(); + let single_requests = single_requests.clone(); + move |req| { + let response = match req { + zakura_state::Request::AwaitUtxos(outpoints) => { + if outpoints.len() == 1 { + single_requests.fetch_add(1, Ordering::SeqCst); + } + zakura_state::Response::Utxos( + outpoints + .into_iter() + .map(|outpoint| (outpoint, expected[&outpoint].clone())) + .collect(), + ) + } + _ => panic!("unexpected state request"), + }; + let stall = stall.load(Ordering::SeqCst); + async move { + if stall { + futures::future::pending::<()>().await; + } + Ok::<_, BoxError>(response) + } + } + }); + let resolver = BlockUtxos::for_block(&block, &known, state.clone()); + for tx in &block.transactions { + let (utxos, outputs, _) = Verifier::< + _, + tower::util::BoxCloneService, + >::spent_utxos( + tx.clone(), + request(tx.clone(), Arc::new(known.clone()), resolver.clone()), + tower::timeout::Timeout::new(state.clone(), super::super::UTXO_LOOKUP_TIMEOUT), + None, + ) + .await + .unwrap(); + for (index, outpoint) in tx.spent_outpoints().enumerate() { + assert_eq!(utxos[&outpoint], expected[&outpoint]); + assert_eq!(outputs[index], expected[&outpoint].output); + } + } + assert_eq!( + single_requests.load(Ordering::SeqCst), + if exceeds_limit { 1024 } else { 0 } + ); + if exceeds_limit { + stall.store(true, Ordering::SeqCst); + tokio::time::advance( + super::super::UTXO_LOOKUP_TIMEOUT - std::time::Duration::from_secs(1), + ) + .await; + let start = tokio::time::Instant::now(); + let tx = block.transactions[0].clone(); + let error = Verifier::< + _, + tower::util::BoxCloneService, + >::spent_utxos( + tx.clone(), + request(tx, Arc::new(known), resolver), + tower::timeout::Timeout::new(state, super::super::UTXO_LOOKUP_TIMEOUT), + None, + ) + .await + .unwrap_err(); + assert_eq!(error, TransactionError::TransparentInputNotFound); + assert_eq!(start.elapsed(), std::time::Duration::from_secs(1)); + } + } +} + +#[tokio::test] +async fn block_resolver_script_size_check_matches_interpreter() { + let (mut block, known, expected) = fixture(2, 0); + block.transactions.truncate(1); + let outpoint = block.transactions[0].spent_outpoints().next().unwrap(); + let mut script = Vec::new(); + // Push and drop 520-byte items without exceeding the stack or opcode limits. + for _ in 0..19 { + script.extend([0x4d, 0x08, 0x02]); + script.extend([0; 520]); + script.push(0x75); + } + script.push(41); + script.extend([0; 41]); + script.extend([0x75, 0x51]); + assert_eq!(script.len(), 10_000); + for too_large in [false, true] { + if too_large { + script.push(0x61); + } + let mut utxo = expected[&outpoint].clone(); + utxo.output.lock_script = transparent::Script::new(&script); + let verifier = zakura_script::CachedFfiTransaction::new( + block.transactions[0].clone(), + Arc::new(vec![utxo.output.clone()]), + NetworkUpgrade::Nu5, + ) + .unwrap(); + assert_eq!(verifier.is_valid(0).is_err(), too_large); + let state = tower::service_fn(move |_| { + let utxo = utxo.clone(); + async move { + Ok::<_, BoxError>(zakura_state::Response::Utxos(HashMap::from([( + outpoint, utxo, + )]))) + } + }); + let resolver = BlockUtxos::for_block(&block, &known, state).unwrap(); + let result = resolver.resolve().await; + if too_large { + assert_eq!( + result.unwrap_err(), + TransactionError::Script(zakura_script::Error::ScriptInvalid) + ); + } else { + assert!(matches!( + result.unwrap(), + super::super::utxo_resolver::ResolvedUtxos::Cached(_) + )); + } + } +} diff --git a/crates/zakura-consensus/src/transaction/utxo_resolver.rs b/crates/zakura-consensus/src/transaction/utxo_resolver.rs new file mode 100644 index 0000000000..967f1af3d8 --- /dev/null +++ b/crates/zakura-consensus/src/transaction/utxo_resolver.rs @@ -0,0 +1,190 @@ +//! One cancellable UTXO resolution shared by transactions in a block. + +use std::{ + collections::{HashMap, HashSet}, + fmt, + sync::Arc, +}; + +use futures::{ + future::{BoxFuture, Shared}, + FutureExt, StreamExt, +}; +use tower::{Service, ServiceExt}; + +use zakura_chain::{block, transparent}; +use zakura_state as zs; + +use crate::{error::TransactionError, BoxError}; + +/// Maximum outstanding batches per block, including missing-output waits. +/// The state separately bounds running database reads across all blocks. +const MAX_IN_FLIGHT_UTXO_BATCHES: usize = 4; + +// This bounds retained script bytes without introducing a consensus limit. +// Larger dependency sets use individual lookups under the same deadline. +pub(super) const MAX_CACHED_UTXO_SCRIPT_BYTES: usize = 8 * 1024 * 1024; + +// Both zcash script interpreters reject scripts above MAX_SCRIPT_SIZE before execution. +const MAX_SPENDABLE_SCRIPT_BYTES: usize = 10_000; + +#[derive(Clone, Debug)] +pub(super) enum ResolvedUtxos { + Cached(Arc>), + Individual(tokio::time::Instant), +} + +/// A block-owned lookup shared by its transactions. +/// +/// The resolver caches only informational outputs for this verification attempt. +/// Contextual validation must still check spends against the block's actual chain. +/// Dropping every clone cancels unresolved requests, but not already-running disk reads. +#[derive(Clone)] +pub struct BlockUtxos(Arc>>>); + +impl fmt::Debug for BlockUtxos { + fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + f.debug_struct("BlockUtxos").finish_non_exhaustive() + } +} + +impl PartialEq for BlockUtxos { + fn eq(&self, other: &Self) -> bool { + Arc::ptr_eq(&self.0, &other.0) + } +} + +impl Eq for BlockUtxos {} + +impl BlockUtxos { + pub(crate) fn for_block( + block: &block::Block, + known_utxos: &HashMap, + state: S, + ) -> Option + where + S: Service + Send + Clone + 'static, + S::Future: Send + 'static, + { + let mut seen = HashSet::new(); + let outpoints: Vec<_> = block + .transactions + .iter() + .flat_map(|tx| tx.spent_outpoints()) + .filter(|outpoint| !known_utxos.contains_key(outpoint) && seen.insert(*outpoint)) + .collect(); + if outpoints.is_empty() { + return None; + } + + let lookup = async move { + let deadline = tokio::time::Instant::now() + super::UTXO_LOOKUP_TIMEOUT; + let resolve = async move { + let mut batches = futures::stream::iter(outpoints) + .chunks(zs::constants::MAX_UTXO_BATCH_SIZE) + .map(move |outpoints| { + let state = state.clone(); + async move { + let response = state + .oneshot(zs::Request::AwaitUtxos(outpoints.clone())) + .await + .map_err(|error| { + match error.downcast::() { + Ok(_) => TransactionError::TransparentInputNotFound, + Err(error) => TransactionError::from(error), + } + })?; + let zs::Response::Utxos(utxos) = response else { + unreachable!("AwaitUtxos returns Utxos") + }; + // Never let an incomplete state response omit a sighash input. + if utxos.len() != outpoints.len() + || outpoints + .iter() + .any(|outpoint| !utxos.contains_key(outpoint)) + { + return Err(TransactionError::TransparentInputNotFound); + } + Ok(utxos) + } + }) + .buffer_unordered(MAX_IN_FLIGHT_UTXO_BATCHES); + let mut resolved = HashMap::new(); + let mut script_bytes = 0usize; + while let Some(batch) = batches.next().await { + let batch = batch?; + for utxo in batch.values() { + check_script_size(utxo)?; + script_bytes = script_bytes + .saturating_add(utxo.output.lock_script.as_raw_bytes().len()); + } + if script_bytes > MAX_CACHED_UTXO_SCRIPT_BYTES { + return Ok(ResolvedUtxos::Individual(deadline)); + } + resolved.extend(batch); + } + Ok(ResolvedUtxos::Cached(Arc::new(resolved))) + }; + // One deadline covers admission, reads, and dependency waits for this block. + // Unlike serial per-input deadlines, later batches do not get extra time. + tokio::time::timeout_at(deadline, resolve) + .await + .map_err(|_| TransactionError::TransparentInputNotFound)? + }; + Some(Self(Arc::new(lookup.boxed().shared()))) + } + + pub(super) async fn resolve(&self) -> Result { + self.0.as_ref().clone().await + } +} + +/// Reject outputs that the script interpreter cannot spend before copying their scripts. +fn check_script_size(utxo: &transparent::Utxo) -> Result<(), TransactionError> { + if utxo.output.lock_script.as_raw_bytes().len() > MAX_SPENDABLE_SCRIPT_BYTES { + return Err(zakura_script::Error::ScriptInvalid.into()); + } + Ok(()) +} + +impl ResolvedUtxos { + pub(super) async fn get( + &self, + outpoint: transparent::OutPoint, + state: S, + ) -> Result + where + S: Service + Send + Clone + 'static, + S::Future: Send + 'static, + { + let deadline = match self { + Self::Cached(utxos) => { + return utxos + .get(&outpoint) + .cloned() + .ok_or(TransactionError::TransparentInputNotFound) + } + Self::Individual(deadline) => deadline, + }; + let response = tokio::time::timeout_at( + *deadline, + state.oneshot(zs::Request::AwaitUtxos(vec![outpoint])), + ) + .await + .map_err(|_| TransactionError::TransparentInputNotFound)? + .map_err( + |error| match error.downcast::() { + Ok(_) => TransactionError::TransparentInputNotFound, + Err(error) => TransactionError::from(error), + }, + )?; + let zs::Response::Utxos(mut utxos) = response else { + unreachable!("AwaitUtxos returns Utxos") + }; + let utxo = utxos + .remove(&outpoint) + .ok_or(TransactionError::TransparentInputNotFound)?; + check_script_size(&utxo)?; + Ok(utxo) + } +} diff --git a/crates/zakura-header-chain/Cargo.toml b/crates/zakura-header-chain/Cargo.toml index e22a6070be..32037d2df4 100644 --- a/crates/zakura-header-chain/Cargo.toml +++ b/crates/zakura-header-chain/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "zakura-header-chain" -version = "2.0.1-rc0" +version = "2.0.2-rc0" authors.workspace = true description = "Fork-aware header-chain domain types and transition engine for Zakura" license.workspace = true diff --git a/crates/zakura-network/Cargo.toml b/crates/zakura-network/Cargo.toml index dadba870e4..e4cc2cf508 100644 --- a/crates/zakura-network/Cargo.toml +++ b/crates/zakura-network/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "zakura-network" -version = "7.1.0-rc0" +version = "7.1.1-rc0" authors.workspace = true description = "Networking code for the Zakura node. Internal crate, published to support cargo install zakura" # # Legal @@ -83,9 +83,9 @@ proptest = { workspace = true, optional = true } proptest-derive = { workspace = true, optional = true } zakura-chain = { path = "../zakura-chain", version = "7.0.0-rc0", features = ["async-error"] } -zakura-header-chain = { path = "../zakura-header-chain", version = "2.0.1-rc0" } +zakura-header-chain = { path = "../zakura-header-chain", version = "2.0.2-rc0" } zakura-jsonl-trace = { path = "../zakura-jsonl-trace", version = "1.2.0" } -zakura-node-services = { path = "../zakura-node-services", version = "3.2.3-rc0" } +zakura-node-services = { path = "../zakura-node-services", version = "3.2.4-rc0" } [dev-dependencies] proptest = { workspace = true } diff --git a/crates/zakura-node-services/Cargo.toml b/crates/zakura-node-services/Cargo.toml index 912034a955..162d322240 100644 --- a/crates/zakura-node-services/Cargo.toml +++ b/crates/zakura-node-services/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "zakura-node-services" -version = "3.2.3-rc0" +version = "3.2.4-rc0" authors.workspace = true description = "The interfaces of some Zakura node services. Internal crate, published to support cargo install zakura" license.workspace = true @@ -35,7 +35,7 @@ rpc-client = [ [dependencies] zakura-chain = { path = "../zakura-chain" , version = "7.0.0-rc0" } -zakura-header-chain = { path = "../zakura-header-chain", version = "2.0.1-rc0" } +zakura-header-chain = { path = "../zakura-header-chain", version = "2.0.2-rc0" } tower = { workspace = true } # Optional dependencies diff --git a/crates/zakura-rpc/Cargo.toml b/crates/zakura-rpc/Cargo.toml index 0caee62147..7fe6b33818 100644 --- a/crates/zakura-rpc/Cargo.toml +++ b/crates/zakura-rpc/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "zakura-rpc" -version = "9.0.1-rc0" +version = "10.0.0-rc0" authors.workspace = true description = "The Zakura node's JSON Remote Procedure Call (JSON-RPC) interface. Internal crate, published to support cargo install zakura" license.workspace = true @@ -111,13 +111,13 @@ zcash_transparent = { workspace = true } zakura-chain = { path = "../zakura-chain", version = "7.0.0-rc0", features = [ "json-conversion", ] } -zakura-consensus = { path = "../zakura-consensus", version = "7.0.2-rc0" } -zakura-network = { path = "../zakura-network", version = "7.1.0-rc0" } -zakura-node-services = { path = "../zakura-node-services", version = "3.2.3-rc0", features = [ +zakura-consensus = { path = "../zakura-consensus", version = "8.0.0-rc0" } +zakura-network = { path = "../zakura-network", version = "7.1.1-rc0" } +zakura-node-services = { path = "../zakura-node-services", version = "3.2.4-rc0", features = [ "rpc-client", ] } -zakura-script = { path = "../zakura-script", version = "3.2.2-rc0" } -zakura-state = { path = "../zakura-state", version = "7.2.0-rc0" } +zakura-script = { path = "../zakura-script", version = "3.2.3-rc0" } +zakura-state = { path = "../zakura-state", version = "8.0.0-rc0" } rustls = { version = "0.23.40", default-features = false, features = ["logging", "ring", "std", "tls12"] } tokio-rustls = { version = "0.26.4", default-features = false, features = ["logging", "ring", "tls12"] } # Only used to read the validity dates of the configured RPC TLS certificates. @@ -139,13 +139,13 @@ tokio = { workspace = true, features = ["full", "tracing", "test-util"] } zakura-chain = { path = "../zakura-chain", version = "7.0.0-rc0", features = [ "proptest-impl", ] } -zakura-consensus = { path = "../zakura-consensus", version = "7.0.2-rc0", features = [ +zakura-consensus = { path = "../zakura-consensus", version = "8.0.0-rc0", features = [ "proptest-impl", ] } -zakura-network = { path = "../zakura-network", version = "7.1.0-rc0", features = [ +zakura-network = { path = "../zakura-network", version = "7.1.1-rc0", features = [ "proptest-impl", ] } -zakura-state = { path = "../zakura-state", version = "7.2.0-rc0", features = [ +zakura-state = { path = "../zakura-state", version = "8.0.0-rc0", features = [ "proptest-impl", ] } diff --git a/crates/zakura-rpc/src/methods.rs b/crates/zakura-rpc/src/methods.rs index ff694fa262..314e9a6d86 100644 --- a/crates/zakura-rpc/src/methods.rs +++ b/crates/zakura-rpc/src/methods.rs @@ -52,7 +52,7 @@ use jsonrpsee_types::{ErrorCode, ErrorObject}; use schemars::JsonSchema; use serde::Deserialize; use tokio::{ - sync::{broadcast, mpsc, watch}, + sync::{broadcast, mpsc, watch, Semaphore}, task::JoinHandle, }; use tower::{Service, ServiceExt}; @@ -1016,6 +1016,9 @@ where /// Handler for the `getblocktemplate` RPC. gbt: GetBlockTemplateHandler, + + /// RPC clones share one proposal slot to bound verification work without proof of work. + proposal_limit: Arc, } /// A type alias for the last event logged by the server. @@ -1113,6 +1116,7 @@ where address_book, last_warn_error_log_rx, gbt, + proposal_limit: Arc::new(Semaphore::new(1)), }; // run the process queue @@ -2552,6 +2556,10 @@ where .as_ref() .and_then(GetBlockTemplateParameters::block_proposal_data) { + let _permit = self.proposal_limit.try_acquire().map_err(|_| { + ErrorObject::owned(0, "block proposal validation is busy", None::<()>) + })?; + return validate_block_proposal( self.gbt.block_verifier_router(), block_proposal_bytes, diff --git a/crates/zakura-rpc/src/methods/tests/snapshot.rs b/crates/zakura-rpc/src/methods/tests/snapshot.rs index 5d05edac44..f24c06dd8d 100644 --- a/crates/zakura-rpc/src/methods/tests/snapshot.rs +++ b/crates/zakura-rpc/src/methods/tests/snapshot.rs @@ -1345,11 +1345,23 @@ pub async fn test_mining_rpcs( ..Default::default() })); - let mock_block_verifier_router_request_handler = async move { - mock_block_verifier_router + let mock_block_verifier_router_request_handler = async { + let response = mock_block_verifier_router .expect_request_that(|req| matches!(req, zakura_consensus::Request::CheckProposal(_))) + .await; + + let busy = rpc_mock_state_verifier + .clone() + .get_block_template(Some(GetBlockTemplateParameters { + mode: GetBlockTemplateRequestMode::Proposal, + data: Some(HexData(BLOCK_MAINNET_1_BYTES.to_vec())), + ..Default::default() + })) .await - .respond(Hash::from([0; 32])); + .expect_err("a pending proposal occupies the shared slot"); + assert_eq!(busy.message(), "block proposal validation is busy"); + + response.respond(Hash::from([0; 32])); }; let (get_block_template, ..) = tokio::join!( @@ -1362,6 +1374,15 @@ pub async fn test_mining_rpcs( snapshot_rpc_getblocktemplate("proposal", get_block_template, None, &settings); + rpc_mock_state_verifier + .get_block_template(Some(GetBlockTemplateParameters { + mode: GetBlockTemplateRequestMode::Proposal, + data: Some(HexData(Vec::new())), + ..Default::default() + })) + .await + .expect("completed proposal releases the slot for the next request"); + // These RPC snapshots use the populated state // `submitblock` diff --git a/crates/zakura-script/Cargo.toml b/crates/zakura-script/Cargo.toml index 4f4dc706f9..6d9c35982b 100644 --- a/crates/zakura-script/Cargo.toml +++ b/crates/zakura-script/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "zakura-script" -version = "3.2.2-rc0" +version = "3.2.3-rc0" authors.workspace = true description = "Zakura script verification wrapping zcashd's zcash_script library. Internal crate, published to support cargo install zakura" license.workspace = true diff --git a/crates/zakura-state/Cargo.toml b/crates/zakura-state/Cargo.toml index ea86bf4db8..4bc041cf7c 100644 --- a/crates/zakura-state/Cargo.toml +++ b/crates/zakura-state/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "zakura-state" -version = "7.2.0-rc0" +version = "8.0.0-rc0" authors.workspace = true description = "State contextual verification and storage code for the Zakura node. Internal crate, published to support cargo install zakura" license.workspace = true @@ -74,9 +74,9 @@ tracing = { workspace = true } sapling-crypto = { workspace = true } zakura-assets = { workspace = true } -zakura-node-services = { path = "../zakura-node-services", version = "3.2.3-rc0" } +zakura-node-services = { path = "../zakura-node-services", version = "3.2.4-rc0" } zakura-chain = { path = "../zakura-chain", version = "7.0.0-rc0", features = ["async-error"] } -zakura-header-chain = { path = "../zakura-header-chain", version = "2.0.1-rc0" } +zakura-header-chain = { path = "../zakura-header-chain", version = "2.0.2-rc0" } # prod feature progress-bar howudoin = { workspace = true, optional = true } @@ -88,7 +88,7 @@ derive-getters.workspace = true derive-new.workspace = true [dev-dependencies] -zakura-header-chain = { path = "../zakura-header-chain", version = "2.0.1-rc0", features = ["test-support"] } +zakura-header-chain = { path = "../zakura-header-chain", version = "2.0.2-rc0", features = ["test-support"] } color-eyre = { workspace = true } once_cell = { workspace = true } diff --git a/crates/zakura-state/src/constants.rs b/crates/zakura-state/src/constants.rs index ab5b73eae2..687f255b82 100644 --- a/crates/zakura-state/src/constants.rs +++ b/crates/zakura-state/src/constants.rs @@ -168,3 +168,12 @@ lazy_static! { ) .expect("regex is valid"); } + +/// Maximum outpoints in one batched UTXO request. +/// +/// This scheduling bound does not change transaction or block validity limits. +pub const MAX_UTXO_BATCH_SIZE: usize = 64; + +/// Maximum batched UTXO read tasks per state instance, shared across service clones. +/// Missing-output waits do not consume this budget. +pub const MAX_CONCURRENT_UTXO_BATCH_READS: usize = 4; diff --git a/crates/zakura-state/src/request.rs b/crates/zakura-state/src/request.rs index 9f699e4a65..0208aae112 100644 --- a/crates/zakura-state/src/request.rs +++ b/crates/zakura-state/src/request.rs @@ -1286,6 +1286,14 @@ pub enum Request { /// Outdated requests are pruned on a regular basis. AwaitUtxo(transparent::OutPoint), + /// Resolve at most [`crate::constants::MAX_UTXO_BATCH_SIZE`] outpoints. + /// + /// This request has the informational semantics of [`Request::AwaitUtxo`]. + /// The state batches available-output reads under a shared work budget, then + /// waits for missing outputs without holding a read permit. Duplicate outpoints + /// resolve once. Callers must apply a timeout. + AwaitUtxos(Vec), + /// Finds the first hash that's in the peer's `known_blocks` and the local best chain. /// Returns a list of hashes that follow that intersection, from the best chain. /// @@ -1402,6 +1410,7 @@ impl Request { Request::CommitSemanticallyVerifiedBlock(_) => "commit_semantically_verified_block", Request::CommitCheckpointVerifiedBlock(_) => "commit_checkpoint_verified_block", Request::AwaitUtxo(_) => "await_utxo", + Request::AwaitUtxos(_) => "await_utxos", Request::Depth(_) => "depth", Request::Tip => "tip", Request::BlockLocator => "block_locator", @@ -2049,9 +2058,11 @@ impl TryFrom for ReadRequest { | Request::InvalidateBlock(_) | Request::ReconsiderBlock(_) => Err("ReadService does not write blocks"), - Request::AwaitUtxo(_) => Err("ReadService does not track pending UTXOs. \ + Request::AwaitUtxo(_) | Request::AwaitUtxos(_) => { + Err("ReadService does not track pending UTXOs. \ Manually convert the request to ReadRequest::AnyChainUtxo, \ - and handle pending UTXOs"), + and handle pending UTXOs") + } Request::KnownBlock(_) => Err("ReadService does not track queued blocks"), diff --git a/crates/zakura-state/src/response.rs b/crates/zakura-state/src/response.rs index 4293db3d48..4c6adc639b 100644 --- a/crates/zakura-state/src/response.rs +++ b/crates/zakura-state/src/response.rs @@ -1,7 +1,7 @@ //! State [`tower::Service`] response types. use std::{ - collections::{BTreeMap, HashSet}, + collections::{BTreeMap, HashMap, HashSet}, sync::Arc, }; @@ -103,6 +103,9 @@ pub enum Response { /// pending unverified blocks, or blocks received after the request was sent. Utxo(transparent::Utxo), + /// All unique outputs requested by [`Request::AwaitUtxos`]. + Utxos(HashMap), + /// The response to a `FindBlockHashes` request. BlockHashes(Vec), diff --git a/crates/zakura-state/src/service.rs b/crates/zakura-state/src/service.rs index 862677446e..fd7bad5aca 100644 --- a/crates/zakura-state/src/service.rs +++ b/crates/zakura-state/src/service.rs @@ -81,6 +81,7 @@ mod pending_utxos; mod queued_blocks; pub(crate) mod read; mod traits; +mod utxo_batch; mod write; #[cfg(any(test, feature = "proptest-impl"))] @@ -238,6 +239,8 @@ pub(crate) struct StateService { /// It allows other async tasks to make progress while concurrently reading data from disk. #[derive(Clone, Debug)] pub struct ReadStateService { + /// Shared across clones; permits live until blocking UTXO reads finish. + utxo_read_budget: utxo_batch::UtxoReadBudget, // Configuration // /// The configured Zcash network. @@ -1328,6 +1331,7 @@ impl ReadStateService { finalized_state::embedded_historical_subtrees(&finalized_state.network()).map(Arc::new); let read_service = Self { + utxo_read_budget: utxo_batch::UtxoReadBudget::default(), network: finalized_state.network(), db: finalized_state.db.clone(), non_finalized_state_receiver, @@ -1663,6 +1667,8 @@ impl Service for StateService { // Uses pending_utxos and non_finalized_state_queued_blocks in the StateService. // If the UTXO isn't in the queued blocks, runs concurrently using the ReadStateService. + Request::AwaitUtxos(outpoints) => self.await_utxos(outpoints), + Request::AwaitUtxo(outpoint) => { let timer = CodeTimer::start(); // Prepare the AwaitUtxo future from PendingUxtos. diff --git a/crates/zakura-state/src/service/finalized_state/disk_db.rs b/crates/zakura-state/src/service/finalized_state/disk_db.rs index 8709e19a54..3beaad6761 100644 --- a/crates/zakura-state/src/service/finalized_state/disk_db.rs +++ b/crates/zakura-state/src/service/finalized_state/disk_db.rs @@ -82,6 +82,34 @@ pub type DBThreadMode = rocksdb::SingleThreaded; /// Also the [`rocksdb::DBAccess`] used by database iterators. pub type DB = rocksdb::DBWithThreadMode; +/// Holds one database generation across dependent batches of typed reads. +pub(super) struct DiskReadSnapshot<'a> { + db: &'a DB, + snapshot: rocksdb::Snapshot<'a>, +} + +impl DiskReadSnapshot<'_> { + /// Reads one column family with RocksDB's pinned, batched MultiGet. + pub(super) fn zs_multi_get( + &self, + cf: &impl rocksdb::AsColumnFamilyRef, + keys: &[K], + ) -> Vec> { + let keys: Vec<_> = keys.iter().map(IntoDisk::as_bytes).collect(); + let mut options = ReadOptions::default(); + options.set_snapshot(&self.snapshot); + self.db + .batched_multi_get_cf_opt(cf, &keys, false, &options) + .into_iter() + .map(|value| { + value + .expect("unexpected database failure") + .map(V::from_bytes) + }) + .collect() + } +} + /// Failure while visiting a raw column family without collecting its rows. pub(crate) enum RawVisitError { /// RocksDB failed while advancing the iterator. @@ -1273,6 +1301,14 @@ impl DiskDb { self.db.cf_handle(cf_name) } + /// Creates a read snapshot without exposing the raw database handle. + pub(super) fn read_snapshot(&self) -> DiskReadSnapshot<'_> { + DiskReadSnapshot { + db: &self.db, + snapshot: self.db.snapshot(), + } + } + /// Read raw bytes from one column family without panicking on RocksDB failure. pub(crate) fn raw_get_cf( &self, diff --git a/crates/zakura-state/src/service/finalized_state/disk_db/tests.rs b/crates/zakura-state/src/service/finalized_state/disk_db/tests.rs index bec028db68..7a07a05b94 100644 --- a/crates/zakura-state/src/service/finalized_state/disk_db/tests.rs +++ b/crates/zakura-state/src/service/finalized_state/disk_db/tests.rs @@ -45,6 +45,46 @@ fn format_bytes_preserves_decimal_unit_boundaries() { assert_eq!(format_bytes(u64::MAX), "18.4 EB"); } +#[test] +fn batched_reads_keep_one_snapshot_and_input_order() { + use crate::{FromDisk, IntoDisk}; + use zakura_chain::block::Height; + + let _init_guard = zakura_test::init(); + let db = DiskDb::new( + &Config::ephemeral(), + "batch-snapshot-test", + &Version::new(1, 0, 0), + &Network::Mainnet, + ["values".to_owned()], + false, + ) + .unwrap(); + let cf = db.cf_handle("values").unwrap(); + db.put_cf(cf, Height(1).as_bytes(), Height(10).as_bytes()) + .unwrap(); + db.put_cf(cf, Height(2).as_bytes(), Height(20).as_bytes()) + .unwrap(); + let snapshot = db.read_snapshot(); + db.put_cf(cf, Height(1).as_bytes(), Height(100).as_bytes()) + .unwrap(); + db.delete_cf(cf, Height(2).as_bytes()).unwrap(); + + let keys = [Height(2), Height(1), Height(2), Height(3)]; + assert_eq!( + snapshot.zs_multi_get::<_, Height>(&cf, &keys), + vec![Some(Height(20)), Some(Height(10)), Some(Height(20)), None] + ); + assert_eq!( + db.read_snapshot().zs_multi_get::<_, Height>(&cf, &keys), + vec![None, Some(Height(100)), None, None] + ); + assert_eq!( + Height::from_bytes(db.get_cf(cf, Height(1).as_bytes()).unwrap().unwrap()), + Height(100) + ); +} + #[test] fn exporting_metrics_refreshes_cached_disk_size() { let _init_guard = zakura_test::init(); diff --git a/crates/zakura-state/src/service/finalized_state/zakura_db/transparent.rs b/crates/zakura-state/src/service/finalized_state/zakura_db/transparent.rs index 4e09d925d4..56ffbb28f1 100644 --- a/crates/zakura-state/src/service/finalized_state/zakura_db/transparent.rs +++ b/crates/zakura-state/src/service/finalized_state/zakura_db/transparent.rs @@ -46,6 +46,9 @@ use crate::{ use super::super::TypedColumnFamily; +#[cfg(test)] +mod tests; + /// The name of the transaction hash by spent outpoints column family. pub const TX_LOC_BY_SPENT_OUT_LOC: &str = "tx_loc_by_spent_out_loc"; @@ -157,6 +160,68 @@ impl ZakuraDb { self.utxo_by_location(output_location) } + /// Reads a bounded UTXO batch from one finalized database snapshot. + /// + /// The first MultiGet resolves transaction locations. The second reads outputs. + /// Both must use the same snapshot so a rollback cannot reuse a location between reads. + pub(crate) fn utxos( + &self, + outpoints: &[transparent::OutPoint], + ) -> HashMap { + if outpoints.is_empty() { + return HashMap::new(); + } + let snapshot = self.db.read_snapshot(); + let tx_loc_by_hash = self + .db + .cf_handle("tx_loc_by_hash") + .expect("the transaction location column family exists"); + let utxo_by_out_loc = self + .db + .cf_handle("utxo_by_out_loc") + .expect("the UTXO column family exists"); + let hashes: Vec<_> = outpoints + .iter() + .map(|outpoint| outpoint.hash) + .collect::>() + .into_iter() + .collect(); + let locations: HashMap<_, TransactionLocation> = hashes + .iter() + .copied() + .zip(snapshot.zs_multi_get(&tx_loc_by_hash, &hashes)) + .filter_map(|(hash, location)| location.map(|location| (hash, location))) + .collect(); + let located: Vec<_> = outpoints + .iter() + .filter_map(|outpoint| { + locations.get(&outpoint.hash).map(|location| { + ( + *outpoint, + OutputLocation::from_outpoint(*location, outpoint), + ) + }) + }) + .collect(); + let keys: Vec<_> = located.iter().map(|(_, location)| *location).collect(); + located + .into_iter() + .zip(snapshot.zs_multi_get::<_, transparent::Output>(&utxo_by_out_loc, &keys)) + .filter_map(|((outpoint, location), output)| { + output.map(|output| { + ( + outpoint, + transparent::Utxo::from_location( + output, + location.height(), + location.transaction_index().as_usize(), + ), + ) + }) + }) + .collect() + } + /// Returns the [`TransactionLocation`] of the transaction that spent the given /// [`transparent::OutPoint`], if it is unspent in the finalized state and its /// spending transaction hash has been indexed. diff --git a/crates/zakura-state/src/service/finalized_state/zakura_db/transparent/tests.rs b/crates/zakura-state/src/service/finalized_state/zakura_db/transparent/tests.rs new file mode 100644 index 0000000000..bb96b14a29 --- /dev/null +++ b/crates/zakura-state/src/service/finalized_state/zakura_db/transparent/tests.rs @@ -0,0 +1,97 @@ +//! Diagnostic comparison of UTXO request scheduling against an ephemeral database. + +use futures::StreamExt; +use tower::{buffer::Buffer, ServiceExt}; + +use super::*; +use crate::{service::StateService, Config, Request, Response}; + +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] +#[ignore = "manual warm-database timing comparison; run with --ignored --nocapture"] +#[allow(clippy::print_stdout)] +async fn utxo_state_path_timing() { + let _init_guard = zakura_test::init(); + let (state, _, _, _) = + StateService::new(Config::ephemeral(), &Network::Mainnet, Height::MAX, 0) + .await + .unwrap(); + let mut outpoints = Vec::new(); + let db = &state.read_service.db.db; + let locations = db.cf_handle("tx_loc_by_hash").unwrap(); + let outputs = db.cf_handle("utxo_by_out_loc").unwrap(); + let mut write = DiskWriteBatch::new(); + // Seed only the indexes these informational reads use. This is not a valid chain. + for index in 0u32..1001 { + let mut hash = [0; 32]; + hash[..4].copy_from_slice(&index.to_le_bytes()); + let outpoint = transparent::OutPoint { + hash: hash.into(), + index: 0, + }; + let location = + TransactionLocation::from_usize(Height(1), usize::try_from(index).unwrap() + 1); + write.zs_insert(&locations, outpoint.hash, location); + write.zs_insert( + &outputs, + OutputLocation::from_outpoint(location, &outpoint), + transparent::Output { + value: Amount::::try_from(u64::from(index) + 1).unwrap(), + lock_script: transparent::Script::new(&[0x51]), + }, + ); + outpoints.push(outpoint); + } + db.write(write).unwrap(); + let state = Buffer::new(state, 64); + + for (input_count, iterations) in [(1, 1000), (8, 1000), (64, 100), (1001, 20)] { + for (name, width, concurrency) in [ + ("serial", 1, 1), + ("per-input-64", 1, 64), + ("batched-64", 64, 4), + ] { + // Warm the same keys before each measured mode. + for outpoint in &outpoints[..input_count] { + state + .clone() + .oneshot(Request::AwaitUtxo(*outpoint)) + .await + .unwrap(); + } + let start = std::time::Instant::now(); + for _ in 0..iterations { + let requests = futures::stream::iter(outpoints[..input_count].chunks(width)) + .map(|chunk| { + let request = if width == 1 { + Request::AwaitUtxo(chunk[0]) + } else { + Request::AwaitUtxos(chunk.to_vec()) + }; + state.clone().oneshot(request) + }) + .buffer_unordered(concurrency); + futures::pin_mut!(requests); + let mut count = 0; + while let Some(response) = requests.next().await { + count += match response.unwrap() { + Response::Utxo(utxo) => { + std::hint::black_box(utxo); + 1 + } + Response::Utxos(utxos) => { + let count = utxos.len(); + std::hint::black_box(utxos); + count + } + _ => unreachable!("UTXO request response"), + }; + } + assert_eq!(count, input_count); + } + println!( + "inputs={input_count} mode={name}: {:?}/set", + start.elapsed() / iterations + ); + } + } +} diff --git a/crates/zakura-state/src/service/pending_utxos.rs b/crates/zakura-state/src/service/pending_utxos.rs index 237f98d0cb..035adeb1c4 100644 --- a/crates/zakura-state/src/service/pending_utxos.rs +++ b/crates/zakura-state/src/service/pending_utxos.rs @@ -1,6 +1,10 @@ //! Pending UTXO tracker for [`AwaitUtxo` requests](crate::Request::AwaitUtxo). -use std::{collections::HashMap, future::Future}; +use std::{ + collections::HashMap, + future::Future, + sync::{Arc, Mutex}, +}; use tokio::sync::broadcast; @@ -8,8 +12,10 @@ use zakura_chain::transparent; use crate::{BoxError, Response}; -#[derive(Debug, Default)] -pub struct PendingUtxos(HashMap>); +#[derive(Clone, Debug, Default)] +pub struct PendingUtxos( + Arc>>>, +); impl PendingUtxos { /// Returns a future that will resolve to the `transparent::Output` pointed @@ -20,6 +26,8 @@ impl PendingUtxos { ) -> impl Future> { let mut receiver = self .0 + .lock() + .expect("pending UTXO lock is not poisoned") .entry(outpoint) .or_insert_with(|| { let (sender, _) = broadcast::channel(1); @@ -41,7 +49,12 @@ impl PendingUtxos { /// arrived. #[inline] pub fn respond(&mut self, outpoint: &transparent::OutPoint, utxo: transparent::Utxo) { - if let Some(sender) = self.0.remove(outpoint) { + let sender = self + .0 + .lock() + .expect("pending UTXO lock is not poisoned") + .remove(outpoint); + if let Some(sender) = sender { // Adding the outpoint as a field lets us cross-reference // with the trace of the verification that made the request. tracing::trace!(?outpoint, "found pending UTXO"); @@ -63,11 +76,17 @@ impl PendingUtxos { /// Scan the set of waiting utxo requests for channels where all receivers /// have been dropped and remove the corresponding sender. pub fn prune(&mut self) { - self.0.retain(|_, chan| chan.receiver_count() > 0); + self.0 + .lock() + .expect("pending UTXO lock is not poisoned") + .retain(|_, chan| chan.receiver_count() > 0); } /// Returns the number of utxos that are being waited on. pub fn len(&self) -> usize { - self.0.len() + self.0 + .lock() + .expect("pending UTXO lock is not poisoned") + .len() } } diff --git a/crates/zakura-state/src/service/queued_blocks.rs b/crates/zakura-state/src/service/queued_blocks.rs index 388c13aa33..66dbf63974 100644 --- a/crates/zakura-state/src/service/queued_blocks.rs +++ b/crates/zakura-state/src/service/queued_blocks.rs @@ -41,7 +41,7 @@ pub struct QueuedBlocks { /// Hashes from `queued_blocks`, indexed by block height. by_height: BTreeMap>, /// Known UTXOs. - known_utxos: HashMap, + known_utxos: HashMap>, } impl QueuedBlocks { @@ -64,7 +64,9 @@ impl QueuedBlocks { // Track known UTXOs in queued blocks. for (outpoint, ordered_utxo) in new.0.new_outputs.iter() { self.known_utxos - .insert(*outpoint, ordered_utxo.utxo.clone()); + .entry(*outpoint) + .or_default() + .insert(new_hash, ordered_utxo.utxo.clone()); } self.blocks.insert(new_hash, new); @@ -115,10 +117,8 @@ impl QueuedBlocks { } } - // TODO: only remove UTXOs if there are no queued blocks with that UTXO - // (known_utxos is best-effort, so this is ok for now) for outpoint in queued.0.new_outputs.keys() { - self.known_utxos.remove(outpoint); + remove_utxo_owner(&mut self.known_utxos, outpoint, &queued.0.hash); } } @@ -189,10 +189,8 @@ impl QueuedBlocks { ) .into())); - // TODO: only remove UTXOs if there are no queued blocks with that UTXO - // (known_utxos is best-effort, so this is ok for now) for outpoint in expired_block.new_outputs.keys() { - self.known_utxos.remove(outpoint); + remove_utxo_owner(&mut self.known_utxos, outpoint, &expired_block.hash); } let parent_list = self @@ -247,7 +245,7 @@ impl QueuedBlocks { /// Try to look up this UTXO in any queued block. #[instrument(skip(self))] pub fn utxo(&self, outpoint: &transparent::OutPoint) -> Option { - self.known_utxos.get(outpoint).cloned() + self.known_utxos.get(outpoint)?.values().next().cloned() } /// Clears known_utxos, by_parent, and by_height, then drains blocks. @@ -279,7 +277,7 @@ pub(crate) struct SentHashes { pub sent: HashMap>, /// Known UTXOs. - known_utxos: HashMap, + known_utxos: HashMap>, /// Whether the hashes in this struct can be used check if the chain can be forked. /// This is set to false until all checkpoint-verified block hashes have been pruned. @@ -310,13 +308,18 @@ impl SentHashes { /// Assumes that blocks are added in the order of their height between `finish_batch` calls /// for efficient pruning. pub fn add(&mut self, block: &SemanticallyVerifiedBlock) { + if self.sent.contains_key(&block.hash) { + return; + } // Track known UTXOs in sent blocks. let outpoints = block .new_outputs .iter() .map(|(outpoint, ordered_utxo)| { self.known_utxos - .insert(*outpoint, ordered_utxo.utxo.clone()); + .entry(*outpoint) + .or_default() + .insert(block.hash, ordered_utxo.utxo.clone()); outpoint }) .cloned() @@ -339,13 +342,18 @@ impl SentHashes { /// /// For more details see `add()`. pub fn add_finalized(&mut self, block: &CheckpointVerifiedBlock) { + if self.sent.contains_key(&block.hash) { + return; + } // Track known UTXOs in sent blocks. let outpoints = block .new_outputs .iter() .map(|(outpoint, ordered_utxo)| { self.known_utxos - .insert(*outpoint, ordered_utxo.utxo.clone()); + .entry(*outpoint) + .or_default() + .insert(block.hash, ordered_utxo.utxo.clone()); outpoint }) .cloned() @@ -360,7 +368,7 @@ impl SentHashes { /// Try to look up this UTXO in any sent block. #[instrument(skip(self))] pub fn utxo(&self, outpoint: &transparent::OutPoint) -> Option { - self.known_utxos.get(outpoint).cloned() + self.known_utxos.get(outpoint)?.values().next().cloned() } /// Finishes the current block batch, and stores it for efficient pruning. @@ -387,10 +395,8 @@ impl SentHashes { buf.push_front((hash, height)); return true; } else if let Some(expired_outpoints) = self.sent.remove(&hash) { - // TODO: only remove UTXOs if there are no queued blocks with that UTXO - // (known_utxos is best-effort, so this is ok for now) for outpoint in expired_outpoints.iter() { - self.known_utxos.remove(outpoint); + remove_utxo_owner(&mut self.known_utxos, outpoint, &hash); } } } @@ -422,7 +428,7 @@ impl SentHashes { }; for outpoint in &outpoints { - self.known_utxos.remove(outpoint); + remove_utxo_owner(&mut self.known_utxos, outpoint, hash); } self.curr_buf.retain(|(h, _)| h != hash); @@ -474,3 +480,17 @@ impl SentHashes { metrics::gauge!("state.memory.sent.cache.batch.count").set(batch_iter().count() as f64); } } + +/// Remove only the departing block's output metadata. +fn remove_utxo_owner( + cache: &mut HashMap>, + outpoint: &transparent::OutPoint, + owner: &block::Hash, +) { + if let std::collections::hash_map::Entry::Occupied(mut entry) = cache.entry(*outpoint) { + entry.get_mut().remove(owner); + if entry.get().is_empty() { + entry.remove(); + } + } +} diff --git a/crates/zakura-state/src/service/queued_blocks/tests/vectors.rs b/crates/zakura-state/src/service/queued_blocks/tests/vectors.rs index 67896281f2..89e1b44f31 100644 --- a/crates/zakura-state/src/service/queued_blocks/tests/vectors.rs +++ b/crates/zakura-state/src/service/queued_blocks/tests/vectors.rs @@ -279,3 +279,55 @@ fn dequeue_descendants_removes_the_complete_failed_subtree() -> Result<()> { Ok(()) } + +#[test] +fn fork_outputs_survive_queue_and_sent_removal() -> Result<()> { + let _init_guard = zakura_test::init(); + let block: Arc = + zakura_test::vectors::BLOCK_MAINNET_419200_BYTES.zcash_deserialize_into()?; + let first = block.prepare(); + let mut second = first.clone(); + second.hash = [42; 32].into(); + second.height = (first.height + 1).unwrap(); + Arc::make_mut(&mut Arc::make_mut(&mut second.block).header).previous_block_hash = + [43; 32].into(); + let outpoint = *first.new_outputs.keys().next().unwrap(); + for ordered in second.new_outputs.values_mut() { + ordered.utxo.height = second.height; + } + for prune in [false, true] { + let mut queue = QueuedBlocks::default(); + queue.queue((first.clone(), oneshot::channel().0)); + queue.queue((second.clone(), oneshot::channel().0)); + if prune { + queue.prune_by_height(first.height); + } else { + queue.dequeue_children(first.block.header.previous_block_hash); + } + assert_eq!( + queue.utxo(&outpoint), + Some(second.new_outputs[&outpoint].utxo.clone()) + ); + queue.prune_by_height(second.height); + assert!(queue.utxo(&outpoint).is_none()); + } + + for prune in [false, true] { + let mut sent = SentHashes::default(); + sent.add(&first); + sent.add(&second); + sent.add(&second); + if prune { + sent.prune_by_height(first.height); + } else { + sent.remove(&first.hash); + } + assert_eq!( + sent.utxo(&outpoint), + Some(second.new_outputs[&outpoint].utxo.clone()) + ); + sent.remove(&second.hash); + assert!(sent.utxo(&outpoint).is_none()); + } + Ok(()) +} diff --git a/crates/zakura-state/src/service/tests.rs b/crates/zakura-state/src/service/tests.rs index 222c2117fe..0fefb2c140 100644 --- a/crates/zakura-state/src/service/tests.rs +++ b/crates/zakura-state/src/service/tests.rs @@ -582,6 +582,7 @@ async fn test_populated_state_responds_correctly( // Spec: transactions in the genesis block are ignored. if height.0 != 0 { + let mut batch_utxos = std::collections::HashMap::new(); for transaction in &block.transactions { let transaction_hash = transaction.hash(); @@ -595,9 +596,26 @@ async fn test_populated_state_responds_correctly( from_coinbase, }; + batch_utxos.insert(outpoint, utxo.clone()); transcript.push((Request::AwaitUtxo(outpoint), Ok(Response::Utxo(utxo)))); } } + for outpoints in batch_utxos + .keys() + .copied() + .collect::>() + .chunks(crate::constants::MAX_UTXO_BATCH_SIZE) + { + transcript.push(( + Request::AwaitUtxos(outpoints.to_vec()), + Ok(Response::Utxos( + outpoints + .iter() + .map(|outpoint| (*outpoint, batch_utxos[outpoint].clone())) + .collect(), + )), + )); + } } let mut append_locator_transcript = |split_ind| { diff --git a/crates/zakura-state/src/service/utxo_batch.rs b/crates/zakura-state/src/service/utxo_batch.rs new file mode 100644 index 0000000000..006790581e --- /dev/null +++ b/crates/zakura-state/src/service/utxo_batch.rs @@ -0,0 +1,535 @@ +//! Batched informational UTXO reads and missing-output dependencies. + +use std::{collections::HashMap, sync::Arc}; + +use futures::{future::BoxFuture, stream::FuturesUnordered, FutureExt, StreamExt}; +use tokio::sync::Semaphore; +use tracing::Span; + +use zakura_chain::{diagnostic::CodeTimer, transparent}; + +use crate::{ + constants::{MAX_CONCURRENT_UTXO_BATCH_READS, MAX_UTXO_BATCH_SIZE}, + request::TimedSpan, + BoxError, Response, +}; + +use super::StateService; + +/// Owns the blocking-read lifetime, including after the request is cancelled. +#[derive(Clone, Debug)] +pub(super) struct UtxoReadBudget(Arc); + +impl Default for UtxoReadBudget { + fn default() -> Self { + Self(Arc::new(Semaphore::new(MAX_CONCURRENT_UTXO_BATCH_READS))) + } +} + +impl UtxoReadBudget { + async fn run( + &self, + read: impl FnOnce() -> Result + Send + 'static, + ) -> Result { + let permit = self + .0 + .clone() + .acquire_owned() + .await + .expect("the UTXO read semaphore is never closed"); + let timed_span = TimedSpan::new(CodeTimer::start_desc("utxo_batch"), Span::current()); + timed_span + .spawn_blocking(move || { + // Dropping the request cannot release a permit held by a running read. + let _permit = permit; + read() + }) + .await + } +} + +impl StateService { + pub(super) fn await_utxos( + &mut self, + outpoints: Vec, + ) -> BoxFuture<'static, Result> { + if outpoints.len() > MAX_UTXO_BATCH_SIZE { + return async { Err("UTXO batch exceeds MAX_UTXO_BATCH_SIZE".into()) }.boxed(); + } + + // Subscribe before reading either queued blocks or committed state. An + // output arriving during the read must still satisfy its dependency. + let mut pending = HashMap::new(); + for outpoint in outpoints { + pending + .entry(outpoint) + .or_insert_with(|| self.pending_utxos.queue(outpoint)); + } + let mut available = HashMap::new(); + pending.retain(|outpoint, _| { + let utxo = self + .non_finalized_state_queued_blocks + .utxo(outpoint) + .or_else(|| self.non_finalized_block_write_sent_hashes.utxo(outpoint)); + if let Some(utxo) = utxo { + self.pending_utxos.respond(outpoint, utxo.clone()); + available.insert(*outpoint, utxo); + false + } else { + true + } + }); + + let state = self.read_service.clone(); + let mut changes = state.non_finalized_state_receiver.clone(); + changes.mark_as_seen(); + let mut notifications = self.pending_utxos.clone(); + async move { + let mut missing: std::collections::HashSet<_> = pending.keys().copied().collect(); + let mut waits = FuturesUnordered::new(); + for (outpoint, response) in pending { + waits.push(async move { (outpoint, response.await) }); + } + while !missing.is_empty() { + let outpoints: Vec<_> = missing.iter().copied().collect(); + let read_state = state.clone(); + let budget = state.utxo_read_budget.clone(); + let read = budget.run(move || { + let chains = read_state.latest_non_finalized_state(); + let mut found = HashMap::new(); + let mut on_disk = Vec::new(); + for outpoint in outpoints { + if let Some(utxo) = chains.any_utxo(&outpoint) { + found.insert(outpoint, utxo); + } else { + on_disk.push(outpoint); + } + } + found.extend(read_state.db.utxos(&on_disk)); + Ok(found) + }); + tokio::pin!(read); + let mut read_finished = false; + loop { + tokio::select! { + biased; + Some((outpoint, response)) = waits.next() => { + let Response::Utxo(utxo) = response? else { + unreachable!("pending UTXOs return Utxo responses") + }; + missing.remove(&outpoint); + available.insert(outpoint, utxo); + if missing.is_empty() { + return Ok(Response::Utxos(available)); + } + } + found = &mut read, if !read_finished => { + read_finished = true; + for (outpoint, utxo) in found? { + notifications.respond(&outpoint, utxo); + } + } + changed = changes.changed(), if read_finished => { + changed?; + break; + } + } + } + } + Ok(Response::Utxos(available)) + } + .boxed() + } +} + +#[cfg(test)] +mod tests { + use std::time::Duration; + + use zakura_chain::{ + block::{Block, Height}, + parameters::Network, + serialization::ZcashDeserializeInto, + }; + + use super::*; + use crate::Config; + + #[tokio::test] + async fn read_budget_survives_request_cancellation_and_is_shared() { + let budget = UtxoReadBudget(Arc::new(Semaphore::new(1))); + let shared_budget = budget.clone(); + let (started, start) = tokio::sync::oneshot::channel(); + let (release, released) = std::sync::mpsc::channel(); + let mut first = Box::pin(budget.run(move || { + started.send(()).unwrap(); + released.recv_timeout(Duration::from_secs(30)).unwrap(); + Ok(()) + })); + assert!(futures::poll!(&mut first).is_pending()); + tokio::time::timeout(Duration::from_secs(30), start) + .await + .unwrap() + .unwrap(); + drop(first); + assert_eq!(budget.0.available_permits(), 0); + + let mut second = Box::pin(shared_budget.run(|| Ok(()))); + assert!(futures::poll!(&mut second).is_pending()); + release.send(()).unwrap(); + tokio::time::timeout(Duration::from_secs(30), second) + .await + .unwrap() + .unwrap(); + assert_eq!(budget.0.available_permits(), 1); + } + + #[tokio::test] + async fn notifications_during_read_admission_are_not_lost() { + let _init_guard = zakura_test::init(); + let (mut state, _, _, _) = + StateService::new(Config::ephemeral(), &Network::Mainnet, Height::MAX, 0) + .await + .unwrap(); + let block: Block = zakura_test::vectors::BLOCK_MAINNET_1_BYTES + .zcash_deserialize_into() + .unwrap(); + let expected: HashMap<_, _> = block.transactions[0] + .outputs() + .iter() + .enumerate() + .map(|(index, output)| { + ( + transparent::OutPoint { + hash: block.transactions[0].hash(), + index: index.try_into().unwrap(), + }, + transparent::Utxo::new(output.clone(), Height(1), true), + ) + }) + .collect(); + let budget = state.read_service.utxo_read_budget.clone(); + let permits = budget + .0 + .clone() + .acquire_many_owned(MAX_CONCURRENT_UTXO_BATCH_READS.try_into().unwrap()) + .await + .unwrap(); + let mut lookup = state.await_utxos(expected.keys().copied().collect()); + assert!(futures::poll!(&mut lookup).is_pending()); + for (outpoint, utxo) in &expected { + state.pending_utxos.respond(outpoint, utxo.clone()); + } + assert_eq!( + tokio::time::timeout(Duration::from_secs(30), lookup) + .await + .unwrap() + .unwrap(), + Response::Utxos(expected) + ); + drop(permits); + } + + #[tokio::test] + async fn missing_outputs_release_reads_and_cancel_dependencies() { + let _init_guard = zakura_test::init(); + let (mut state, _, _, _) = + StateService::new(Config::ephemeral(), &Network::Mainnet, Height::MAX, 0) + .await + .unwrap(); + let outpoint = transparent::OutPoint { + hash: [42; 32].into(), + index: 0, + }; + let budget = state.read_service.utxo_read_budget.clone(); + let mut lookup = state.await_utxos(vec![outpoint, outpoint]); + assert_eq!( + state.pending_utxos.len(), + 1, + "duplicate dependencies must share one subscription" + ); + assert!(futures::poll!(&mut lookup).is_pending()); + tokio::time::timeout(Duration::from_secs(30), async { + loop { + assert!(futures::poll!(&mut lookup).is_pending()); + if budget.0.available_permits() == MAX_CONCURRENT_UTXO_BATCH_READS { + break; + } + tokio::task::yield_now().await; + } + }) + .await + .unwrap(); + // An independent read can complete while this batch still waits for an output. + tokio::time::timeout(Duration::from_secs(30), budget.run(|| Ok(()))) + .await + .unwrap() + .unwrap(); + drop(lookup); + state.pending_utxos.prune(); + assert_eq!(state.pending_utxos.len(), 0); + + assert!(state + .await_utxos(vec![outpoint; MAX_UTXO_BATCH_SIZE + 1]) + .await + .is_err()); + assert_eq!( + state.pending_utxos.len(), + 0, + "reject oversized batches before subscribing" + ); + let _permits = budget + .0 + .clone() + .acquire_many_owned(MAX_CONCURRENT_UTXO_BATCH_READS.try_into().unwrap()) + .await + .unwrap(); + assert_eq!( + state.await_utxos(Vec::new()).await.unwrap(), + Response::Utxos(HashMap::new()) + ); + } + + #[tokio::test] + async fn committed_outputs_wake_waiters_and_batch_hits_broadcast() { + use crate::CheckpointVerifiedBlock; + + let _init_guard = zakura_test::init(); + let (mut state, _, _, _) = + StateService::new(Config::ephemeral(), &Network::Mainnet, Height::MAX, 0) + .await + .unwrap(); + let block: Arc = zakura_test::vectors::BLOCK_MAINNET_1_BYTES + .zcash_deserialize_into() + .unwrap(); + let outpoint = transparent::OutPoint::from_usize(block.transactions[0].hash(), 0); + let expected = + transparent::Utxo::new(block.transactions[0].outputs()[0].clone(), Height(1), true); + // Admission happened before this request subscribed. + state.pending_utxos.respond(&outpoint, expected.clone()); + let mut lookup = state.await_utxos(vec![outpoint]); + let budget = state.read_service.utxo_read_budget.clone(); + assert!(futures::poll!(&mut lookup).is_pending()); + tokio::time::timeout(Duration::from_secs(30), async { + loop { + assert!(futures::poll!(&mut lookup).is_pending()); + if budget.0.available_permits() == MAX_CONCURRENT_UTXO_BATCH_READS { + break; + } + tokio::task::yield_now().await; + } + }) + .await + .unwrap(); + for (_, bytes) in zakura_test::vectors::MAINNET_BLOCKS.range(0..=1) { + let block: Arc = bytes.zcash_deserialize_into().unwrap(); + state + .queue_and_commit_to_finalized_state(CheckpointVerifiedBlock::from(block)) + .await + .unwrap() + .unwrap(); + } + assert_eq!( + tokio::time::timeout(Duration::from_secs(30), lookup) + .await + .unwrap() + .unwrap(), + Response::Utxos(HashMap::from([(outpoint, expected.clone())])) + ); + assert_eq!(state.pending_utxos.len(), 0); + + let earlier = state.pending_utxos.queue(outpoint); + state.await_utxos(vec![outpoint]).await.unwrap(); + assert_eq!( + tokio::time::timeout(Duration::from_secs(30), earlier) + .await + .unwrap() + .unwrap(), + Response::Utxo(expected) + ); + assert_eq!(state.pending_utxos.len(), 0); + } + + #[tokio::test] + async fn reconsidered_outputs_wake_existing_waiters() { + use crate::{arbitrary::Prepare, CheckpointVerifiedBlock}; + + let _init_guard = zakura_test::init(); + tokio::time::timeout(Duration::from_secs(30), async { + use zakura_chain::{ + parameters::testnet::{ + ConfiguredActivationHeights, ConfiguredCheckpoints, Parameters, + }, + transaction::Transaction, + }; + let genesis: Arc = zakura_test::vectors::BLOCK_MAINNET_GENESIS_BYTES + .zcash_deserialize_into() + .unwrap(); + let network = Parameters::build() + .with_genesis_hash(genesis.hash()) + .unwrap() + .with_checkpoints(ConfiguredCheckpoints::HeightsAndHashes( + [(Height(0), genesis.hash())].into_iter().collect(), + )) + .unwrap() + .with_activation_heights(ConfiguredActivationHeights { + canopy: Some(1), + ..Default::default() + }) + .unwrap() + .with_disable_pow(true) + .clear_funding_streams() + .to_network() + .unwrap(); + let (mut state, _, _, _) = + StateService::new(Config::ephemeral(), &network, Height(0), 0) + .await + .unwrap(); + state + .queue_and_commit_to_finalized_state(CheckpointVerifiedBlock::from(genesis)) + .await + .unwrap() + .unwrap(); + let block: Arc = zakura_test::vectors::BLOCK_MAINNET_1_BYTES + .zcash_deserialize_into() + .unwrap(); + let mut block = (*block).clone(); + let original = &block.transactions[0]; + block.transactions[0] = Arc::new(Transaction::V4 { + inputs: original.inputs().to_vec(), + outputs: original.outputs().to_vec(), + lock_time: zakura_chain::transaction::LockTime::unlocked(), + expiry_height: Height(0), + joinsplit_data: None, + sapling_shielded_data: None, + }); + Arc::make_mut(&mut block.header).merkle_root = block.transactions.iter().collect(); + let block = Arc::new(block); + state + .queue_and_commit_to_non_finalized_state(block.clone().prepare()) + .await + .unwrap() + .unwrap(); + state + .send_invalidate_block(block.hash()) + .await + .unwrap() + .unwrap(); + state + .non_finalized_block_write_sent_hashes + .remove(&block.hash()); + let outpoint = transparent::OutPoint::from_usize(block.transactions[0].hash(), 0); + let mut lookup = state.await_utxos(vec![outpoint]); + assert!(futures::poll!(&mut lookup).is_pending()); + loop { + assert!(futures::poll!(&mut lookup).is_pending()); + if state.read_service.utxo_read_budget.0.available_permits() + == MAX_CONCURRENT_UTXO_BATCH_READS + { + break; + } + tokio::task::yield_now().await; + } + state + .send_reconsider_block(block.hash()) + .await + .unwrap() + .unwrap(); + assert_eq!( + lookup.await.unwrap(), + Response::Utxos(HashMap::from([( + outpoint, + transparent::Utxo::new( + block.transactions[0].outputs()[0].clone(), + Height(1), + true + ), + )])) + ); + assert_eq!(state.pending_utxos.len(), 0); + }) + .await + .unwrap(); + } + + #[tokio::test] + async fn finalized_batches_match_serial_lookups_and_queued_outputs_skip_reads() { + use crate::{arbitrary::Prepare, CheckpointVerifiedBlock}; + + let _init_guard = zakura_test::init(); + let (mut state, _, _, _) = + StateService::new(Config::ephemeral(), &Network::Mainnet, Height::MAX, 0) + .await + .unwrap(); + let mut outpoints = Vec::new(); + for (_, bytes) in zakura_test::vectors::MAINNET_BLOCKS.range(0..=10) { + let block: Arc = bytes.zcash_deserialize_into().unwrap(); + for tx in &block.transactions { + outpoints.extend( + tx.outputs() + .iter() + .enumerate() + .map(|(index, _)| transparent::OutPoint::from_usize(tx.hash(), index)), + ); + } + state + .queue_and_commit_to_finalized_state(CheckpointVerifiedBlock::from(block)) + .await + .unwrap() + .unwrap(); + } + outpoints.extend_from_within(..); + outpoints.push(transparent::OutPoint { + hash: [43; 32].into(), + index: 0, + }); + outpoints.push(transparent::OutPoint { + hash: outpoints[1].hash, + index: 100, + }); + let expected: HashMap<_, _> = outpoints + .iter() + .filter_map(|outpoint| { + state + .read_service + .db + .utxo(outpoint) + .map(|utxo| (*outpoint, utxo.utxo)) + }) + .collect(); + assert!(!expected.is_empty()); + assert_eq!(state.read_service.db.utxos(&outpoints), expected); + + let queued: Arc = zakura_test::vectors::BLOCK_MAINNET_419201_BYTES + .zcash_deserialize_into() + .unwrap(); + let queued = queued.prepare(); + let expected: HashMap<_, _> = queued + .new_outputs + .iter() + .take(MAX_UTXO_BATCH_SIZE) + .map(|(outpoint, utxo)| (*outpoint, utxo.utxo.clone())) + .collect(); + let (send, _receive) = tokio::sync::oneshot::channel(); + state + .non_finalized_state_queued_blocks + .queue((queued, send)); + let budget = state.read_service.utxo_read_budget.clone(); + let _permits = budget + .0 + .clone() + .acquire_many_owned(MAX_CONCURRENT_UTXO_BATCH_READS.try_into().unwrap()) + .await + .unwrap(); + assert_eq!( + tokio::time::timeout( + Duration::from_secs(30), + state.await_utxos(expected.keys().copied().collect()) + ) + .await + .unwrap() + .unwrap(), + Response::Utxos(expected) + ); + } +} diff --git a/crates/zakura-state/src/service/write.rs b/crates/zakura-state/src/service/write.rs index 6c4ecd8127..24c17f2d1d 100644 --- a/crates/zakura-state/src/service/write.rs +++ b/crates/zakura-state/src/service/write.rs @@ -2249,6 +2249,8 @@ impl WriteBlockWorkerTask { }, ) { Ok((finalized, note_commitment_trees)) => { + // Wake informational UTXO reads that ran before this commit became visible. + non_finalized_state_sender.send_modify(|_| {}); // Whether this successful commit consumed header-carried // tree-aux roots to skip the note-commitment frontier rebuild. if next_block_took_vct_path { diff --git a/crates/zakura-utils/Cargo.toml b/crates/zakura-utils/Cargo.toml index 4cedb7ac1c..6a3b9d6dc3 100644 --- a/crates/zakura-utils/Cargo.toml +++ b/crates/zakura-utils/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "zakura-utils" -version = "2.2.3-rc0" +version = "2.2.4-rc0" authors.workspace = true description = "Developer tools for Zakura maintenance and testing. Internal crate, published to support cargo install zakura" license.workspace = true @@ -74,14 +74,14 @@ tracing-error = { workspace = true } tracing-subscriber = { workspace = true } thiserror = { workspace = true, optional = true } -zakura-node-services = { path = "../zakura-node-services", version = "3.2.3-rc0", optional = true } +zakura-node-services = { path = "../zakura-node-services", version = "3.2.4-rc0", optional = true } zakura-chain = { path = "../zakura-chain", version = "7.0.0-rc0", optional = true } # These crates are needed for the zakura-checkpoints binary itertools = { workspace = true, optional = true } # This crate is needed for the zakura-checkpoints offline export mode -zakura-state = { path = "../zakura-state", version = "7.2.0-rc0", optional = true } +zakura-state = { path = "../zakura-state", version = "8.0.0-rc0", optional = true } # This crate is needed for the zakura-checkpoints binary tokio = { workspace = true, features = ["macros", "rt-multi-thread"], optional = true } diff --git a/crates/zakurad/Cargo.toml b/crates/zakurad/Cargo.toml index 41ead7322d..5dc5db7d46 100644 --- a/crates/zakurad/Cargo.toml +++ b/crates/zakurad/Cargo.toml @@ -169,21 +169,21 @@ comparison-interpreter = ["zakura-script/comparison-interpreter"] [dependencies] zakura-chain = { path = "../zakura-chain", version = "7.0.0-rc0" } -zakura-consensus = { path = "../zakura-consensus", version = "7.0.2-rc0" } -zakura-header-chain = { path = "../zakura-header-chain", version = "2.0.1-rc0" } +zakura-consensus = { path = "../zakura-consensus", version = "8.0.0-rc0" } +zakura-header-chain = { path = "../zakura-header-chain", version = "2.0.2-rc0" } zakura-jsonl-trace = { path = "../zakura-jsonl-trace", version = "1.2.0" } -zakura-network = { path = "../zakura-network", version = "7.1.0-rc0" } -zakura-node-services = { path = "../zakura-node-services", version = "3.2.3-rc0", features = ["rpc-client"] } -zakura-rpc = { path = "../zakura-rpc", version = "9.0.1-rc0" } -zakura-state = { path = "../zakura-state", version = "7.2.0-rc0" } +zakura-network = { path = "../zakura-network", version = "7.1.1-rc0" } +zakura-node-services = { path = "../zakura-node-services", version = "3.2.4-rc0", features = ["rpc-client"] } +zakura-rpc = { path = "../zakura-rpc", version = "10.0.0-rc0" } +zakura-state = { path = "../zakura-state", version = "8.0.0-rc0" } # zakura-script is not used directly, but we list it here to enable the # "comparison-interpreter" feature. (Feature unification will take care of # enabling it in the other imports of zcash-script.) -zakura-script = { path = "../zakura-script", version = "3.2.2-rc0" } +zakura-script = { path = "../zakura-script", version = "3.2.3-rc0" } zcash_script = { workspace = true } # Required for crates.io publishing, but it's only used in tests -zakura-utils = { path = "../zakura-utils", version = "2.2.3-rc0", optional = true } +zakura-utils = { path = "../zakura-utils", version = "2.2.4-rc0", optional = true } abscissa_core = { workspace = true, default-features = false, features = ["application"] } clap = { workspace = true, features = ["cargo"] } @@ -306,10 +306,10 @@ proptest = { workspace = true } proptest-derive = { workspace = true } zakura-chain = { path = "../zakura-chain", version = "7.0.0-rc0", features = ["proptest-impl"] } -zakura-consensus = { path = "../zakura-consensus", version = "7.0.2-rc0", features = ["proptest-impl"] } -zakura-header-chain = { path = "../zakura-header-chain", version = "2.0.1-rc0", features = ["test-support"] } -zakura-network = { path = "../zakura-network", version = "7.1.0-rc0", features = ["proptest-impl", "zakura-testkit"] } -zakura-state = { path = "../zakura-state", version = "7.2.0-rc0", features = ["proptest-impl"] } +zakura-consensus = { path = "../zakura-consensus", version = "8.0.0-rc0", features = ["proptest-impl"] } +zakura-header-chain = { path = "../zakura-header-chain", version = "2.0.2-rc0", features = ["test-support"] } +zakura-network = { path = "../zakura-network", version = "7.1.1-rc0", features = ["proptest-impl", "zakura-testkit"] } +zakura-state = { path = "../zakura-state", version = "8.0.0-rc0", features = ["proptest-impl"] } zakura-test = { path = "../zakura-test", version = "2.1.0" } @@ -322,7 +322,7 @@ zakura-test = { path = "../zakura-test", version = "2.1.0" } # When `-Z bindeps` is stabilised, enable this binary dependency instead: # https://github.com/rust-lang/cargo/issues/9096 # zakura-utils { path = "../zakura-utils", artifact = "bin:zakura-checkpoints" } -zakura-utils = { path = "../zakura-utils", version = "2.2.3-rc0" } +zakura-utils = { path = "../zakura-utils", version = "2.2.4-rc0" } [package.metadata.cargo-udeps.ignore] # These dependencies are false positives - they are actually used diff --git a/docs/changelog/params.md b/docs/changelog/params.md index 20c52a47a5..001c44dc16 100644 --- a/docs/changelog/params.md +++ b/docs/changelog/params.md @@ -32,6 +32,12 @@ Keep entries **newest-first**. Each row records: | Parameter | Location | Old → New | PR | Why | | --- | --- | --- | --- | --- | +| `MAX_CACHED_UTXO_SCRIPT_BYTES` | `crates/zakura-consensus/src/transaction/utxo_resolver.rs` | unbounded → `8 MiB` per block resolver | [#919](https://github.com/zakura-core/zakura/pull/919) | Fall back to individual lookups under the shared deadline when historical scripts exceed the cache budget. | +| `proposal_limit` | `crates/zakura-rpc/src/methods.rs` | unbounded → `1` concurrent proposal per RPC server | [#919](https://github.com/zakura-core/zakura/pull/919) | Reject excess proposal requests before verification instead of accumulating work without proof of work. | +| `MAX_UTXO_BATCH_SIZE` | `crates/zakura-state/src/constants.rs` | one outpoint per request → at most `64` per batch | [#919](https://github.com/zakura-core/zakura/pull/919) | Bound request work while using native database MultiGet. | +| `MAX_CONCURRENT_UTXO_BATCH_READS` | `crates/zakura-state/src/constants.rs` | no shared block-lookup read cap → `4` batch-read tasks per state instance | [#919](https://github.com/zakura-core/zakura/pull/919) | Bound running batched reads across blocks; missing-output waits release read permits. | +| `MAX_IN_FLIGHT_UTXO_BATCHES` | `crates/zakura-consensus/src/transaction/utxo_resolver.rs` | per-transaction lookups → at most `4` pending batches per block | [#919](https://github.com/zakura-core/zakura/pull/919) | Bound each block's outstanding state requests, including dependency waits. | +| `UTXO_LOOKUP_TIMEOUT` scope | `crates/zakura-consensus/src/transaction/utxo_resolver.rs` | `6 min` per external input → `6 min` per block resolver | [#919](https://github.com/zakura-core/zakura/pull/919) | Bound read admission and missing-output waits with one deadline; later batches can expire earlier than with serial lookups. Standalone and mempool paths retain their existing timeouts. | | `LEGACY_FALLBACK_APPLY_DRAIN_DEADLINE` | `crates/zakurad/src/commands/start/zakura/coordinator.rs` | new → `30 min` | [#831](https://github.com/zakura-core/zakura/pull/831) | Terminate the node when native block applies prevent legacy fallback from acquiring exclusive ownership, instead of leaving the fallback handoff pending forever. | | `MAX_CANDIDATE_TIPS_V1` | `crates/zakura-header-chain/src/config.rs` | `10` → `11` | [#831](https://github.com/zakura-core/zakura/pull/831) | Retain ten full-state fork tips plus one independent selected header tip, so header candidate pressure cannot evict a branch that full state still owns. | | `VCT_LOCAL_OPERATION_FATAL_AFTER` | `crates/zakura-network/src/zakura/header_sync/reactor.rs` | new → `30 min` | [#821](https://github.com/zakura-core/zakura/pull/821) | Terminate a node whose local VCT repair prepare or apply operation remains pending, while allowing slow valid operations substantially more time than the existing one-minute stall diagnostic. | diff --git a/docs/changelog/unreleased/919.md b/docs/changelog/unreleased/919.md new file mode 100644 index 0000000000..ec1ed06158 --- /dev/null +++ b/docs/changelog/unreleased/919.md @@ -0,0 +1,15 @@ +## Fixed + +- Batched external UTXO lookups during block verification under a shared database-read limit + ([#919](https://github.com/zakura-core/zakura/pull/919)). +- Limited RPC block proposal validation to one concurrent request per server + ([#919](https://github.com/zakura-core/zakura/pull/919)). +- Limited the block resolver's retained locking scripts to 8 MiB, with per-input + lookups for larger dependency sets. Rejected duplicate spends and scripts that + exceed the interpreter's existing size limit before copying historical outputs + ([#919](https://github.com/zakura-core/zakura/pull/919)). +- Woke UTXO waiters on batch read hits and committed state updates. Allowed output + notifications to complete lookups while disk reads wait for admission + ([#919](https://github.com/zakura-core/zakura/pull/919)). +- Preserved cached outputs while another queued or sent fork block still owns them + ([#919](https://github.com/zakura-core/zakura/pull/919)).