Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
66 changes: 49 additions & 17 deletions src/chain/bitcoind.rs
Original file line number Diff line number Diff line change
Expand Up @@ -791,14 +791,14 @@ pub enum BitcoindClient {
rpc_client: Arc<RpcClient>,
latest_mempool_timestamp: AtomicU64,
mempool_entries_cache: tokio::sync::Mutex<HashMap<Txid, MempoolEntry>>,
mempool_txs_cache: tokio::sync::Mutex<HashMap<Txid, (Transaction, u64)>>,
mempool_txs_cache: tokio::sync::Mutex<HashMap<Txid, Transaction>>,
},
Rest {
rest_client: Arc<RestClient>,
rpc_client: Arc<RpcClient>,
latest_mempool_timestamp: AtomicU64,
mempool_entries_cache: tokio::sync::Mutex<HashMap<Txid, MempoolEntry>>,
mempool_txs_cache: tokio::sync::Mutex<HashMap<Txid, (Transaction, u64)>>,
mempool_txs_cache: tokio::sync::Mutex<HashMap<Txid, Transaction>>,
},
}

Expand Down Expand Up @@ -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<HashMap<Txid, MempoolEntry>>,
mempool_txs_cache: &tokio::sync::Mutex<HashMap<Txid, (Transaction, u64)>>,
mempool_txs_cache: &tokio::sync::Mutex<HashMap<Txid, Transaction>>,
best_processed_height: u32,
) -> Result<Vec<(Transaction, u64)>, 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;

Expand Down Expand Up @@ -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;
Expand All @@ -1304,17 +1307,15 @@ impl BitcoindClient {
&self, bdk_unconfirmed_txids: Vec<Txid>,
) -> Result<Vec<(Txid, u64)>, 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,
)
Expand All @@ -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<HashMap<Txid, MempoolEntry>>,
bdk_unconfirmed_txids: Vec<Txid>,
) -> Result<Vec<(Txid, u64)>, 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)
}
Expand Down Expand Up @@ -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};
Expand All @@ -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);
Expand Down
103 changes: 103 additions & 0 deletions src/wallet/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<DynStore> = 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));
}
}
Loading