diff --git a/Cargo.lock b/Cargo.lock index 50e3b5b41c..9748b76b83 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -4143,6 +4143,7 @@ dependencies = [ "sp-runtime-interface", "sp-trie", "spv_validation", + "tempfile", "testcontainers", "timed-map", "tokio", diff --git a/Cargo.toml b/Cargo.toml index eac8eac5fd..1ce0e81764 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -191,6 +191,7 @@ sp-trie = { version = "6.0", default-features = false } sql-builder = "3.1.1" syn = "1.0" sysinfo = "0.28" +tempfile = "3.4.0" # using the same version as cosmrs tendermint-rpc = { version = "0.35", default-features = false } testcontainers = "0.15.0" diff --git a/mm2src/coins/Cargo.toml b/mm2src/coins/Cargo.toml index 1042df4966..c099395605 100644 --- a/mm2src/coins/Cargo.toml +++ b/mm2src/coins/Cargo.toml @@ -61,6 +61,7 @@ hex.workspace = true http.workspace = true itertools = { workspace = true, features = ["use_std"] } jsonrpc-core.workspace = true +jubjub.workspace = true keys = { path = "../mm2_bitcoin/keys" } lazy_static.workspace = true libc.workspace = true diff --git a/mm2src/coins/eth/eth_tests.rs b/mm2src/coins/eth/eth_tests.rs index ce5d58f58a..7cde52e740 100644 --- a/mm2src/coins/eth/eth_tests.rs +++ b/mm2src/coins/eth/eth_tests.rs @@ -1,7 +1,6 @@ use super::*; use crate::IguanaPrivKey; use common::block_on; -use futures_util::future; use mm2_core::mm_ctx::MmCtxBuilder; cfg_native!( @@ -10,6 +9,7 @@ cfg_native!( use common::{now_sec, block_on_f01}; use ethkey::{Generator, Random}; + use futures_util::future; use mm2_test_helpers::for_tests::{ETH_MAINNET_CHAIN_ID, ETH_MAINNET_NODES, ETH_SEPOLIA_CHAIN_ID, ETH_SEPOLIA_NODES, ETH_SEPOLIA_TOKEN_CONTRACT}; use mocktopus::mocking::*; diff --git a/mm2src/coins/utxo/utxo_tests.rs b/mm2src/coins/utxo/utxo_tests.rs index c53d81045d..7d8d760452 100644 --- a/mm2src/coins/utxo/utxo_tests.rs +++ b/mm2src/coins/utxo/utxo_tests.rs @@ -12,10 +12,10 @@ use crate::rpc_command::init_scan_for_new_addresses::{InitScanAddressesRpcOps, S ScanAddressesResponse}; use crate::utxo::qtum::{qtum_coin_with_priv_key, QtumCoin, QtumDelegationOps, QtumDelegationRequest}; #[cfg(not(target_arch = "wasm32"))] -use crate::utxo::rpc_clients::{BlockHashOrHeight, NativeUnspent}; +use crate::utxo::rpc_clients::{BlockHashOrHeight, ElectrumClientSettings, NativeUnspent}; use crate::utxo::rpc_clients::{ElectrumBalance, ElectrumBlockHeader, ElectrumClient, ElectrumClientImpl, - ElectrumClientSettings, GetAddressInfoRes, ListSinceBlockRes, NativeClient, - NativeClientImpl, NetworkInfo, UtxoRpcClientOps, ValidateAddressRes, VerboseBlock}; + GetAddressInfoRes, ListSinceBlockRes, NativeClient, NativeClientImpl, NetworkInfo, + UtxoRpcClientOps, ValidateAddressRes, VerboseBlock}; use crate::utxo::spv::SimplePaymentVerification; #[cfg(not(target_arch = "wasm32"))] use crate::utxo::utxo_block_header_storage::{BlockHeaderStorage, SqliteBlockHeadersStorage}; @@ -43,6 +43,7 @@ use futures::future::{join_all, Either, FutureExt, TryFutureExt}; use hex::FromHex; use keys::prefixes::*; use mm2_core::mm_ctx::MmCtxBuilder; +#[cfg(not(target_arch = "wasm32"))] use mm2_event_stream::StreamingManager; use mm2_number::bigdecimal::{BigDecimal, Signed}; use mm2_number::MmNumber; diff --git a/mm2src/coins/z_coin.rs b/mm2src/coins/z_coin.rs index e31a8f8fa0..243d6f89e6 100644 --- a/mm2src/coins/z_coin.rs +++ b/mm2src/coins/z_coin.rs @@ -14,18 +14,17 @@ use crate::my_tx_history_v2::{MyTxHistoryErrorV2, MyTxHistoryRequestV2, MyTxHist use crate::rpc_command::init_withdraw::{InitWithdrawCoin, WithdrawInProgressStatus, WithdrawTaskHandleShared}; use crate::utxo::rpc_clients::{ElectrumConnectionSettings, UnspentInfo, UtxoRpcClientEnum, UtxoRpcError, UtxoRpcFut, UtxoRpcResult}; -use crate::utxo::utxo_builder::UtxoCoinBuildError; -use crate::utxo::utxo_builder::{UtxoCoinBuilder, UtxoCoinBuilderCommonOps, UtxoFieldsWithGlobalHDBuilder, - UtxoFieldsWithHardwareWalletBuilder, UtxoFieldsWithIguanaSecretBuilder}; -use crate::utxo::utxo_common::{addresses_from_script, big_decimal_from_sat}; -use crate::utxo::utxo_common::{big_decimal_from_sat_unsigned, payment_script}; +use crate::utxo::utxo_builder::{UtxoCoinBuildError, UtxoCoinBuilder, UtxoCoinBuilderCommonOps, + UtxoFieldsWithGlobalHDBuilder, UtxoFieldsWithHardwareWalletBuilder, + UtxoFieldsWithIguanaSecretBuilder}; +use crate::utxo::utxo_common::{addresses_from_script, big_decimal_from_sat, big_decimal_from_sat_unsigned, + payment_script}; use crate::utxo::{sat_from_big_decimal, utxo_common, ActualFeeRate, AdditionalTxData, AddrFromStrError, Address, BroadcastTxErr, FeePolicy, GetUtxoListOps, HistoryUtxoTx, HistoryUtxoTxMap, MatureUnspentList, - RecentlySpentOutPointsGuard, UtxoActivationParams, UtxoAddressFormat, UtxoArc, UtxoCoinFields, - UtxoCommonOps, UtxoRpcMode, UtxoTxBroadcastOps, UtxoTxGenerationOps, VerboseTransactionFrom}; -use crate::utxo::{UnsupportedAddr, UtxoFeeDetails}; -use crate::z_coin::storage::{BlockDbImpl, WalletDbShared}; - + RecentlySpentOutPointsGuard, UnsupportedAddr, UtxoActivationParams, UtxoAddressFormat, UtxoArc, + UtxoCoinFields, UtxoCommonOps, UtxoFeeDetails, UtxoRpcMode, UtxoTxBroadcastOps, UtxoTxGenerationOps, + VerboseTransactionFrom}; +use crate::z_coin::storage::{BlockDbImpl, LockedNotesStorage, WalletDbShared}; use crate::z_coin::z_tx_history::{fetch_tx_history_from_db, ZCoinTxHistoryItem}; use crate::{BalanceError, BalanceFut, CheckIfMyPaymentSentArgs, CoinBalance, ConfirmPaymentInput, DexFee, FeeApproxStage, FoundSwapTxSpend, HistorySyncState, MarketCoinOps, MmCoin, NegotiateSwapContractAddrErr, @@ -38,15 +37,18 @@ use crate::{BalanceError, BalanceFut, CheckIfMyPaymentSentArgs, CoinBalance, Con ValidatePaymentError, ValidatePaymentInput, VerificationError, VerificationResult, WaitForHTLCTxSpendArgs, WatcherOps, WeakSpawner, WithdrawError, WithdrawFut, WithdrawRequest}; +use crate::z_coin::storage::z_locked_notes::LockedNote; use async_trait::async_trait; use bitcrypto::dhash256; use chain::constants::SEQUENCE_FINAL; use chain::{Transaction as UtxoTx, TransactionOutput}; -use common::executor::{AbortableSystem, AbortedError}; +use common::executor::{AbortableSystem, AbortedError, SpawnFuture}; +use common::log::info; use common::{calc_total_pages, log}; use crypto::privkey::{key_pair_from_secret, secp_privkey_from_hash}; use crypto::HDPathToCoin; use crypto::{Bip32DerPathOps, GlobalHDAccountArc}; +use futures::channel::oneshot; use futures::compat::Future01CompatExt; use futures::lock::Mutex as AsyncMutex; use futures::{FutureExt, TryFutureExt}; @@ -62,11 +64,10 @@ use script::{Builder as ScriptBuilder, Opcode, Script, TransactionInputSigner}; use serde_json::Value as Json; use serialization::CoinVariant; use std::collections::{HashMap, HashSet}; -use std::convert::TryInto; +use std::convert::{TryFrom, TryInto}; use std::iter; use std::num::NonZeroU32; use std::num::TryFromIntError; -use std::path::PathBuf; use std::sync::Arc; pub use z_coin_errors::*; pub use z_htlc::z_send_dex_fee; @@ -79,8 +80,10 @@ use zcash_client_backend::wallet::{AccountId, SpendableNote}; use zcash_extras::WalletRead; use zcash_primitives::consensus::{BlockHeight, BranchId, NetworkUpgrade, Parameters, H0}; use zcash_primitives::memo::MemoBytes; +use zcash_primitives::sapling::keys::prf_expand; use zcash_primitives::sapling::keys::OutgoingViewingKey; use zcash_primitives::sapling::note_encryption::try_sapling_output_recovery; +use zcash_primitives::sapling::Rseed; use zcash_primitives::transaction::builder::Builder as ZTxBuilder; use zcash_primitives::transaction::components::{Amount, OutputDescription, TxOut}; use zcash_primitives::transaction::Transaction as ZTransaction; @@ -91,8 +94,7 @@ use zcash_proofs::prover::LocalTxProver; cfg_native!( use common::{async_blocking, sha256_digest}; - use zcash_client_sqlite::error::SqliteClientError as ZcashClientError; - use zcash_client_sqlite::wallet::get_balance; + use std::path::PathBuf; use zcash_proofs::default_params_folder; use z_rpc::init_native_client; ); @@ -100,7 +102,6 @@ cfg_native!( cfg_wasm32!( use crate::z_coin::storage::ZcashParamsWasmImpl; use common::executor::AbortOnDropHandle; - use futures::channel::oneshot; use rand::rngs::OsRng; use zcash_primitives::transaction::builder::TransactionMetadata; pub use z_coin_errors::ZCoinBalanceError; @@ -207,6 +208,7 @@ pub struct ZCoinFields { light_wallet_db: WalletDbShared, consensus_params: ZcoinConsensusParams, sync_state_connector: AsyncMutex, + locked_notes_db: LockedNotesStorage, } impl Transaction for ZTransaction { @@ -262,6 +264,13 @@ pub struct ZcoinTxDetails { internal_id: i64, } +struct GenTxData<'a> { + tx: ZTransaction, + data: AdditionalTxData, + sync_guard: SaplingSyncGuard<'a>, + rseeds: Vec, +} + impl ZCoin { #[inline] pub fn utxo_rpc_client(&self) -> &UtxoRpcClientEnum { &self.utxo_arc.rpc_client } @@ -325,25 +334,7 @@ impl ZCoin { }) } - #[cfg(not(target_arch = "wasm32"))] - async fn my_balance_sat(&self) -> Result> { - let wallet_db = self.z_fields.light_wallet_db.clone(); - async_blocking(move || { - let db_guard = wallet_db.db.inner(); - let db_guard = db_guard.lock().unwrap(); - let balance = get_balance(&db_guard, AccountId::default())?.into(); - Ok(balance) - }) - .await - } - - #[cfg(target_arch = "wasm32")] - async fn my_balance_sat(&self) -> Result> { - let wallet_db = self.z_fields.light_wallet_db.clone(); - Ok(wallet_db.db.get_balance(AccountId::default()).await?.into()) - } - - async fn get_spendable_notes(&self) -> Result, MmError> { + async fn get_wallet_notes(&self) -> Result, MmError> { let wallet_db = self.z_fields.light_wallet_db.clone(); let db_guard = wallet_db.db; let latest_db_block = match db_guard @@ -362,8 +353,8 @@ impl ZCoin { } /// Returns spendable notes - async fn spendable_notes_ordered(&self) -> Result, MmError> { - let mut unspents = self.get_spendable_notes().await?; + async fn wallet_notes_ordered(&self) -> Result, MmError> { + let mut unspents = self.get_wallet_notes().await?; unspents.sort_unstable_by(|a, b| a.note_value.cmp(&b.note_value)); Ok(unspents) @@ -383,27 +374,28 @@ impl ZCoin { &self, t_outputs: Vec, z_outputs: Vec, - ) -> Result<(ZTransaction, AdditionalTxData, SaplingSyncGuard<'_>), MmError> { + ) -> Result, MmError> { + // Wait for chain to sync before selecting spendable notes or waiting for locked_notes to become + // available. let sync_guard = self.wait_for_gen_tx_blockchain_sync().await?; - + drop(sync_guard); let tx_fee = self.get_one_kbyte_tx_fee().await?; let t_output_sat: u64 = t_outputs.iter().fold(0, |cur, out| cur + u64::from(out.value)); let z_output_sat: u64 = z_outputs.iter().fold(0, |cur, out| cur + u64::from(out.amount)); let total_output_sat = t_output_sat + z_output_sat; let total_output = big_decimal_from_sat_unsigned(total_output_sat, self.utxo_arc.decimals); let total_required = &total_output + &tx_fee; + let spendable_notes = wait_for_spendable_balance_spawner(self, &total_required).await?; + + // Recreate sync_guard + let sync_guard = self.wait_for_gen_tx_blockchain_sync().await?; - let spendable_notes = self - .spendable_notes_ordered() - .await - .mm_err(|err| GenTxError::SpendableNotesError(err.to_string()))?; let mut total_input_amount = BigDecimal::from(0); let mut change = BigDecimal::from(0); - let mut received_by_me = 0u64; - let mut tx_builder = ZTxBuilder::new(self.consensus_params(), sync_guard.respawn_guard.current_block()); + let mut rseeds: Vec = vec![]; for spendable_note in spendable_notes { total_input_amount += big_decimal_from_sat_unsigned(spendable_note.note_value.into(), self.decimals()); @@ -422,6 +414,8 @@ impl ZCoin { .or_mm_err(|| GenTxError::FailedToGetMerklePath)?, )?; + rseeds.push(rseed_to_string(&spendable_note.rseed)); + if total_input_amount >= total_required { change = &total_input_amount - &total_required; break; @@ -444,19 +438,21 @@ impl ZCoin { tx_builder.add_sapling_output(z_out.viewing_key, z_out.to_addr, z_out.amount, z_out.memo)?; } + // add change to tx output + let change_sat = sat_from_big_decimal(&change, self.utxo_arc.decimals)?; if change > BigDecimal::from(0u8) { - let change_sat = sat_from_big_decimal(&change, self.utxo_arc.decimals)?; received_by_me += change_sat; + let change_amount = Amount::from_u64(change_sat).map_to_mm(|_| { + GenTxError::NumConversion(NumConversError(format!( + "Failed to get ZCash amount from {}", + change_sat + ))) + })?; tx_builder.add_sapling_output( Some(self.z_fields.evk.fvk.ovk), self.z_fields.my_z_addr.clone(), - Amount::from_u64(change_sat).map_to_mm(|_| { - GenTxError::NumConversion(NumConversError(format!( - "Failed to get ZCash amount from {}", - change_sat - ))) - })?, + change_amount, None, )?; } @@ -478,13 +474,19 @@ impl ZCoin { .await? .tx_result?; - let additional_data = AdditionalTxData { + let data = AdditionalTxData { received_by_me, spent_by_me: sat_from_big_decimal(&total_input_amount, self.decimals())?, fee_amount: sat_from_big_decimal(&tx_fee, self.decimals())?, kmd_rewards: None, }; - Ok((tx, additional_data, sync_guard)) + + Ok(GenTxData { + tx, + data, + sync_guard, + rseeds, + }) } pub async fn send_outputs( @@ -492,7 +494,12 @@ impl ZCoin { t_outputs: Vec, z_outputs: Vec, ) -> Result> { - let (tx, _, mut sync_guard) = self.gen_tx(t_outputs, z_outputs).await?; + let GenTxData { + tx, + data, + rseeds, + mut sync_guard, + } = self.gen_tx(t_outputs, z_outputs).await?; let mut tx_bytes = Vec::with_capacity(1024); tx.write(&mut tx_bytes).expect("Write should not fail"); @@ -501,6 +508,25 @@ impl ZCoin { .compat() .await?; + // TODO: Execute updates to `locked_notes_db` and `wallet_db` in a single transaction. + // This will be possible with a newer librustzcash that supports both spent notes and unconfirmed change tracking. + // See: https://github.com/KomodoPlatform/komodo-defi-framework/pull/2331#pullrequestreview-2883773336 + for rseed in rseeds { + self.z_fields + .locked_notes_db + .insert_spent_note(tx.txid().to_string(), rseed) + .await + .mm_err(|err| SendOutputsErr::InternalError(err.to_string()))?; + } + + if data.received_by_me > 0 { + self.z_fields + .locked_notes_db + .insert_change_note(tx.txid().to_string(), data.received_by_me) + .await + .mm_err(|err| SendOutputsErr::InternalError(err.to_string()))?; + } + sync_guard.respawn_guard.watch_for_tx(tx.txid()); Ok(tx) } @@ -678,7 +704,7 @@ impl ZCoin { else { return Ok(false); }; - if &address == expected_address { + if &address != expected_address { return Ok(false); } if note.value != amount_sat { @@ -918,10 +944,13 @@ impl<'a> UtxoCoinBuilder for ZCoinBuilder<'a> { let z_tx_prover = self.z_tx_prover().await?; let blocks_db = self.init_blocks_db().await?; + let locked_notes_db = LockedNotesStorage::new(self.ctx, self.my_z_addr_encoded.clone()).await?; let (sync_state_connector, light_wallet_db) = match &self.z_coin_params.mode { #[cfg(not(target_arch = "wasm32"))] - ZcoinRpcMode::Native => init_native_client(&self, self.native_client()?, blocks_db).await?, + ZcoinRpcMode::Native => { + init_native_client(&self, self.native_client()?, blocks_db, locked_notes_db.clone()).await? + }, ZcoinRpcMode::Light { light_wallet_d_servers, sync_params, @@ -934,6 +963,7 @@ impl<'a> UtxoCoinBuilder for ZCoinBuilder<'a> { blocks_db, sync_params, skip_sync_params.unwrap_or_default(), + locked_notes_db.clone(), ) .await? }, @@ -950,6 +980,7 @@ impl<'a> UtxoCoinBuilder for ZCoinBuilder<'a> { light_wallet_db, consensus_params: self.protocol_info.consensus_params, sync_state_connector, + locked_notes_db, }); Ok(ZCoin { utxo_arc, z_fields }) @@ -1030,12 +1061,7 @@ impl<'a> ZCoinBuilder<'a> { let ctx = &self.ctx; let ticker = self.ticker.to_string(); - #[cfg(target_arch = "wasm32")] - let cache_db_path = PathBuf::new(); - #[cfg(not(target_arch = "wasm32"))] - let cache_db_path = self.ctx.global_dir().join(format!("{}_cache.db", self.ticker)); - - BlockDbImpl::new(ctx, ticker, cache_db_path) + BlockDbImpl::new(ctx, ticker) .await .mm_err(|err| ZcoinClientInitError::ZcoinStorageError(err.to_string())) } @@ -1150,15 +1176,62 @@ impl MarketCoinOps for ZCoin { )) } + /// Calculates the wallet balance, divided into spendable and unspendable portions. + /// Unspendable balance consists of notes that are locked in the wallet. + /// TODO: Track unconfirmed change outputs in a dedicated DB/table (similar to locked_notes_db). + /// - Include them in the unspendable portion of the balance until confirmed. + /// - This will improve spendable/unspendable accuracy. fn my_balance(&self) -> BalanceFut { let coin = self.clone(); let fut = async move { - let sat = coin - .my_balance_sat() + let locked_notes = coin + .z_fields + .locked_notes_db + .load_all_notes() .await .mm_err(|e| BalanceError::WalletStorageError(e.to_string()))?; - Ok(CoinBalance::new(big_decimal_from_sat_unsigned(sat, coin.decimals()))) + + // Locked (unconfirmed) spent notes are not counted as spendable. + let spent_rseeds: HashSet<_> = locked_notes + .iter() + .filter_map(|n| { + if let LockedNote::Spent { rseed, .. } = n { + Some(rseed.clone()) + } else { + None + } + }) + .collect(); + + // Locked (unconfirmed) change notes are counted as unspendable. + let unspendable_change_sat: u64 = locked_notes + .iter() + .filter_map(|n| { + if let LockedNote::Change { value, .. } = n { + Some(*value) + } else { + None + } + }) + .sum(); + + let wallet_notes = coin + .get_wallet_notes() + .await + .map_err(|err| BalanceError::WalletStorageError(err.to_string()))?; + + let spendable_amount = wallet_notes + .iter() + .filter(|n| !spent_rseeds.contains(&rseed_to_string(&n.rseed))) + .fold(Amount::zero(), |acc, n| acc + n.note_value); + + let spendable_sat = + u64::try_from(spendable_amount).map_to_mm(|err| BalanceError::Internal(err.to_string()))?; + let unspendable = big_decimal_from_sat_unsigned(unspendable_change_sat, coin.decimals()); + let spendable = big_decimal_from_sat_unsigned(spendable_sat, coin.decimals()); + Ok(CoinBalance { spendable, unspendable }) }; + Box::new(fut.boxed().compat()) } @@ -1409,7 +1482,7 @@ impl SwapOps for ZCoin { /// TODO: when all mm2 nodes upgrade to support the burn account then disable validation of the Standard option async fn validate_fee(&self, validate_fee_args: ValidateFeeArgs<'_>) -> ValidatePaymentResult<()> { let z_tx = match validate_fee_args.fee_tx { - TransactionEnum::ZTransaction(t) => t.clone(), + TransactionEnum::ZTransaction(t) => t, fee_tx => { return MmError::err(ValidatePaymentError::InternalError(format!( "Invalid fee tx type. fee tx: {:?}", @@ -1916,7 +1989,7 @@ impl InitWithdrawCoin for ZCoin { memo, }; - let (tx, data, _sync_guard) = self.gen_tx(vec![], vec![z_output]).await?; + let GenTxData { tx, data, .. } = self.gen_tx(vec![], vec![z_output]).await?; let mut tx_bytes = Vec::with_capacity(1024); tx.write(&mut tx_bytes) .map_to_mm(|e| WithdrawError::InternalError(e.to_string()))?; @@ -1925,9 +1998,10 @@ impl InitWithdrawCoin for ZCoin { let received_by_me = big_decimal_from_sat_unsigned(data.received_by_me, self.decimals()); let spent_by_me = big_decimal_from_sat_unsigned(data.spent_by_me, self.decimals()); + let tx_hash_hex = hex::encode(&tx_hash); Ok(TransactionDetails { - tx: TransactionData::new_signed(tx_bytes.into(), hex::encode(&tx_hash)), + tx: TransactionData::new_signed(tx_bytes.into(), tx_hash_hex), from: vec![self.z_fields.my_z_addr_encoded.clone()], to: vec![req.to], my_balance_change: &received_by_me - &spent_by_me, @@ -1949,6 +2023,108 @@ impl InitWithdrawCoin for ZCoin { } } +/// Waits until there are enough _unlocked_ Sapling notes to cover `total_required`. +/// TODO: Consider adding `wait_until` argument. +/// TODO: Integrate this into `light_wallet_db_sync_loop` instead of having a separate function. +/// Can be addressed when migrating to a newer librustzcash which supports spent note tracking. +/// See: https://github.com/KomodoPlatform/komodo-defi-framework/pull/2331#pullrequestreview-2883773336 +async fn wait_for_spendable_balance_impl( + selfi: ZCoin, + total_required: BigDecimal, +) -> Result, MmError> { + const MAX_RETRIES: usize = 40; + const RETRY_DELAY: f64 = 15.0; + + let mut retries = 0; + + loop { + let wallet_notes = selfi + .wallet_notes_ordered() + .await + .map_err(|e| GenTxError::SpendableNotesError(e.to_string()))?; + let wallet_notes_len = wallet_notes.len(); + + let locked_notes = selfi.z_fields.locked_notes_db.load_all_notes().await?; + + let unlocked_notes: Vec = if locked_notes.is_empty() { + wallet_notes + } else { + let unconfirmed_spent_rseeds: HashSet = locked_notes + .iter() + .filter_map(|n| { + if let LockedNote::Spent { rseed, .. } = n { + Some(rseed.clone()) + } else { + None + } + }) + .collect(); + + wallet_notes + .into_iter() + .filter(|note| !unconfirmed_spent_rseeds.contains(&rseed_to_string(¬e.rseed))) + .collect() + }; + let unlocked_notes_len = unlocked_notes.len(); + + let sum_available = unlocked_notes.iter().map(|n| n.note_value).sum::(); + let sum_available = u64::try_from(sum_available).map_to_mm(|err| GenTxError::Internal(err.to_string()))?; + let sum_available = big_decimal_from_sat_unsigned(sum_available, selfi.decimals()); + + // Reteurn InsufficientBalance error when all notes are unlocked but amount is insufficient. + if sum_available < total_required && unlocked_notes_len == wallet_notes_len { + return MmError::err(GenTxError::InsufficientBalance { + coin: selfi.ticker().to_string(), + available: sum_available, + required: total_required, + }); + } + + // Returns available notes when either sufficient funds exist or all notes are unlocked. + // Otherwise, waits for locked notes to become available up to MAX_RETRIES. + if sum_available >= total_required || unlocked_notes_len == wallet_notes_len { + return Ok(unlocked_notes.into_iter()); + } + + if retries >= MAX_RETRIES { + return MmError::err(GenTxError::Internal(format!( + "Locked notes did not become available after {} retries", + MAX_RETRIES + ))); + } + + info!( + "Locked notes present; retrying in {}s (attempt {}/{})", + RETRY_DELAY, + retries + 1, + MAX_RETRIES + ); + common::executor::Timer::sleep(RETRY_DELAY).await; + retries += 1; + } +} + +async fn wait_for_spendable_balance_spawner( + selfi: &ZCoin, + total_required: &BigDecimal, +) -> Result, MmError> { + let coin = selfi.clone(); + let required = total_required.clone(); + let (tx, rx) = oneshot::channel(); + + selfi.spawner().spawn(async move { + let result = wait_for_spendable_balance_impl(coin, required).await; + let _ = tx.send(result); + }); + + match rx.await { + Ok(res) => res, + Err(_) => MmError::err(GenTxError::Internal( + "wait_for_spendable_balance task was cancelled".into(), + )), + } +} + /// Interpret a string or hex-encoded memo, and return a Memo object. /// Inspired by https://github.com/adityapk00/zecwallet-light-cli/blob/v1.7.20/lib/src/lightwallet/utils.rs#L23 #[allow(clippy::result_large_err)] @@ -2012,6 +2188,16 @@ fn extended_spending_key_from_global_hd_account( Ok(spending_key) } +#[inline] +fn rseed_to_string(rseed: &Rseed) -> String { + const INPUT: [u8; 1] = [0x04]; + + match rseed { + Rseed::BeforeZip212(rcm) => rcm.to_string(), + Rseed::AfterZip212(rseed) => jubjub::Fr::from_bytes_wide(prf_expand(rseed, &INPUT).as_array()).to_string(), + } +} + #[test] fn derive_z_key_from_mm_seed() { use crypto::privkey::key_pair_from_seed; diff --git a/mm2src/coins/z_coin/storage.rs b/mm2src/coins/z_coin/storage.rs index e2534281b7..e8b1b6e971 100644 --- a/mm2src/coins/z_coin/storage.rs +++ b/mm2src/coins/z_coin/storage.rs @@ -1,16 +1,19 @@ use crate::z_coin::{ValidateBlocksError, ZcoinConsensusParams, ZcoinStorageError}; use mm2_event_stream::StreamingManager; -pub mod blockdb; -pub use blockdb::*; +pub(crate) mod blockdb; +pub(crate) use blockdb::*; + +pub(crate) mod walletdb; +pub(crate) use walletdb::*; + +pub(crate) mod z_locked_notes; +pub(crate) use z_locked_notes::{LockedNotesStorage, LockedNotesStorageError}; -pub mod walletdb; #[cfg(target_arch = "wasm32")] mod z_params; #[cfg(target_arch = "wasm32")] pub(crate) use z_params::ZcashParamsWasmImpl; -pub use walletdb::*; - use mm2_err_handle::mm_error::MmResult; #[cfg(target_arch = "wasm32")] use walletdb::wasm::storage::DataConnStmtCacheWasm; @@ -118,6 +121,7 @@ pub async fn scan_cached_block( data: &DataConnStmtCacheWrapper, params: &ZcoinConsensusParams, block: &CompactBlock, + locked_notes_db: &LockedNotesStorage, last_height: &mut BlockHeight, ) -> Result>, ValidateBlocksError> { let mut data_guard = data.inner().clone(); @@ -159,7 +163,6 @@ pub async fn scan_cached_block( // To enforce that all roots match, // see -> https://github.com/KomodoPlatform/librustzcash/blob/e92443a7bbd1c5e92e00e6deb45b5a33af14cea4/zcash_client_backend/src/data_api/chain.rs#L304-L326 - let new_witnesses = data_guard .advance_by_block( &(PrunedBlock { @@ -186,5 +189,15 @@ pub async fn scan_cached_block( witnesses.extend(new_witnesses); *last_height = current_height; + // TODO: Execute updates to `locked_notes_db` and `wallet_db` in a single transaction. + // This will be possible with a newer librustzcash that supports both spent notes and unconfirmed change tracking. + // See: https://github.com/KomodoPlatform/komodo-defi-framework/pull/2331#pullrequestreview-2883773336 + for tx in &txs { + locked_notes_db + .remove_notes_for_txid(tx.txid.to_string()) + .await + .map_err(|err| ValidateBlocksError::DbError(err.to_string()))?; + } + Ok(txs) } diff --git a/mm2src/coins/z_coin/storage/blockdb/blockdb_idb_storage.rs b/mm2src/coins/z_coin/storage/blockdb/blockdb_idb_storage.rs index 826ed52bdd..b95af8c4cd 100644 --- a/mm2src/coins/z_coin/storage/blockdb/blockdb_idb_storage.rs +++ b/mm2src/coins/z_coin/storage/blockdb/blockdb_idb_storage.rs @@ -1,5 +1,5 @@ use crate::z_coin::storage::{scan_cached_block, validate_chain, BlockDbImpl, BlockProcessingMode, CompactBlockRow, - ZcoinConsensusParams, ZcoinStorageRes}; + LockedNotesStorage, ZcoinConsensusParams, ZcoinStorageRes}; use crate::z_coin::tx_history_events::ZCoinTxHistoryEventStreamer; use crate::z_coin::z_balance_streaming::ZCoinBalanceEventStreamer; use crate::z_coin::z_coin_errors::ZcoinStorageError; @@ -10,7 +10,6 @@ use mm2_db::indexed_db::{BeBigUint, ConstructibleDb, DbIdentifier, DbInstance, D IndexedDbBuilder, InitDbResult, MultiIndex, OnUpgradeResult, TableSignature}; use mm2_err_handle::prelude::*; use protobuf::Message; -use std::path::PathBuf; use zcash_client_backend::proto::compact_formats::CompactBlock; use zcash_extras::WalletRead; use zcash_primitives::block::BlockHash; @@ -68,7 +67,7 @@ impl BlockDbInner { } impl BlockDbImpl { - pub async fn new(ctx: &MmArc, ticker: String, _path: PathBuf) -> ZcoinStorageRes { + pub async fn new(ctx: &MmArc, ticker: String) -> ZcoinStorageRes { Ok(Self { db: ConstructibleDb::new(ctx).into_shared(), ticker, @@ -221,6 +220,7 @@ impl BlockDbImpl { mode: BlockProcessingMode, validate_from: Option<(BlockHeight, BlockHash)>, limit: Option, + locked_notes_db: &LockedNotesStorage, ) -> ZcoinStorageRes<()> { let ticker = self.ticker.to_owned(); let mut from_height = match &mode { @@ -254,7 +254,7 @@ impl BlockDbImpl { validate_chain(block, &mut prev_height, &mut prev_hash).await?; }, BlockProcessingMode::Scan(data, streaming_manager) => { - let txs = scan_cached_block(data, ¶ms, &block, &mut from_height).await?; + let txs = scan_cached_block(data, ¶ms, &block, locked_notes_db, &mut from_height).await?; if !txs.is_empty() { // Stream out the new transactions. streaming_manager diff --git a/mm2src/coins/z_coin/storage/blockdb/blockdb_sql_storage.rs b/mm2src/coins/z_coin/storage/blockdb/blockdb_sql_storage.rs index 8dd4dd39f7..6083751df6 100644 --- a/mm2src/coins/z_coin/storage/blockdb/blockdb_sql_storage.rs +++ b/mm2src/coins/z_coin/storage/blockdb/blockdb_sql_storage.rs @@ -1,5 +1,5 @@ use crate::z_coin::storage::{scan_cached_block, validate_chain, BlockDbImpl, BlockProcessingMode, CompactBlockRow, - ZcoinStorageRes}; + LockedNotesStorage, ZcoinStorageRes}; use crate::z_coin::tx_history_events::ZCoinTxHistoryEventStreamer; use crate::z_coin::z_balance_streaming::ZCoinBalanceEventStreamer; use crate::z_coin::z_coin_errors::ZcoinStorageError; @@ -12,7 +12,6 @@ use itertools::Itertools; use mm2_core::mm_ctx::MmArc; use mm2_err_handle::prelude::*; use protobuf::Message; -use std::path::PathBuf; use std::sync::{Arc, Mutex}; use zcash_client_backend::data_api::error::Error as ChainError; use zcash_client_backend::proto::compact_formats::CompactBlock; @@ -46,7 +45,8 @@ impl From> for ZcoinStorageError { impl BlockDbImpl { #[cfg(not(test))] - pub async fn new(_ctx: &MmArc, ticker: String, path: PathBuf) -> ZcoinStorageRes { + pub async fn new(ctx: &MmArc, ticker: String) -> ZcoinStorageRes { + let path = ctx.global_dir().join(format!("{}_cache.db", ticker)); async_blocking(move || { mm2_io::fs::create_parents(&path).map_err(|err| ZcoinStorageError::IoError(err.to_string()))?; let conn = Connection::open(path).map_to_mm(|err| ZcoinStorageError::DbError(err.to_string()))?; @@ -70,7 +70,7 @@ impl BlockDbImpl { } #[cfg(test)] - pub(crate) async fn new(ctx: &MmArc, ticker: String, _path: PathBuf) -> ZcoinStorageRes { + pub(crate) async fn new(ctx: &MmArc, ticker: String) -> ZcoinStorageRes { let ctx = ctx.clone(); async_blocking(move || { let conn = ctx @@ -192,6 +192,7 @@ impl BlockDbImpl { mode: BlockProcessingMode, validate_from: Option<(BlockHeight, BlockHash)>, limit: Option, + locked_notes_db: &LockedNotesStorage, ) -> ZcoinStorageRes<()> { let ticker = self.ticker.to_owned(); let mut from_height = match &mode { @@ -230,7 +231,7 @@ impl BlockDbImpl { validate_chain(block, &mut prev_height, &mut prev_hash).await?; }, BlockProcessingMode::Scan(data, streaming_manager) => { - let txs = scan_cached_block(data, ¶ms, &block, &mut from_height).await?; + let txs = scan_cached_block(data, ¶ms, &block, locked_notes_db, &mut from_height).await?; if !txs.is_empty() { // Stream out the new transactions. streaming_manager diff --git a/mm2src/coins/z_coin/storage/blockdb/mod.rs b/mm2src/coins/z_coin/storage/blockdb/mod.rs index bc41c4de00..e1b9c13d5b 100644 --- a/mm2src/coins/z_coin/storage/blockdb/mod.rs +++ b/mm2src/coins/z_coin/storage/blockdb/mod.rs @@ -25,7 +25,6 @@ pub struct BlockDbImpl { mod block_db_storage_tests { use crate::z_coin::storage::BlockDbImpl; use common::log::info; - use std::path::PathBuf; use mm2_test_helpers::for_tests::mm_ctx_with_custom_db; @@ -37,9 +36,7 @@ mod block_db_storage_tests { pub(crate) async fn test_insert_block_and_get_latest_block_impl() { let ctx = mm_ctx_with_custom_db(); - let db = BlockDbImpl::new(&ctx, TICKER.to_string(), PathBuf::new()) - .await - .unwrap(); + let db = BlockDbImpl::new(&ctx, TICKER.to_string()).await.unwrap(); // insert block for header in HEADERS.iter() { db.insert_block(header.0, hex::decode(header.1).unwrap()).await.unwrap(); @@ -52,9 +49,7 @@ mod block_db_storage_tests { pub(crate) async fn test_rewind_to_height_impl() { let ctx = mm_ctx_with_custom_db(); - let db = BlockDbImpl::new(&ctx, TICKER.to_string(), PathBuf::new()) - .await - .unwrap(); + let db = BlockDbImpl::new(&ctx, TICKER.to_string()).await.unwrap(); // insert block for header in HEADERS.iter() { db.insert_block(header.0, hex::decode(header.1).unwrap()).await.unwrap(); @@ -77,9 +72,7 @@ mod block_db_storage_tests { #[allow(unused)] pub(crate) async fn test_process_blocks_with_mode_impl() { let ctx = mm_ctx_with_custom_db(); - let db = BlockDbImpl::new(&ctx, TICKER.to_string(), PathBuf::new()) - .await - .unwrap(); + let db = BlockDbImpl::new(&ctx, TICKER.to_string()).await.unwrap(); // insert block for header in HEADERS.iter() { let inserted_id = db.insert_block(header.0, hex::decode(header.1).unwrap()).await.unwrap(); diff --git a/mm2src/coins/z_coin/storage/walletdb/wasm/mod.rs b/mm2src/coins/z_coin/storage/walletdb/wasm/mod.rs index c1ffdfb0a2..5fb428db4f 100644 --- a/mm2src/coins/z_coin/storage/walletdb/wasm/mod.rs +++ b/mm2src/coins/z_coin/storage/walletdb/wasm/mod.rs @@ -68,6 +68,7 @@ fn to_spendable_note(note: SpendableNoteConstructor) -> MmResult LocalTxProver { let (spend_buf, output_buf) = wagyu_zcash_parameters::load_sapling_parameters(); @@ -206,6 +208,9 @@ mod wasm_test { async fn test_valid_chain_state() { // init blocks_db let ctx = mm_ctx_with_custom_db(); + let locked_notes_db = LockedNotesStorage::new(ctx.clone(), MY_ADDRESS.to_string()) + .await + .unwrap(); let blockdb = BlockDbImpl::new(&ctx, TICKER.to_string(), PathBuf::new()) .await .unwrap(); @@ -226,6 +231,7 @@ mod wasm_test { BlockProcessingMode::Validate, walletdb.get_max_height_hash().await.unwrap(), None, + &locked_notes_db, ) .await .unwrap(); @@ -247,6 +253,7 @@ mod wasm_test { BlockProcessingMode::Validate, walletdb.get_max_height_hash().await.unwrap(), None, + &locked_notes_db, ) .await .unwrap(); @@ -259,6 +266,7 @@ mod wasm_test { BlockProcessingMode::Scan(scan, StreamingManager::default()), None, None, + &locked_notes_db, ) .await .unwrap(); @@ -271,6 +279,7 @@ mod wasm_test { BlockProcessingMode::Validate, max_height_hash, None, + &locked_notes_db, ) .await .unwrap(); @@ -292,6 +301,7 @@ mod wasm_test { BlockProcessingMode::Validate, walletdb.get_max_height_hash().await.unwrap(), None, + &locked_notes_db, ) .await .unwrap(); @@ -304,6 +314,7 @@ mod wasm_test { BlockProcessingMode::Scan(scan, StreamingManager::default()), None, None, + &locked_notes_db, ) .await .unwrap(); @@ -315,6 +326,7 @@ mod wasm_test { BlockProcessingMode::Validate, walletdb.get_max_height_hash().await.unwrap(), None, + &locked_notes_db, ) .await .unwrap(); @@ -324,6 +336,9 @@ mod wasm_test { async fn invalid_chain_cache_disconnected() { // init blocks_db let ctx = mm_ctx_with_custom_db(); + let locked_notes_db = LockedNotesStorage::new(ctx.clone(), MY_ADDRESS.to_string()) + .await + .unwrap(); let blockdb = BlockDbImpl::new(&ctx, TICKER.to_string(), PathBuf::new()) .await .unwrap(); @@ -363,6 +378,7 @@ mod wasm_test { BlockProcessingMode::Scan(scan, StreamingManager::default()), None, None, + &locked_notes_db, ) .await .unwrap(); @@ -374,6 +390,7 @@ mod wasm_test { BlockProcessingMode::Validate, walletdb.get_max_height_hash().await.unwrap(), None, + &locked_notes_db, ) .await .unwrap(); @@ -403,6 +420,7 @@ mod wasm_test { BlockProcessingMode::Validate, walletdb.get_max_height_hash().await.unwrap(), None, + &locked_notes_db, ) .await .unwrap_err(); @@ -418,6 +436,9 @@ mod wasm_test { async fn test_invalid_chain_reorg() { // init blocks_db let ctx = mm_ctx_with_custom_db(); + let locked_notes_db = LockedNotesStorage::new(ctx.clone(), MY_ADDRESS.to_string()) + .await + .unwrap(); let blockdb = BlockDbImpl::new(&ctx, TICKER.to_string(), PathBuf::new()) .await .unwrap(); @@ -457,6 +478,7 @@ mod wasm_test { BlockProcessingMode::Scan(scan, StreamingManager::default()), None, None, + &locked_notes_db, ) .await .unwrap(); @@ -468,6 +490,7 @@ mod wasm_test { BlockProcessingMode::Validate, walletdb.get_max_height_hash().await.unwrap(), None, + &locked_notes_db, ) .await .unwrap(); @@ -497,6 +520,7 @@ mod wasm_test { BlockProcessingMode::Validate, walletdb.get_max_height_hash().await.unwrap(), None, + &locked_notes_db, ) .await .unwrap_err(); @@ -512,6 +536,9 @@ mod wasm_test { async fn test_data_db_rewinding() { // init blocks_db let ctx = mm_ctx_with_custom_db(); + let locked_notes_db = LockedNotesStorage::new(ctx.clone(), MY_ADDRESS.to_string()) + .await + .unwrap(); let blockdb = BlockDbImpl::new(&ctx, TICKER.to_string(), PathBuf::new()) .await .unwrap(); @@ -546,6 +573,7 @@ mod wasm_test { BlockProcessingMode::Scan(scan, StreamingManager::default()), None, None, + &locked_notes_db, ) .await .unwrap(); @@ -576,6 +604,7 @@ mod wasm_test { BlockProcessingMode::Scan(scan, StreamingManager::default()), None, None, + &locked_notes_db, ) .await .unwrap(); @@ -588,6 +617,9 @@ mod wasm_test { async fn test_scan_cached_blocks_requires_sequential_blocks() { // init blocks_db let ctx = mm_ctx_with_custom_db(); + let locked_notes_db = LockedNotesStorage::new(ctx.clone(), MY_ADDRESS.to_string()) + .await + .unwrap(); let blockdb = BlockDbImpl::new(&ctx, TICKER.to_string(), PathBuf::new()) .await .unwrap(); @@ -615,6 +647,7 @@ mod wasm_test { BlockProcessingMode::Scan(scan, StreamingManager::default()), None, None, + &locked_notes_db, ) .await .unwrap(); @@ -633,6 +666,7 @@ mod wasm_test { BlockProcessingMode::Scan(scan, StreamingManager::default()), None, None, + &locked_notes_db, ) .await .unwrap_err(); @@ -656,7 +690,8 @@ mod wasm_test { consensus_params.clone(), BlockProcessingMode::Scan(scan, StreamingManager::default()), None, - None + None, + &locked_notes_db ) .await .is_ok()); @@ -671,6 +706,9 @@ mod wasm_test { async fn test_scan_cached_blokcs_finds_received_notes() { // init blocks_db let ctx = mm_ctx_with_custom_db(); + let locked_notes_db = LockedNotesStorage::new(ctx.clone(), MY_ADDRESS.to_string()) + .await + .unwrap(); let blockdb = BlockDbImpl::new(&ctx, TICKER.to_string(), PathBuf::new()) .await .unwrap(); @@ -700,7 +738,8 @@ mod wasm_test { consensus_params.clone(), BlockProcessingMode::Scan(scan, StreamingManager::default()), None, - None + None, + &locked_notes_db ) .await .is_ok()); @@ -721,7 +760,8 @@ mod wasm_test { consensus_params.clone(), BlockProcessingMode::Scan(scan, StreamingManager::default()), None, - None + None, + &locked_notes_db ) .await .is_ok()); @@ -734,6 +774,9 @@ mod wasm_test { async fn test_scan_cached_blocks_finds_change_notes() { // init blocks_db let ctx = mm_ctx_with_custom_db(); + let locked_notes_db = LockedNotesStorage::new(ctx.clone(), MY_ADDRESS.to_string()) + .await + .unwrap(); let blockdb = BlockDbImpl::new(&ctx, TICKER.to_string(), PathBuf::new()) .await .unwrap(); @@ -763,7 +806,8 @@ mod wasm_test { consensus_params.clone(), BlockProcessingMode::Scan(scan, StreamingManager::default()), None, - None + None, + &locked_notes_db ) .await .is_ok()); @@ -794,6 +838,7 @@ mod wasm_test { BlockProcessingMode::Scan(scan, StreamingManager::default()), None, None, + &locked_notes_db, ) .await; assert!(scan.is_ok()); @@ -810,6 +855,7 @@ mod wasm_test { // async fn create_to_address_fails_on_unverified_notes() { // // init blocks_db // let ctx = mm_ctx_with_custom_db(); + // let locked_notes_db = LockedNotesStorage::new(ctx.clone(), MY_ADDRESS.to_string()).await.unwrap(); // let blockdb = BlockDbImpl::new(&ctx, TICKER.to_string(), PathBuf::new()).await.unwrap(); // // // init walletdb. @@ -853,7 +899,7 @@ mod wasm_test { // // Scan the cache // let scan = DataConnStmtCacheWrapper::new(DataConnStmtCacheWasm(walletdb.clone())); // assert!(blockdb - // .process_blocks_with_mode(consensus_params.clone(), BlockProcessingMode::Scan(scan, StreamingManager::default()), None, None) + // .process_blocks_with_mode(consensus_params.clone(), BlockProcessingMode::Scan(scan, StreamingManager::default()), None, None, &locked_notes_db) // .await // .is_ok()); // @@ -898,7 +944,7 @@ mod wasm_test { // // Scan the cache // let scan = DataConnStmtCacheWrapper::new(DataConnStmtCacheWasm(walletdb.clone())); // assert!(blockdb - // .process_blocks_with_mode(consensus_params.clone(), BlockProcessingMode::Scan(scan, StreamingManager::default()), None, None) + // .process_blocks_with_mode(consensus_params.clone(), BlockProcessingMode::Scan(scan, StreamingManager::default()), None, None, &locked_notes_db) // .await // .is_ok()); // @@ -929,7 +975,7 @@ mod wasm_test { // // Scan the cache // let scan = DataConnStmtCacheWrapper::new(DataConnStmtCacheWasm(walletdb.clone())); // assert!(blockdb - // .process_blocks_with_mode(consensus_params.clone(), BlockProcessingMode::Scan(scan, StreamingManager::default()), None, None) + // .process_blocks_with_mode(consensus_params.clone(), BlockProcessingMode::Scan(scan, StreamingManager::default()), None, None, &locked_notes_db) // .await // .is_ok()); // @@ -1079,6 +1125,7 @@ mod wasm_test { // // // init blocks_db // let ctx = mm_ctx_with_custom_db(); + // let locked_notes_db = LockedNotesStorage::new(ctx.clone(), MY_ADDRESS.to_string()).await.unwrap(); // let blockdb = BlockDbImpl::new(&ctx, TICKER.to_string(), PathBuf::new()).await.unwrap(); // // // init walletdb. @@ -1099,7 +1146,7 @@ mod wasm_test { // // Scan the cache // let scan = DataConnStmtCacheWrapper::new(DataConnStmtCacheWasm(walletdb.clone())); // blockdb - // .process_blocks_with_mode(consensus_params.clone(), BlockProcessingMode::Scan(scan, StreamingManager::default()), None, None) + // .process_blocks_with_mode(consensus_params.clone(), BlockProcessingMode::Scan(scan, StreamingManager::default()), None, None, &locked_notes_db) // .await // .unwrap(); // assert_eq!(walletdb.get_balance(AccountId(0)).await.unwrap(), value); @@ -1156,7 +1203,7 @@ mod wasm_test { // // Scan the cache // let scan = DataConnStmtCacheWrapper::new(DataConnStmtCacheWasm(walletdb.clone())); // blockdb - // .process_blocks_with_mode(consensus_params.clone(), BlockProcessingMode::Scan(scan, StreamingManager::default()), None, None) + // .process_blocks_with_mode(consensus_params.clone(), BlockProcessingMode::Scan(scan, StreamingManager::default()), None, None, &locked_notes_db) // .await // .unwrap(); // @@ -1192,7 +1239,7 @@ mod wasm_test { // // Scan the cache // let scan = DataConnStmtCacheWrapper::new(DataConnStmtCacheWasm(walletdb.clone())); // blockdb - // .process_blocks_with_mode(consensus_params.clone(), BlockProcessingMode::Scan(scan, StreamingManager::default()), None, None) + // .process_blocks_with_mode(consensus_params.clone(), BlockProcessingMode::Scan(scan, StreamingManager::default()), None, None, &locked_notes_db) // .await // .unwrap(); // diff --git a/mm2src/coins/z_coin/storage/walletdb/wasm/storage.rs b/mm2src/coins/z_coin/storage/walletdb/wasm/storage.rs index 06d06c4d32..78d8a6fbaa 100644 --- a/mm2src/coins/z_coin/storage/walletdb/wasm/storage.rs +++ b/mm2src/coins/z_coin/storage/walletdb/wasm/storage.rs @@ -1082,7 +1082,7 @@ impl WalletRead for WalletIndexedDb { let matching_tx = maybe_txs.iter().find(|(id_tx, _tx)| id_tx.to_bigint() == note.spent); if let Some((_, tx)) = matching_tx { - if tx.block.is_none() { + if tx.block.is_some() { nullifiers.push(( AccountId( note.account @@ -1095,18 +1095,6 @@ impl WalletRead for WalletIndexedDb { .unwrap(), )); } - } else { - nullifiers.push(( - AccountId( - note.account - .to_u32() - .ok_or_else(|| ZcoinStorageError::GetFromStorageError("Invalid amount".to_string()))?, - ), - Nullifier::from_slice(¬e.nf.clone().ok_or_else(|| { - ZcoinStorageError::GetFromStorageError("Error while putting tx_meta".to_string()) - })?) - .unwrap(), - )); } } diff --git a/mm2src/coins/z_coin/storage/walletdb/wasm/tables.rs b/mm2src/coins/z_coin/storage/walletdb/wasm/tables.rs index 53471571d4..92ce152837 100644 --- a/mm2src/coins/z_coin/storage/walletdb/wasm/tables.rs +++ b/mm2src/coins/z_coin/storage/walletdb/wasm/tables.rs @@ -141,11 +141,6 @@ impl WalletDbReceivedNotesTable { pub const TICKER_ACCOUNT_INDEX: &'static str = "ticker_account_index"; /// A **unique** index that consists of the following properties: /// * ticker - /// * note_id - /// * nf - pub const TICKER_NOTES_ID_NF_INDEX: &'static str = "ticker_note_id_nf_index"; - /// A **unique** index that consists of the following properties: - /// * ticker /// * tx /// * output_index pub const TICKER_TX_OUTPUT_INDEX: &'static str = "ticker_tx_output_index"; diff --git a/mm2src/coins/z_coin/storage/z_locked_notes/mod.rs b/mm2src/coins/z_coin/storage/z_locked_notes/mod.rs new file mode 100644 index 0000000000..fdbd9b8b76 --- /dev/null +++ b/mm2src/coins/z_coin/storage/z_locked_notes/mod.rs @@ -0,0 +1,155 @@ +use enum_derives::EnumFromStringify; + +cfg_native!( + pub(crate) mod sqlite; + + use db_common::async_sql_conn::{AsyncConnError, AsyncConnection}; + use futures::lock::Mutex; + use std::sync::Arc; +); + +cfg_wasm32!( + pub(crate) mod wasm; + + use self::wasm::LockedNoteDbInner; + use mm2_db::indexed_db::{DbTransactionError, InitDbError, SharedDb}; +); + +/// Represents a shielded note temporarily locked due to a pending transaction. +/// Locked notes are excluded from the spendable balance until confirmed or cleared. +#[derive(Debug, Clone)] +pub(crate) enum LockedNote { + /// A note being spent by a pending shielded transaction (`rseed` is the note's randomness). + Spent { rseed: String }, + + /// A pending change output from an unconfirmed shielded transaction (`value` is the expected amount). + Change { value: u64 }, +} + +/// A wrapper for the db connection to the change note cache database in native and browser. +#[derive(Clone)] +pub struct LockedNotesStorage { + #[cfg(not(target_arch = "wasm32"))] + pub db: Arc>, + #[cfg(target_arch = "wasm32")] + pub db: SharedDb, + #[allow(unused)] + address: String, +} + +#[derive(Clone, Debug, Display, Eq, PartialEq, EnumFromStringify)] +pub(crate) enum LockedNotesStorageError { + #[cfg(not(target_arch = "wasm32"))] + #[display(fmt = "Sqlite Error: {_0}")] + #[from_stringify("AsyncConnError", "db_common::sqlite::rusqlite::Error")] + SqliteError(String), + #[cfg(target_arch = "wasm32")] + #[display(fmt = "IndexedDb Error: {_0}")] + #[from_stringify("InitDbError", "DbTransactionError")] + IndexedDbError(String), +} + +#[cfg(any(test, target_arch = "wasm32"))] +pub(super) mod locked_notes_test { + use crate::z_coin::storage::z_locked_notes::{LockedNote, LockedNotesStorage}; + use common::cross_test; + use mm2_test_helpers::for_tests::mm_ctx_with_custom_db; + + common::cfg_wasm32! { + use wasm_bindgen_test::*; + wasm_bindgen_test::wasm_bindgen_test_configure!(run_in_browser); + } + + const MY_ADDRESS: &str = "my_address"; + + cross_test!(test_insert_and_remove_note, { + let ctx = mm_ctx_with_custom_db(); + let db = LockedNotesStorage::new(&ctx, MY_ADDRESS.to_string()).await.unwrap(); + + // Insert a pending spent note + let spent_txid = "0x18b1acd8ceae8d71a2ae8b7e4a3e48ceb39dc237f0aa38c468425b88dc8d5f3e".to_string(); + let spent_rseed = "0xcfec34a81e67e85aa1ce1a6666f92f9bc5606f0795be555bb3c9f9ac089aa4f7".to_string(); + db.insert_spent_note(spent_txid.clone(), spent_rseed.clone()) + .await + .unwrap(); + + // Insert a pending change note + let change_txid = "0xdeadbeefdeadbeefdeadbeefdeadbeefdeadbeefdeadbeefdeadbeefdeadbeef".to_string(); + let change_value = 123456; + db.insert_change_note(change_txid.clone(), change_value).await.unwrap(); + + // Remove by txid + db.remove_notes_for_txid(spent_txid.clone()).await.unwrap(); + db.remove_notes_for_txid(change_txid.clone()).await.unwrap(); + + let notes = db.load_all_notes().await.unwrap(); + assert!(notes.is_empty()); + + // Insert both again but using same txid + db.insert_spent_note(spent_txid.clone(), spent_rseed.clone()) + .await + .unwrap(); + db.insert_change_note(spent_txid.clone(), change_value).await.unwrap(); + + // Remove by txid (removes both input and output if same txid) + db.remove_notes_for_txid(spent_txid.clone()).await.unwrap(); + + let notes = db.load_all_notes().await.unwrap(); + assert!(notes.is_empty()); + }); + + cross_test!(test_load_all_notes, { + let ctx = mm_ctx_with_custom_db(); + let db = LockedNotesStorage::new(&ctx, MY_ADDRESS.to_string()).await.unwrap(); + + let spent_txid = "0x01".to_string(); + let spent_rseed = "0xcafe000000000000000000000000000000000000000000000000000000000000".to_string(); + let change_txid = "0x02".to_string(); + let change_value = 123456789; + + db.insert_spent_note(spent_txid.clone(), spent_rseed.clone()) + .await + .unwrap(); + db.insert_change_note(change_txid.clone(), change_value).await.unwrap(); + + let notes = db.load_all_notes().await.unwrap(); + + assert_eq!(notes.len(), 2); + + match ¬es[0] { + LockedNote::Spent { rseed } => { + assert_eq!(rseed, &spent_rseed); + }, + _ => panic!("First note should be a Spent note"), + } + match ¬es[1] { + LockedNote::Change { value } => { + assert_eq!(*value, change_value); + }, + _ => panic!("Second note should be a Change note"), + } + }); + + cross_test!(test_sum_changes, { + let ctx = mm_ctx_with_custom_db(); + let db = LockedNotesStorage::new(&ctx, MY_ADDRESS.to_string()).await.unwrap(); + + db.insert_change_note("txid1".to_string(), 1000).await.unwrap(); + db.insert_change_note("txid2".to_string(), 2000).await.unwrap(); + db.insert_spent_note("0xinputrseed".to_string(), "txid3".to_string()) + .await + .unwrap(); + + let notes = db.load_all_notes().await.unwrap(); + + // Only sum Output note values + let sum: u64 = notes + .iter() + .filter_map(|n| match n { + LockedNote::Change { value, .. } => Some(*value), + _ => None, + }) + .sum(); + assert_eq!(sum, 3000); + }); +} diff --git a/mm2src/coins/z_coin/storage/z_locked_notes/sqlite.rs b/mm2src/coins/z_coin/storage/z_locked_notes/sqlite.rs new file mode 100644 index 0000000000..1738c82025 --- /dev/null +++ b/mm2src/coins/z_coin/storage/z_locked_notes/sqlite.rs @@ -0,0 +1,158 @@ +use super::{LockedNote, LockedNotesStorage, LockedNotesStorageError}; +use db_common::async_sql_conn::{AsyncConnError, AsyncConnection}; +use db_common::sqlite::run_optimization_pragmas; +use db_common::sqlite::rusqlite::params; +use futures::lock::Mutex; +use itertools::Itertools; +use mm2_core::mm_ctx::MmArc; +use mm2_err_handle::prelude::*; +use std::convert::TryInto; +use std::sync::Arc; + +const TABLE_NAME: &str = "locked_notes_cache"; + +async fn create_table(conn: Arc>) -> Result<(), AsyncConnError> { + let conn = conn.lock().await; + conn.call(move |conn| { + run_optimization_pragmas(conn)?; + conn.execute( + &format!( + "CREATE TABLE IF NOT EXISTS {TABLE_NAME} ( + variant TEXT NOT NULL, -- 'Spent' or 'Change' + txid VARCHAR NOT NULL, + rseed VARCHAR, -- only for Spent + value INTEGER, -- only for Change + UNIQUE (variant, txid, rseed, value) + )" + ), + [], + )?; + Ok(()) + }).await +} + +impl LockedNotesStorage { + #[cfg(not(any(test, feature = "run-docker-tests")))] + pub(crate) async fn new(ctx: &MmArc, address: String) -> MmResult { + let path = ctx.wallet_dir().join(format!("{}_locked_notes_cache.db", address)); + let db = AsyncConnection::open(path) + .await + .map_to_mm(|err| LockedNotesStorageError::SqliteError(err.to_string()))?; + let db = Arc::new(Mutex::new(db)); + + create_table(db.clone()).await?; + + Ok(Self { db, address }) + } + + #[cfg(any(test, feature = "run-docker-tests"))] + pub(crate) async fn new(ctx: &MmArc, address: String) -> MmResult { + #[cfg(feature = "run-docker-tests")] + let db = { + let path = ctx.wallet_dir().join(format!("{}_locked_notes_cache.db", address)); + mm2_io::fs::create_parents_async(&path) + .await + .map_err(|err| LockedNotesStorageError::SqliteError(err.to_string()))?; + Arc::new(Mutex::new( + AsyncConnection::open(path) + .await + .map_to_mm(|err| LockedNotesStorageError::SqliteError(err.to_string()))?, + )) + }; + #[cfg(all(test, not(feature = "run-docker-tests")))] + let db = { + let test_conn = Arc::new(Mutex::new(AsyncConnection::open_in_memory().await.unwrap())); + ctx.async_sqlite_connection.get().cloned().unwrap_or(test_conn) + }; + + create_table(db.clone()).await?; + + Ok(Self { db, address }) + } + + pub(crate) async fn insert_spent_note( + &self, + txid: String, + rseed: String, + ) -> MmResult<(), LockedNotesStorageError> { + let db = self.db.lock().await; + Ok(db.call(move |conn| { + conn.prepare(&format!( + "INSERT OR REPLACE INTO {TABLE_NAME} (variant, txid, rseed, value) VALUES (?, ?, ?, NULL)" + ))? + .execute(params!["Spent", txid, rseed])?; + Ok(()) + }).await?) + } + + pub(crate) async fn insert_change_note( + &self, + txid: String, + value: u64, + ) -> MmResult<(), LockedNotesStorageError> { + let db = self.db.lock().await; + Ok(db.call(move |conn| { + conn.prepare(&format!( + "INSERT OR REPLACE INTO {TABLE_NAME} (variant, txid, rseed, value) VALUES (?, ?, NULL, ?)" + ))? + .execute(params!["Change", txid, value as i64])?; + Ok(()) + }).await?) + } + + pub(crate) async fn remove_notes_for_txid(&self, txid: String) -> MmResult<(), LockedNotesStorageError> { + let db = self.db.lock().await; + Ok(db + .call(move |conn| { + conn.execute( + &format!("DELETE FROM {TABLE_NAME} WHERE txid=?"), + [&txid], + )?; + Ok(()) + }) + .await?) + } + + pub(crate) async fn load_all_notes(&self) -> MmResult, LockedNotesStorageError> { + let db = self.db.lock().await; + Ok(db.call(move |conn| { + let mut stmt = conn.prepare(&format!( + "SELECT variant, txid, rseed, value FROM {TABLE_NAME};" + ))?; + let rows = stmt.query_map(params![], |row| { + let variant: String = row.get(0)?; + let rseed: Option = row.get(2)?; + let value: Option = row.get(3)?; + + match variant.as_str() { + "Spent" => { + let rseed = rseed.ok_or_else(|| db_common::sqlite::rusqlite::Error::FromSqlConversionFailure( + 2, // Column index for "rseed" + db_common::sqlite::rusqlite::types::Type::Text, + "NULL value found for required rseed field".into() + ))?; + Ok(LockedNote::Spent { rseed }) + }, + "Change" => { + let i64_value = value.ok_or_else(|| db_common::sqlite::rusqlite::Error::FromSqlConversionFailure( + 3, // Column index for "value" + db_common::sqlite::rusqlite::types::Type::Integer, + "NULL value found for required value field".into() + ))?; + + let value = i64_value.try_into() + .map_err(|_| db_common::sqlite::rusqlite::Error::IntegralValueOutOfRange(3, i64_value))?; + + Ok(LockedNote::Change { value }) + }, + unexpected => Err(db_common::sqlite::rusqlite::Error::FromSqlConversionFailure( + 0, // Column index for "variant" + db_common::sqlite::rusqlite::types::Type::Text, + format!("Unexpected variant value: {}", unexpected).into() + )), + } + })?; + Ok(rows.flatten().collect_vec()) + }).await?) + } +} diff --git a/mm2src/coins/z_coin/storage/z_locked_notes/wasm.rs b/mm2src/coins/z_coin/storage/z_locked_notes/wasm.rs new file mode 100644 index 0000000000..dc72ca9d4e --- /dev/null +++ b/mm2src/coins/z_coin/storage/z_locked_notes/wasm.rs @@ -0,0 +1,158 @@ +use super::{LockedNote, LockedNotesStorage, LockedNotesStorageError}; + +use mm2_core::mm_ctx::MmArc; +use mm2_db::indexed_db::{ConstructibleDb, DbIdentifier, DbInstance, DbLocked, DbUpgrader, IndexedDb, IndexedDbBuilder, + InitDbResult, OnUpgradeResult, TableSignature, OnUpgradeError}; +use mm2_err_handle::prelude::*; + +const DB_NAME: &str = "z_change_note_storage"; +const DB_VERSION: u32 = 1; + +pub type LockedNotesDbInnerLocked<'a> = DbLocked<'a, LockedNoteDbInner>; + +#[derive(Debug, Clone, Deserialize, Serialize)] +pub struct LockedNoteTable { + address: String, + variant: String, // "Spent" or "Change" + txid: String, + rseed: Option, // Only for Spent + value: Option, // Only for Change +} + +impl TableSignature for LockedNoteTable { + const TABLE_NAME: &'static str = "change_notes"; + + fn on_upgrade_needed(upgrader: &DbUpgrader, mut old_version: u32, new_version: u32) -> OnUpgradeResult<()> { + while old_version < new_version { + match old_version { + 0 => { + let table = upgrader.create_table(Self::TABLE_NAME)?; + table.create_index("address", false)?; + table.create_index("variant", false)?; + table.create_index("txid", false)?; + table.create_index("rseed", false)?; + table.create_index("value", false)?; + } + unsupported_version => { + return MmError::err(OnUpgradeError::UnsupportedVersion { + unsupported_version, + old_version, + new_version, + }); + } + } + old_version += 1; + } + Ok(()) + } +} + +pub struct LockedNoteDbInner(IndexedDb); + +#[async_trait::async_trait] +impl DbInstance for LockedNoteDbInner { + const DB_NAME: &'static str = DB_NAME; + + async fn init(db_id: DbIdentifier) -> InitDbResult { + let inner = IndexedDbBuilder::new(db_id) + .with_version(DB_VERSION) + .with_table::() + .build() + .await?; + + Ok(Self(inner)) + } +} + +impl LockedNoteDbInner { + pub fn get_inner(&self) -> &IndexedDb { &self.0 } +} + +impl LockedNotesStorage { + async fn lockdb(&self) -> MmResult { + Ok(self.db.get_or_initialize().await?) + } +} + +impl LockedNotesStorage { + pub(crate) async fn new(ctx: &MmArc, address: String) -> Result { + let db = ConstructibleDb::new(ctx).into_shared(); + Ok(Self { address, db }) + } + + pub(crate) async fn insert_spent_note( + &self, + txid: String, + rseed: String, + ) -> MmResult<(), LockedNotesStorageError> { + let db = self.lockdb().await?; + let address = self.address.clone(); + let transaction = db.get_inner().transaction().await?; + let change_note_table = transaction.table::().await?; + + let change_note = LockedNoteTable { + address, + variant: "Spent".to_owned(), + txid, + rseed: Some(rseed), + value: None, + }; + Ok(change_note_table + .add_item(&change_note) + .await + .map(|_| ())?) + } + + pub(crate) async fn insert_change_note( + &self, + txid: String, + value: u64, + ) -> MmResult<(), LockedNotesStorageError> { + let db = self.lockdb().await?; + let address = self.address.clone(); + let transaction = db.get_inner().transaction().await?; + let change_note_table = transaction.table::().await?; + + let change_note = LockedNoteTable { + address, + variant: "Change".to_owned(), + txid, + rseed: None, + value: Some(value), + }; + Ok(change_note_table + .add_item(&change_note) + .await + .map(|_| ())?) + } + + pub(crate) async fn remove_notes_for_txid(&self, txid: String) -> MmResult<(), LockedNotesStorageError> { + let db = self.lockdb().await?; + let transaction = db.get_inner().transaction().await?; + let change_note_table = transaction.table::().await?; + change_note_table.delete_items_by_index("txid", &txid).await?; + + Ok(()) + } + + pub(crate) async fn load_all_notes(&self) -> MmResult, LockedNotesStorageError> { + let db = self.lockdb().await?; + let transaction = db.get_inner().transaction().await?; + let change_note_table = transaction.table::().await?; + let records = change_note_table.get_items("address", &self.address).await?; + Ok(records + .into_iter() + .filter_map(|(_, n)| { + match n.variant.as_str() { + "Spent" => n.rseed.clone().map(|rseed| LockedNote::Spent { + rseed, + }), + "Change" => n.value.map(|value| LockedNote::Change { + value, + }), + _ => None, + } + }) + .collect()) + } +} diff --git a/mm2src/coins/z_coin/z_coin_errors.rs b/mm2src/coins/z_coin/z_coin_errors.rs index 3c6e44a32a..23b1664a5b 100644 --- a/mm2src/coins/z_coin/z_coin_errors.rs +++ b/mm2src/coins/z_coin/z_coin_errors.rs @@ -1,3 +1,4 @@ +use super::storage::LockedNotesStorageError; use crate::my_tx_history_v2::MyTxHistoryErrorV2; use crate::utxo::rpc_clients::UtxoRpcError; use crate::utxo::utxo_builder::UtxoCoinBuildError; @@ -9,6 +10,7 @@ use common::jsonrpc_client::JsonRpcError; #[cfg(not(target_arch = "wasm32"))] use db_common::sqlite::rusqlite::Error as SqliteError; use derive_more::Display; +use enum_derives::EnumFromStringify; use http::uri::InvalidUri; #[cfg(target_arch = "wasm32")] use mm2_db::indexed_db::cursor_prelude::*; @@ -100,7 +102,7 @@ pub enum UrlIterError { ConnectionFailure(tonic::transport::Error), } -#[derive(Debug, Display)] +#[derive(Debug, Display, EnumFromStringify)] pub enum GenTxError { DecryptedOutputNotFound, GetWitnessErr(GetUnspentWitnessErr), @@ -130,6 +132,8 @@ pub enum GenTxError { FailedToCreateNote, SpendableNotesError(String), Internal(String), + #[from_stringify("LockedNotesStorageError")] + SaveLockedNotesError(String), } impl From for GenTxError { @@ -177,7 +181,8 @@ impl From for WithdrawError { | GenTxError::LightClientErr(_) | GenTxError::SpendableNotesError(_) | GenTxError::FailedToCreateNote - | GenTxError::Internal(_) => WithdrawError::InternalError(gen_tx.to_string()), + | GenTxError::Internal(_) + | GenTxError::SaveLockedNotesError(_) => WithdrawError::InternalError(gen_tx.to_string()), } } } @@ -231,10 +236,11 @@ impl From for GetUnspentWitnessErr { fn from(err: SqliteError) -> GetUnspentWitnessErr { GetUnspentWitnessErr::ZcashDBError(err.to_string()) } } -#[derive(Debug, Display)] +#[derive(Debug, Display, EnumFromStringify)] pub enum ZCoinBuildError { UtxoBuilderError(UtxoCoinBuildError), GetAddressError, + #[from_stringify("LockedNotesStorageError")] ZcashDBError(String), Rpc(UtxoRpcError), #[display(fmt = "Sapling cache DB does not exist at {}. Please download it.", path)] diff --git a/mm2src/coins/z_coin/z_rpc.rs b/mm2src/coins/z_coin/z_rpc.rs index aba2d0d55f..9b35081f7c 100644 --- a/mm2src/coins/z_coin/z_rpc.rs +++ b/mm2src/coins/z_coin/z_rpc.rs @@ -1,9 +1,10 @@ use super::{z_coin_errors::*, BlockDbImpl, CheckPointBlockInfo, WalletDbShared, ZCoinBuilder, ZcoinConsensusParams}; -use crate::utxo::rpc_clients::NO_TX_ERROR_CODE; use crate::utxo::utxo_builder::{UtxoCoinBuilderCommonOps, DAY_IN_SECONDS}; +use crate::z_coin::storage::z_locked_notes::LockedNotesStorage; use crate::z_coin::storage::{BlockProcessingMode, DataConnStmtCacheWrapper}; use crate::z_coin::SyncStartPoint; use crate::RpcCommonOps; + use async_trait::async_trait; use common::executor::Timer; use common::executor::{spawn_abortable, AbortOnDropHandle}; @@ -338,13 +339,11 @@ impl ZRpcOps for LightRpcClient { match client.get_transaction(request).await { Ok(_) => break, Err(e) => { - error!("Error on getting tx {}", tx_id); - if e.message().contains(NO_TX_ERROR_CODE) { - if attempts >= 3 { - return false; - } - attempts += 1; + error!("Error on getting tx {}: err: {}", tx_id, e.to_string()); + if attempts >= 5 { + return false; } + attempts += 1; Timer::sleep(30.).await; }, } @@ -475,16 +474,16 @@ impl ZRpcOps for NativeClient { async fn check_tx_existence(&self, tx_id: TxId) -> bool { let mut attempts = 0; loop { - match self.get_raw_transaction_bytes(&H256Json::from(tx_id.0)).compat().await { + let tx_hash = H256Json::from(tx_id.0).reversed(); + let tx = self.get_raw_transaction_bytes(&tx_hash).compat().await; + match tx { Ok(_) => break, Err(e) => { - error!("Error on getting tx {}", tx_id); - if e.to_string().contains(NO_TX_ERROR_CODE) { - if attempts >= 3 { - return false; - } - attempts += 1; + error!("Error on getting tx {}: err: {}", tx_id, e.to_string()); + if attempts >= 5 { + return false; } + attempts += 1; Timer::sleep(30.).await; }, } @@ -507,6 +506,7 @@ pub(super) async fn init_light_client<'a>( blocks_db: BlockDbImpl, sync_params: &Option, skip_sync_params: bool, + locked_notes_db: LockedNotesStorage, ) -> Result<(AsyncMutex, WalletDbShared), MmError> { let coin = builder.ticker.to_string(); let (sync_status_notifier, sync_watcher) = channel(1); @@ -568,6 +568,7 @@ pub(super) async fn init_light_client<'a>( scan_interval_ms: builder.z_coin_params.scan_interval_ms, first_sync_block: first_sync_block.clone(), streaming_manager: builder.ctx.event_stream_manager.clone(), + locked_notes_db, }; let abort_handle = spawn_abortable(light_wallet_db_sync_loop(sync_handle, Box::new(light_rpc_clients))); @@ -583,6 +584,7 @@ pub(super) async fn init_native_client<'a>( builder: &ZCoinBuilder<'a>, native_client: NativeClient, blocks_db: BlockDbImpl, + locked_notes_db: LockedNotesStorage, ) -> Result<(AsyncMutex, WalletDbShared), MmError> { let coin = builder.ticker.to_string(); let (sync_status_notifier, sync_watcher) = channel(1); @@ -614,6 +616,7 @@ pub(super) async fn init_native_client<'a>( scan_interval_ms: builder.z_coin_params.scan_interval_ms, first_sync_block: first_sync_block.clone(), streaming_manager: builder.ctx.event_stream_manager.clone(), + locked_notes_db, }; let abort_handle = spawn_abortable(light_wallet_db_sync_loop(sync_handle, Box::new(native_client))); @@ -700,6 +703,7 @@ pub struct SaplingSyncLoopHandle { current_block: BlockHeight, blocks_db: BlockDbImpl, wallet_db: WalletDbShared, + locked_notes_db: LockedNotesStorage, consensus_params: ZcoinConsensusParams, /// Notifies about sync status without stopping the loop, e.g. on coin activation sync_status_notifier: AsyncSender, @@ -800,6 +804,7 @@ impl SaplingSyncLoopHandle { BlockProcessingMode::Validate, wallet_ops.get_max_height_hash().await?, None, + &self.locked_notes_db, ) .await { @@ -844,6 +849,7 @@ impl SaplingSyncLoopHandle { BlockProcessingMode::Scan(scan, self.streaming_manager.clone()), None, Some(self.scan_blocks_per_iteration), + &self.locked_notes_db, ) .await?; @@ -914,7 +920,7 @@ async fn light_wallet_db_sync_loop(mut sync_handle: SaplingSyncLoopHandle, mut c let walletdb = &sync_handle.wallet_db; if let Ok(is_tx_imported) = walletdb.is_tx_imported(tx_id).await { if !is_tx_imported { - info!("Tx {} is not imported yet", tx_id); + error!("Tx {} is not imported yet", tx_id); Timer::sleep(10.).await; continue; } diff --git a/mm2src/mm2_main/Cargo.toml b/mm2src/mm2_main/Cargo.toml index 560156e148..b2fe520cfb 100644 --- a/mm2src/mm2_main/Cargo.toml +++ b/mm2src/mm2_main/Cargo.toml @@ -98,6 +98,7 @@ serialization_derive = { path = "../mm2_bitcoin/serialization_derive" } spv_validation = { path = "../mm2_bitcoin/spv_validation" } sp-runtime-interface.workspace = true sp-trie.workspace = true +tempfile.workspace = true trie-db.workspace = true trie-root.workspace = true uuid.workspace = true diff --git a/mm2src/mm2_main/tests/docker_tests/docker_tests_common.rs b/mm2src/mm2_main/tests/docker_tests/docker_tests_common.rs index d513ada2d2..6a483767af 100644 --- a/mm2src/mm2_main/tests/docker_tests/docker_tests_common.rs +++ b/mm2src/mm2_main/tests/docker_tests/docker_tests_common.rs @@ -239,7 +239,7 @@ impl CoinDockerOps for ZCoinAssetDockerOps { impl ZCoinAssetDockerOps { pub fn new() -> ZCoinAssetDockerOps { - let (ctx, coin) = block_on(z_coin_from_spending_key("secret-extended-key-main1q0k2ga2cqqqqpq8m8j6yl0say83cagrqp53zqz54w38ezs8ly9ly5ptamqwfpq85u87w0df4k8t2lwyde3n9v0gcr69nu4ryv60t0kfcsvkr8h83skwqex2nf0vr32794fmzk89cpmjptzc22lgu5wfhhp8lgf3f5vn2l3sge0udvxnm95k6dtxj2jwlfyccnum7nz297ecyhmd5ph526pxndww0rqq0qly84l635mec0x4yedf95hzn6kcgq8yxts26k98j9g32kjc8y83fe")); + let (ctx, coin) = block_on(z_coin_from_spending_key("secret-extended-key-main1q0k2ga2cqqqqpq8m8j6yl0say83cagrqp53zqz54w38ezs8ly9ly5ptamqwfpq85u87w0df4k8t2lwyde3n9v0gcr69nu4ryv60t0kfcsvkr8h83skwqex2nf0vr32794fmzk89cpmjptzc22lgu5wfhhp8lgf3f5vn2l3sge0udvxnm95k6dtxj2jwlfyccnum7nz297ecyhmd5ph526pxndww0rqq0qly84l635mec0x4yedf95hzn6kcgq8yxts26k98j9g32kjc8y83fe", "fe")); ZCoinAssetDockerOps { ctx, coin } } diff --git a/mm2src/mm2_main/tests/docker_tests/z_coin_docker_tests.rs b/mm2src/mm2_main/tests/docker_tests/z_coin_docker_tests.rs index f9cd8e99af..1fa139ab31 100644 --- a/mm2src/mm2_main/tests/docker_tests/z_coin_docker_tests.rs +++ b/mm2src/mm2_main/tests/docker_tests/z_coin_docker_tests.rs @@ -1,6 +1,7 @@ use bitcrypto::dhash160; use coins::z_coin::{z_coin_from_conf_and_params_with_docker, z_send_dex_fee, ZCoin, ZcoinActivationParams, ZcoinRpcMode}; +use coins::DexFeeBurnDestination; use coins::{coin_errors::ValidatePaymentError, CoinProtocol, DexFee, PrivKeyBuildPolicy, RefundPaymentArgs, SendPaymentArgs, SpendPaymentArgs, SwapOps, SwapTxTypeWithSecretHash, ValidateFeeArgs}; use common::now_sec; @@ -8,22 +9,33 @@ use lazy_static::lazy_static; use mm2_core::mm_ctx::{MmArc, MmCtxBuilder}; use mm2_number::MmNumber; use mm2_test_helpers::for_tests::zombie_conf_for_docker; +use tempfile::TempDir; use tokio::sync::Mutex; // https://github.com/KomodoPlatform/librustzcash/blob/4e030a0f44cc17f100bf5f019563be25c5b8755f/zcash_client_backend/src/data_api/wallet.rs#L72-L73 lazy_static! { - static ref TEST_MUTEX: Mutex<()> = Mutex::new(()); + /// For secret....fe + static ref GEN_TX_LOCK_MUTEX: Mutex<()> = Mutex::new(()); + /// For secret....we + static ref GEN_TX_LOCK_MUTEX_ADDR2: Mutex<()> = Mutex::new(()); + /// This `TempDir` is created once on first use and cleaned up when the process exits. + static ref TEMP_DIR: Mutex = Mutex::new(TempDir::new().unwrap()); } -/// Build asset `ZCoin` from ticker and spendingkey str without filling the balance. -pub async fn z_coin_from_spending_key(spending_key: &str) -> (MmArc, ZCoin) { - let ctx = MmCtxBuilder::new().into_mm_arc(); +/// Build asset `ZCoin` from ticker and spending_key. +pub async fn z_coin_from_spending_key<'a>(spending_key: &str, path: &'a str) -> (MmArc, ZCoin) { + let tmp = TEMP_DIR.lock().await; + let db_path = tmp.path().join(format!("ZOMBIE_DB_{path}")); + std::fs::create_dir_all(&db_path).unwrap(); + let ctx = MmCtxBuilder::new().with_conf(json!({ "dbdir": db_path})).into_mm_arc(); + let mut conf = zombie_conf_for_docker(); let params = ZcoinActivationParams { mode: ZcoinRpcMode::Native, ..Default::default() }; let pk_data = [1; 32]; + let protocol_info = match serde_json::from_value::(conf["protocol"].take()).unwrap() { CoinProtocol::ZHTLC(protocol_info) => protocol_info, other_protocol => panic!("Failed to get protocol from config: {:?}", other_protocol), @@ -44,17 +56,24 @@ pub async fn z_coin_from_spending_key(spending_key: &str) -> (MmArc, ZCoin) { (ctx, coin) } -#[ignore] -#[tokio::test(flavor = "multi_thread")] +#[tokio::test(flavor = "current_thread")] +async fn prepare_zombie_sapling_cache() { + let _lock = GEN_TX_LOCK_MUTEX.lock().await; + let (_ctx, coin) = z_coin_from_spending_key("secret-extended-key-main1q0k2ga2cqqqqpq8m8j6yl0say83cagrqp53zqz54w38ezs8ly9ly5ptamqwfpq85u87w0df4k8t2lwyde3n9v0gcr69nu4ryv60t0kfcsvkr8h83skwqex2nf0vr32794fmzk89cpmjptzc22lgu5wfhhp8lgf3f5vn2l3sge0udvxnm95k6dtxj2jwlfyccnum7nz297ecyhmd5ph526pxndww0rqq0qly84l635mec0x4yedf95hzn6kcgq8yxts26k98j9g32kjc8y83fe", "fe").await; + assert!(coin.is_sapling_state_synced().await); + drop(_lock) +} + +#[tokio::test(flavor = "current_thread")] async fn zombie_coin_send_and_refund_maker_payment() { - let _lock = TEST_MUTEX.lock().await; + let _lock = GEN_TX_LOCK_MUTEX.lock().await; + let (_ctx, coin) = z_coin_from_spending_key("secret-extended-key-main1q0k2ga2cqqqqpq8m8j6yl0say83cagrqp53zqz54w38ezs8ly9ly5ptamqwfpq85u87w0df4k8t2lwyde3n9v0gcr69nu4ryv60t0kfcsvkr8h83skwqex2nf0vr32794fmzk89cpmjptzc22lgu5wfhhp8lgf3f5vn2l3sge0udvxnm95k6dtxj2jwlfyccnum7nz297ecyhmd5ph526pxndww0rqq0qly84l635mec0x4yedf95hzn6kcgq8yxts26k98j9g32kjc8y83fe", "fe").await; + + assert!(coin.is_sapling_state_synced().await); - let (_ctx, coin) = z_coin_from_spending_key("secret-extended-key-main1q0k2ga2cqqqqpq8m8j6yl0say83cagrqp53zqz54w38ezs8ly9ly5ptamqwfpq85u87w0df4k8t2lwyde3n9v0gcr69nu4ryv60t0kfcsvkr8h83skwqex2nf0vr32794fmzk89cpmjptzc22lgu5wfhhp8lgf3f5vn2l3sge0udvxnm95k6dtxj2jwlfyccnum7nz297ecyhmd5ph526pxndww0rqq0qly84l635mec0x4yedf95hzn6kcgq8yxts26k98j9g32kjc8y83fe").await; let time_lock = now_sec() - 3600; let secret_hash = [0; 20]; - let maker_uniq_data = [3; 32]; - let taker_uniq_data = [5; 32]; let taker_key_pair = coin.derive_htlc_key_pair(taker_uniq_data.as_slice()); let taker_pub = taker_key_pair.public(); @@ -90,12 +109,12 @@ async fn zombie_coin_send_and_refund_maker_payment() { drop(_lock); } -#[ignore] -#[tokio::test(flavor = "multi_thread")] +#[tokio::test(flavor = "current_thread")] async fn zombie_coin_send_and_spend_maker_payment() { - let _lock = TEST_MUTEX.lock().await; + let _lock = GEN_TX_LOCK_MUTEX_ADDR2.lock().await; + let (_ctx, coin) = z_coin_from_spending_key("secret-extended-key-main1qvqstxphqyqqpqqnh3hstqpdjzkpadeed6u7fz230jmm2mxl0aacrtu9vt7a7rmr2w5az5u79d24t0rudak3newknrz5l0m3dsd8m4dffqh5xwyldc5qwz8pnalrnhlxdzf900x83jazc52y25e9hvyd4kepaze6nlcvk8sd8a4qjh3e9j5d6730t7ctzhhrhp0zljjtwuptadnksxf8a8y5axwdhass5pjaxg0hzhg7z25rx0rll7a6txywl32s6cda0s5kexr03uqdtelwe", "we").await; - let (_ctx, coin) = z_coin_from_spending_key("secret-extended-key-main1q0k2ga2cqqqqpq8m8j6yl0say83cagrqp53zqz54w38ezs8ly9ly5ptamqwfpq85u87w0df4k8t2lwyde3n9v0gcr69nu4ryv60t0kfcsvkr8h83skwqex2nf0vr32794fmzk89cpmjptzc22lgu5wfhhp8lgf3f5vn2l3sge0udvxnm95k6dtxj2jwlfyccnum7nz297ecyhmd5ph526pxndww0rqq0qly84l635mec0x4yedf95hzn6kcgq8yxts26k98j9g32kjc8y83fe").await; + assert!(coin.is_sapling_state_synced().await); let lock_time = now_sec() - 1000; let secret = [0; 32]; @@ -139,37 +158,55 @@ async fn zombie_coin_send_and_spend_maker_payment() { drop(_lock); } -#[ignore] -#[tokio::test(flavor = "multi_thread")] -async fn prepare_zombie_sapling_cache() { - let (_ctx, coin) = z_coin_from_spending_key("secret-extended-key-main1q0k2ga2cqqqqpq8m8j6yl0say83cagrqp53zqz54w38ezs8ly9ly5ptamqwfpq85u87w0df4k8t2lwyde3n9v0gcr69nu4ryv60t0kfcsvkr8h83skwqex2nf0vr32794fmzk89cpmjptzc22lgu5wfhhp8lgf3f5vn2l3sge0udvxnm95k6dtxj2jwlfyccnum7nz297ecyhmd5ph526pxndww0rqq0qly84l635mec0x4yedf95hzn6kcgq8yxts26k98j9g32kjc8y83fe").await; +#[tokio::test(flavor = "current_thread")] +async fn zombie_coin_send_standard_dex_fee() { + let _lock = GEN_TX_LOCK_MUTEX_ADDR2.lock().await; + let (_ctx, coin) = z_coin_from_spending_key("secret-extended-key-main1qvqstxphqyqqpqqnh3hstqpdjzkpadeed6u7fz230jmm2mxl0aacrtu9vt7a7rmr2w5az5u79d24t0rudak3newknrz5l0m3dsd8m4dffqh5xwyldc5qwz8pnalrnhlxdzf900x83jazc52y25e9hvyd4kepaze6nlcvk8sd8a4qjh3e9j5d6730t7ctzhhrhp0zljjtwuptadnksxf8a8y5axwdhass5pjaxg0hzhg7z25rx0rll7a6txywl32s6cda0s5kexr03uqdtelwe", "we").await; assert!(coin.is_sapling_state_synced().await); + + let tx = z_send_dex_fee(&coin, DexFee::Standard("0.01".into()), &[1; 16]) + .await + .unwrap(); + log!("dex fee tx {}", tx.txid()); + drop(_lock) } -#[ignore] -#[tokio::test(flavor = "multi_thread")] +#[tokio::test(flavor = "current_thread")] async fn zombie_coin_send_dex_fee() { - let _lock = TEST_MUTEX.lock().await; + let _lock = GEN_TX_LOCK_MUTEX_ADDR2.lock().await; + let (_ctx, coin) = z_coin_from_spending_key("secret-extended-key-main1qvqstxphqyqqpqqnh3hstqpdjzkpadeed6u7fz230jmm2mxl0aacrtu9vt7a7rmr2w5az5u79d24t0rudak3newknrz5l0m3dsd8m4dffqh5xwyldc5qwz8pnalrnhlxdzf900x83jazc52y25e9hvyd4kepaze6nlcvk8sd8a4qjh3e9j5d6730t7ctzhhrhp0zljjtwuptadnksxf8a8y5axwdhass5pjaxg0hzhg7z25rx0rll7a6txywl32s6cda0s5kexr03uqdtelwe", "we").await; - let (_ctx, coin) = z_coin_from_spending_key("secret-extended-key-main1q0k2ga2cqqqqpq8m8j6yl0say83cagrqp53zqz54w38ezs8ly9ly5ptamqwfpq85u87w0df4k8t2lwyde3n9v0gcr69nu4ryv60t0kfcsvkr8h83skwqex2nf0vr32794fmzk89cpmjptzc22lgu5wfhhp8lgf3f5vn2l3sge0udvxnm95k6dtxj2jwlfyccnum7nz297ecyhmd5ph526pxndww0rqq0qly84l635mec0x4yedf95hzn6kcgq8yxts26k98j9g32kjc8y83fe").await; + assert!(coin.is_sapling_state_synced().await); - let dex_fee = DexFee::Standard("0.01".into()); + let dex_fee = DexFee::WithBurn { + fee_amount: "0.0075".into(), + burn_amount: "0.0025".into(), + burn_destination: DexFeeBurnDestination::PreBurnAccount, + }; let tx = z_send_dex_fee(&coin, dex_fee, &[1; 16]).await.unwrap(); log!("dex fee tx {}", tx.txid()); drop(_lock); } -#[ignore] -#[tokio::test(flavor = "multi_thread")] +#[tokio::test(flavor = "current_thread")] async fn zombie_coin_validate_dex_fee() { - let _lock = TEST_MUTEX.lock().await; - let (_ctx, coin) = z_coin_from_spending_key("secret-extended-key-main1q0k2ga2cqqqqpq8m8j6yl0say83cagrqp53zqz54w38ezs8ly9ly5ptamqwfpq85u87w0df4k8t2lwyde3n9v0gcr69nu4ryv60t0kfcsvkr8h83skwqex2nf0vr32794fmzk89cpmjptzc22lgu5wfhhp8lgf3f5vn2l3sge0udvxnm95k6dtxj2jwlfyccnum7nz297ecyhmd5ph526pxndww0rqq0qly84l635mec0x4yedf95hzn6kcgq8yxts26k98j9g32kjc8y83fe").await; + let _lock = GEN_TX_LOCK_MUTEX.lock().await; + let (_ctx, coin) = z_coin_from_spending_key("secret-extended-key-main1q0k2ga2cqqqqpq8m8j6yl0say83cagrqp53zqz54w38ezs8ly9ly5ptamqwfpq85u87w0df4k8t2lwyde3n9v0gcr69nu4ryv60t0kfcsvkr8h83skwqex2nf0vr32794fmzk89cpmjptzc22lgu5wfhhp8lgf3f5vn2l3sge0udvxnm95k6dtxj2jwlfyccnum7nz297ecyhmd5ph526pxndww0rqq0qly84l635mec0x4yedf95hzn6kcgq8yxts26k98j9g32kjc8y83fe", "fe").await; - // let balance = coin.my_balance().compat().await; + assert!(coin.is_sapling_state_synced().await); - let dex_fee = DexFee::Standard("0.01".into()); - let tx = z_send_dex_fee(&coin, dex_fee, &[1; 16]).await.unwrap(); + let tx = z_send_dex_fee( + &coin, + DexFee::WithBurn { + fee_amount: "0.0075".into(), + burn_amount: "0.0025".into(), + burn_destination: DexFeeBurnDestination::PreBurnAccount, + }, + &[1; 16], + ) + .await + .unwrap(); log!("dex fee tx {}", tx.txid()); let tx = tx.into(); @@ -177,52 +214,78 @@ async fn zombie_coin_validate_dex_fee() { fee_tx: &tx, expected_sender: &[], dex_fee: &DexFee::Standard(MmNumber::from("0.001")), - min_block_number: 4, + min_block_number: 12000, uuid: &[1; 16], }; // Invalid amount should return an error let err = coin.validate_fee(validate_fee_args).await.unwrap_err().into_inner(); match err { - ValidatePaymentError::WrongPaymentTx(err) => assert!(err.contains("Dex fee has invalid amount")), + ValidatePaymentError::WrongPaymentTx(err) => assert!(err.contains("invalid amount")), _ => panic!("Expected `WrongPaymentTx`: {:?}", err), } // Invalid memo should return an error + let expected_fee = DexFee::WithBurn { + fee_amount: "0.0075".into(), + burn_amount: "0.0025".into(), + burn_destination: DexFeeBurnDestination::PreBurnAccount, + }; + let validate_fee_args = ValidateFeeArgs { fee_tx: &tx, expected_sender: &[], - dex_fee: &DexFee::Standard(MmNumber::from("0.01")), - min_block_number: 10, + dex_fee: &expected_fee, + min_block_number: 12000, uuid: &[2; 16], }; + let err = coin.validate_fee(validate_fee_args).await.unwrap_err().into_inner(); match err { - ValidatePaymentError::WrongPaymentTx(err) => assert!(err.contains("Dex fee has invalid memo")), + ValidatePaymentError::WrongPaymentTx(err) => assert!(err.contains("invalid memo")), _ => panic!("Expected `WrongPaymentTx`: {:?}", err), } - // Confirmed before min block + // Success validation let validate_fee_args = ValidateFeeArgs { fee_tx: &tx, expected_sender: &[], - dex_fee: &DexFee::Standard(MmNumber::from("0.01")), - min_block_number: 20000, + dex_fee: &expected_fee, + min_block_number: 12000, + uuid: &[1; 16], + }; + coin.validate_fee(validate_fee_args).await.unwrap(); + + // Test old standard dex fee with no burn output + // TODO: disable when the upgrade transition period ends + let tx_2 = z_send_dex_fee(&coin, DexFee::Standard("0.00879999".into()), &[1; 16]) + .await + .unwrap(); + log!("dex fee tx {}", tx_2.txid()); + let tx_2 = tx_2.into(); + + // Success validation + let validate_fee_args = ValidateFeeArgs { + fee_tx: &tx_2, + expected_sender: &[], + dex_fee: &DexFee::Standard("0.00999999".into()), + min_block_number: 12000, uuid: &[1; 16], }; let err = coin.validate_fee(validate_fee_args).await.unwrap_err().into_inner(); match err { - ValidatePaymentError::WrongPaymentTx(err) => assert!(err.contains("confirmed before min block")), + ValidatePaymentError::WrongPaymentTx(err) => assert!(err.contains("invalid amount")), _ => panic!("Expected `WrongPaymentTx`: {:?}", err), } // Success validation + let expected_std_fee = DexFee::Standard("0.00879999".into()); let validate_fee_args = ValidateFeeArgs { - fee_tx: &tx, + fee_tx: &tx_2, expected_sender: &[], - dex_fee: &DexFee::Standard(MmNumber::from("0.01")), - min_block_number: 20, + dex_fee: &expected_std_fee, + min_block_number: 12000, uuid: &[1; 16], }; coin.validate_fee(validate_fee_args).await.unwrap(); - drop(_lock); + drop(_lock) } diff --git a/mm2src/mm2_test_helpers/src/for_tests.rs b/mm2src/mm2_test_helpers/src/for_tests.rs index 5d7c1dad6b..c8f989206c 100644 --- a/mm2src/mm2_test_helpers/src/for_tests.rs +++ b/mm2src/mm2_test_helpers/src/for_tests.rs @@ -478,11 +478,11 @@ pub enum Mm2InitPrivKeyPolicy { GlobalHDAccount, } -pub fn zombie_conf() -> Json { zombie_conf_inner(None) } +pub fn zombie_conf() -> Json { zombie_conf_inner(None, 0) } -pub fn zombie_conf_for_docker() -> Json { zombie_conf_inner(Some(10)) } +pub fn zombie_conf_for_docker() -> Json { zombie_conf_inner(Some(10), 1) } -pub fn zombie_conf_inner(custom_blocktime: Option) -> Json { +pub fn zombie_conf_inner(custom_blocktime: Option, required_confirmations: u8) -> Json { json!({ "coin":"ZOMBIE", "asset":"ZOMBIE", @@ -509,7 +509,7 @@ pub fn zombie_conf_inner(custom_blocktime: Option) -> Json { "z_derivation_path": "m/32'/133'", } }, - "required_confirmations":0, + "required_confirmations": required_confirmations, "derivation_path": "m/44'/133'", }) }