diff --git a/src/chain/bitcoind.rs b/src/chain/bitcoind.rs index 9ed38c2128..213f4a414d 100644 --- a/src/chain/bitcoind.rs +++ b/src/chain/bitcoind.rs @@ -791,14 +791,14 @@ pub enum BitcoindClient { rpc_client: Arc, latest_mempool_timestamp: AtomicU64, mempool_entries_cache: tokio::sync::Mutex>, - mempool_txs_cache: tokio::sync::Mutex>, + mempool_txs_cache: tokio::sync::Mutex>, }, Rest { rest_client: Arc, rpc_client: Arc, latest_mempool_timestamp: AtomicU64, mempool_entries_cache: tokio::sync::Mutex>, - mempool_txs_cache: tokio::sync::Mutex>, + mempool_txs_cache: tokio::sync::Mutex>, }, } @@ -1238,9 +1238,12 @@ impl BitcoindClient { async fn get_mempool_transactions_and_timestamp_at_height_inner( &self, latest_mempool_timestamp: &AtomicU64, mempool_entries_cache: &tokio::sync::Mutex>, - mempool_txs_cache: &tokio::sync::Mutex>, + mempool_txs_cache: &tokio::sync::Mutex>, best_processed_height: u32, ) -> Result, BitcoindClientError> { + // Date `last_seen` on our own clock: the entry-time watermark below resets on restart. + let observed_at = + SystemTime::now().duration_since(UNIX_EPOCH).map(|d| d.as_secs()).unwrap_or(0); let prev_mempool_time = latest_mempool_timestamp.load(Ordering::Relaxed); let mut latest_time = prev_mempool_time; @@ -1273,15 +1276,15 @@ impl BitcoindClient { continue; } - if let Some((cached_tx, cached_time)) = mempool_txs_cache.get(txid) { - txs_to_emit.push((cached_tx.clone(), *cached_time)); + if let Some(cached_tx) = mempool_txs_cache.get(txid) { + txs_to_emit.push((cached_tx.clone(), observed_at)); continue; } match self.get_raw_transaction(&entry.txid).await { Ok(Some(tx)) => { - mempool_txs_cache.insert(entry.txid, (tx.clone(), entry.time)); - txs_to_emit.push((tx, entry.time)); + mempool_txs_cache.insert(entry.txid, tx.clone()); + txs_to_emit.push((tx, observed_at)); }, Ok(None) => { continue; @@ -1304,17 +1307,15 @@ impl BitcoindClient { &self, bdk_unconfirmed_txids: Vec, ) -> Result, BitcoindClientError> { match self { - BitcoindClient::Rpc { latest_mempool_timestamp, mempool_entries_cache, .. } => { + BitcoindClient::Rpc { mempool_entries_cache, .. } => { Self::get_evicted_mempool_txids_and_timestamp_inner( - latest_mempool_timestamp, mempool_entries_cache, bdk_unconfirmed_txids, ) .await }, - BitcoindClient::Rest { latest_mempool_timestamp, mempool_entries_cache, .. } => { + BitcoindClient::Rest { mempool_entries_cache, .. } => { Self::get_evicted_mempool_txids_and_timestamp_inner( - latest_mempool_timestamp, mempool_entries_cache, bdk_unconfirmed_txids, ) @@ -1324,16 +1325,17 @@ impl BitcoindClient { } async fn get_evicted_mempool_txids_and_timestamp_inner( - latest_mempool_timestamp: &AtomicU64, mempool_entries_cache: &tokio::sync::Mutex>, bdk_unconfirmed_txids: Vec, ) -> Result, BitcoindClientError> { - let latest_mempool_timestamp = latest_mempool_timestamp.load(Ordering::Relaxed); + // BDK only drops the transaction once this reaches its persisted `last_seen`. + let observed_at = + SystemTime::now().duration_since(UNIX_EPOCH).map(|d| d.as_secs()).unwrap_or(0); let mempool_entries_cache = mempool_entries_cache.lock().await; let evicted_txids = bdk_unconfirmed_txids .into_iter() .filter(|txid| !mempool_entries_cache.contains_key(txid)) - .map(|txid| (txid, latest_mempool_timestamp)) + .map(|txid| (txid, observed_at)) .collect(); Ok(evicted_txids) } @@ -1618,9 +1620,9 @@ impl std::error::Error for BitcoindClientError {} #[cfg(test)] mod tests { - use std::collections::HashSet; + use std::collections::{HashMap, HashSet}; use std::sync::Mutex; - use std::time::Duration; + use std::time::{Duration, SystemTime, UNIX_EPOCH}; use bitcoin::hashes::Hash; use bitcoin::{FeeRate, OutPoint, ScriptBuf, Transaction, TxIn, TxOut, Txid, Witness}; @@ -1631,12 +1633,42 @@ mod tests { use serde_json::json; use crate::chain::bitcoind::{ - acquire_initial_wallet_sync_guard, FeeResponse, GetMempoolEntryResponse, + acquire_initial_wallet_sync_guard, BitcoindClient, FeeResponse, GetMempoolEntryResponse, GetRawMempoolResponse, GetRawTransactionResponse, MempoolMinFeeResponse, }; use crate::chain::{WalletSyncGuard, WalletSyncStatus}; use crate::Error; + /// An absence must be dated with when we observed it, not with the newest mempool entry + /// time we have seen. That watermark resets to zero on restart while the transaction's + /// `last_seen` persists, and BDK keeps an evicted transaction canonical while its + /// `evicted_at` predates its `last_seen`. + #[tokio::test] + async fn mempool_absence_is_dated_with_the_observation_time() { + let before = SystemTime::now().duration_since(UNIX_EPOCH).map(|d| d.as_secs()).unwrap_or(0); + + // An empty mempool: nothing here could advance an entry-time watermark past the + // transaction's `last_seen`. + let mempool_entries_cache = tokio::sync::Mutex::new(HashMap::new()); + let txid = Txid::from_byte_array([23u8; 32]); + + let evicted = BitcoindClient::get_evicted_mempool_txids_and_timestamp_inner( + &mempool_entries_cache, + vec![txid], + ) + .await + .unwrap(); + + assert_eq!(evicted.len(), 1); + assert_eq!(evicted[0].0, txid); + assert!( + evicted[0].1 >= before, + "an absence must be dated with when we observed it, got {} before {}", + evicted[0].1, + before, + ); + } + #[tokio::test] async fn initial_sync_waits_for_in_progress_sync() { let status = Mutex::new(WalletSyncStatus::Completed); diff --git a/src/wallet/mod.rs b/src/wallet/mod.rs index 13a8ef4e00..ef51ee9afb 100644 --- a/src/wallet/mod.rs +++ b/src/wallet/mod.rs @@ -4414,4 +4414,107 @@ mod tests { ); assert_ne!(locked_wallet.next_unused_address(KeychainKind::Internal).index, 0); } + + /// Pins down the eviction ordering our bitcoind mempool producer relies on: BDK keeps an + /// evicted transaction canonical for as long as the absence we reported predates the + /// `last_seen` we reported, and only drops it once the absence catches up. The producer has + /// to date absences on a clock that can clear a `last_seen` reloaded from the store, which + /// is why it reads the local clock rather than Bitcoin Core's mempool entry times — see + /// `chain::bitcoind::observation_time_secs`. + #[cfg(feature = "chain-bitcoind")] + #[tokio::test] + async fn mempool_eviction_needs_an_absence_dated_after_last_seen() { + // The mempool entry time our outgoing transaction was emitted with, i.e. the + // `last_seen` the wallet persists for it. + const LAST_SEEN: u64 = 200; + + // A confirmed 100_000 sat input, spent by an unconfirmed 60_000 sat payment that paid a + // 1_000 sat fee and left 39_000 sats of change. + let store: Arc = Arc::new(DynStoreWrapper(InMemoryStore::new())); + let spend_txid = { + let wallet = new_test_wallet(Arc::clone(&store), false).await; + + let (funding_tx, block_id) = { + let mut locked_wallet = wallet.inner.lock().unwrap(); + let script_pubkey = locked_wallet + .reveal_next_address(KeychainKind::External) + .address + .script_pubkey(); + let funding_tx = Transaction { + version: bitcoin::transaction::Version::TWO, + lock_time: LockTime::ZERO, + input: Vec::new(), + output: vec![TxOut { value: Amount::from_sat(100_000), script_pubkey }], + }; + let block_id = BlockId { + height: locked_wallet.latest_checkpoint().height() + 1, + hash: bitcoin::BlockHash::from_byte_array([23; 32]), + }; + (funding_tx, block_id) + }; + let funding_txid = funding_tx.compute_txid(); + let mut tx_update = TxUpdate::default(); + tx_update.txs = vec![Arc::new(funding_tx)]; + tx_update.anchors = + [(ConfirmationBlockTime { block_id, confirmation_time: 1 }, funding_txid)].into(); + let chain = CheckPoint::from_block_ids([ + wallet.inner.lock().unwrap().latest_checkpoint().block_id(), + block_id, + ]) + .unwrap(); + wallet + .apply_update(Update { tx_update, chain: Some(chain), ..Default::default() }) + .await + .unwrap(); + assert_eq!(wallet.get_balances(0).unwrap(), (100_000, 100_000)); + + let change_script_pubkey = wallet + .inner + .lock() + .unwrap() + .reveal_next_address(KeychainKind::Internal) + .address + .script_pubkey(); + let spend_tx = Transaction { + version: bitcoin::transaction::Version::TWO, + lock_time: LockTime::ZERO, + input: vec![bitcoin::TxIn { + previous_output: OutPoint { txid: funding_txid, vout: 0 }, + ..Default::default() + }], + output: vec![ + TxOut { + value: Amount::from_sat(60_000), + script_pubkey: ScriptBuf::from_bytes(vec![0x51]), + }, + TxOut { value: Amount::from_sat(39_000), script_pubkey: change_script_pubkey }, + ], + }; + let spend_txid = spend_tx.compute_txid(); + wallet.apply_mempool_txs(vec![(spend_tx, LAST_SEEN)], Vec::new()).await.unwrap(); + assert_eq!(wallet.get_balances(0).unwrap(), (39_000, 39_000)); + + spend_txid + }; + + // Restart: the wallet reloads the transaction and its `last_seen` from the store. + let wallet = new_test_wallet(Arc::clone(&store), true).await; + assert_eq!(wallet.get_unconfirmed_txids(), vec![spend_txid]); + assert_eq!(wallet.get_balances(0).unwrap(), (39_000, 39_000)); + + // The transaction is gone from the mempool, with nothing confirming, replacing or + // descending from it. Absences predating `last_seen` leave it canonical: the input + // stays spent and the stale change keeps counting. + for stale in [0, LAST_SEEN - 1] { + wallet.apply_mempool_txs(Vec::new(), vec![(spend_txid, stale)]).await.unwrap(); + assert_eq!(wallet.get_unconfirmed_txids(), vec![spend_txid]); + assert_eq!(wallet.get_balances(0).unwrap(), (39_000, 39_000)); + } + + // Once the absence reaches `last_seen`, the transaction stops being canonical: the + // confirmed input comes back and the stale change stops counting. + wallet.apply_mempool_txs(Vec::new(), vec![(spend_txid, LAST_SEEN)]).await.unwrap(); + assert!(wallet.get_unconfirmed_txids().is_empty()); + assert_eq!(wallet.get_balances(0).unwrap(), (100_000, 100_000)); + } }