From 4c3b31cf0ca45dc662b5bc89fdb2a25910201ce8 Mon Sep 17 00:00:00 2001 From: Elias Rohrer Date: Mon, 5 Oct 2026 10:21:00 +0200 Subject: [PATCH] Date mempool absences with our own clock The bitcoind chain source dated a transaction's absence from the mempool with the newest mempool entry time it had seen since startup. That watermark restarts at zero while the transaction's `last_seen` persists with the wallet, and BDK keeps an evicted transaction canonical until its eviction catches up with `last_seen`, so a transaction that vanished while we were down kept its input marked spent and its change counted as ours. Date the observations we hand BDK with our local clock instead, as the Esplora and Electrum sources and upstream `bdk_bitcoind_rpc` do, and keep the entry-time watermark for emission deduplication only. This only bites after a restart, and only while the mempool holds nothing newer than the transaction: otherwise the first poll re-emits the whole mempool and advances the watermark before that same poll reports any eviction. In practice that means an empty mempool on signet or regtest, or a similarly idle backend, and the next transaction to arrive resolves it anyway. Co-Authored-By: HAL 9000 --- src/chain/bitcoind.rs | 66 ++++++++++++++++++++------- src/wallet/mod.rs | 103 ++++++++++++++++++++++++++++++++++++++++++ 2 files changed, 152 insertions(+), 17 deletions(-) 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)); + } }