Skip to content

Commit 4c3b31c

Browse files
committed
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
1 parent 144a158 commit 4c3b31c

2 files changed

Lines changed: 152 additions & 17 deletions

File tree

‎src/chain/bitcoind.rs‎

Lines changed: 49 additions & 17 deletions
Original file line numberDiff line numberDiff line change
@@ -791,14 +791,14 @@ pub enum BitcoindClient {
791791
rpc_client: Arc<RpcClient>,
792792
latest_mempool_timestamp: AtomicU64,
793793
mempool_entries_cache: tokio::sync::Mutex<HashMap<Txid, MempoolEntry>>,
794-
mempool_txs_cache: tokio::sync::Mutex<HashMap<Txid, (Transaction, u64)>>,
794+
mempool_txs_cache: tokio::sync::Mutex<HashMap<Txid, Transaction>>,
795795
},
796796
Rest {
797797
rest_client: Arc<RestClient>,
798798
rpc_client: Arc<RpcClient>,
799799
latest_mempool_timestamp: AtomicU64,
800800
mempool_entries_cache: tokio::sync::Mutex<HashMap<Txid, MempoolEntry>>,
801-
mempool_txs_cache: tokio::sync::Mutex<HashMap<Txid, (Transaction, u64)>>,
801+
mempool_txs_cache: tokio::sync::Mutex<HashMap<Txid, Transaction>>,
802802
},
803803
}
804804

@@ -1238,9 +1238,12 @@ impl BitcoindClient {
12381238
async fn get_mempool_transactions_and_timestamp_at_height_inner(
12391239
&self, latest_mempool_timestamp: &AtomicU64,
12401240
mempool_entries_cache: &tokio::sync::Mutex<HashMap<Txid, MempoolEntry>>,
1241-
mempool_txs_cache: &tokio::sync::Mutex<HashMap<Txid, (Transaction, u64)>>,
1241+
mempool_txs_cache: &tokio::sync::Mutex<HashMap<Txid, Transaction>>,
12421242
best_processed_height: u32,
12431243
) -> Result<Vec<(Transaction, u64)>, BitcoindClientError> {
1244+
// Date `last_seen` on our own clock: the entry-time watermark below resets on restart.
1245+
let observed_at =
1246+
SystemTime::now().duration_since(UNIX_EPOCH).map(|d| d.as_secs()).unwrap_or(0);
12441247
let prev_mempool_time = latest_mempool_timestamp.load(Ordering::Relaxed);
12451248
let mut latest_time = prev_mempool_time;
12461249

@@ -1273,15 +1276,15 @@ impl BitcoindClient {
12731276
continue;
12741277
}
12751278

1276-
if let Some((cached_tx, cached_time)) = mempool_txs_cache.get(txid) {
1277-
txs_to_emit.push((cached_tx.clone(), *cached_time));
1279+
if let Some(cached_tx) = mempool_txs_cache.get(txid) {
1280+
txs_to_emit.push((cached_tx.clone(), observed_at));
12781281
continue;
12791282
}
12801283

12811284
match self.get_raw_transaction(&entry.txid).await {
12821285
Ok(Some(tx)) => {
1283-
mempool_txs_cache.insert(entry.txid, (tx.clone(), entry.time));
1284-
txs_to_emit.push((tx, entry.time));
1286+
mempool_txs_cache.insert(entry.txid, tx.clone());
1287+
txs_to_emit.push((tx, observed_at));
12851288
},
12861289
Ok(None) => {
12871290
continue;
@@ -1304,17 +1307,15 @@ impl BitcoindClient {
13041307
&self, bdk_unconfirmed_txids: Vec<Txid>,
13051308
) -> Result<Vec<(Txid, u64)>, BitcoindClientError> {
13061309
match self {
1307-
BitcoindClient::Rpc { latest_mempool_timestamp, mempool_entries_cache, .. } => {
1310+
BitcoindClient::Rpc { mempool_entries_cache, .. } => {
13081311
Self::get_evicted_mempool_txids_and_timestamp_inner(
1309-
latest_mempool_timestamp,
13101312
mempool_entries_cache,
13111313
bdk_unconfirmed_txids,
13121314
)
13131315
.await
13141316
},
1315-
BitcoindClient::Rest { latest_mempool_timestamp, mempool_entries_cache, .. } => {
1317+
BitcoindClient::Rest { mempool_entries_cache, .. } => {
13161318
Self::get_evicted_mempool_txids_and_timestamp_inner(
1317-
latest_mempool_timestamp,
13181319
mempool_entries_cache,
13191320
bdk_unconfirmed_txids,
13201321
)
@@ -1324,16 +1325,17 @@ impl BitcoindClient {
13241325
}
13251326

13261327
async fn get_evicted_mempool_txids_and_timestamp_inner(
1327-
latest_mempool_timestamp: &AtomicU64,
13281328
mempool_entries_cache: &tokio::sync::Mutex<HashMap<Txid, MempoolEntry>>,
13291329
bdk_unconfirmed_txids: Vec<Txid>,
13301330
) -> Result<Vec<(Txid, u64)>, BitcoindClientError> {
1331-
let latest_mempool_timestamp = latest_mempool_timestamp.load(Ordering::Relaxed);
1331+
// BDK only drops the transaction once this reaches its persisted `last_seen`.
1332+
let observed_at =
1333+
SystemTime::now().duration_since(UNIX_EPOCH).map(|d| d.as_secs()).unwrap_or(0);
13321334
let mempool_entries_cache = mempool_entries_cache.lock().await;
13331335
let evicted_txids = bdk_unconfirmed_txids
13341336
.into_iter()
13351337
.filter(|txid| !mempool_entries_cache.contains_key(txid))
1336-
.map(|txid| (txid, latest_mempool_timestamp))
1338+
.map(|txid| (txid, observed_at))
13371339
.collect();
13381340
Ok(evicted_txids)
13391341
}
@@ -1618,9 +1620,9 @@ impl std::error::Error for BitcoindClientError {}
16181620

16191621
#[cfg(test)]
16201622
mod tests {
1621-
use std::collections::HashSet;
1623+
use std::collections::{HashMap, HashSet};
16221624
use std::sync::Mutex;
1623-
use std::time::Duration;
1625+
use std::time::{Duration, SystemTime, UNIX_EPOCH};
16241626

16251627
use bitcoin::hashes::Hash;
16261628
use bitcoin::{FeeRate, OutPoint, ScriptBuf, Transaction, TxIn, TxOut, Txid, Witness};
@@ -1631,12 +1633,42 @@ mod tests {
16311633
use serde_json::json;
16321634

16331635
use crate::chain::bitcoind::{
1634-
acquire_initial_wallet_sync_guard, FeeResponse, GetMempoolEntryResponse,
1636+
acquire_initial_wallet_sync_guard, BitcoindClient, FeeResponse, GetMempoolEntryResponse,
16351637
GetRawMempoolResponse, GetRawTransactionResponse, MempoolMinFeeResponse,
16361638
};
16371639
use crate::chain::{WalletSyncGuard, WalletSyncStatus};
16381640
use crate::Error;
16391641

1642+
/// An absence must be dated with when we observed it, not with the newest mempool entry
1643+
/// time we have seen. That watermark resets to zero on restart while the transaction's
1644+
/// `last_seen` persists, and BDK keeps an evicted transaction canonical while its
1645+
/// `evicted_at` predates its `last_seen`.
1646+
#[tokio::test]
1647+
async fn mempool_absence_is_dated_with_the_observation_time() {
1648+
let before = SystemTime::now().duration_since(UNIX_EPOCH).map(|d| d.as_secs()).unwrap_or(0);
1649+
1650+
// An empty mempool: nothing here could advance an entry-time watermark past the
1651+
// transaction's `last_seen`.
1652+
let mempool_entries_cache = tokio::sync::Mutex::new(HashMap::new());
1653+
let txid = Txid::from_byte_array([23u8; 32]);
1654+
1655+
let evicted = BitcoindClient::get_evicted_mempool_txids_and_timestamp_inner(
1656+
&mempool_entries_cache,
1657+
vec![txid],
1658+
)
1659+
.await
1660+
.unwrap();
1661+
1662+
assert_eq!(evicted.len(), 1);
1663+
assert_eq!(evicted[0].0, txid);
1664+
assert!(
1665+
evicted[0].1 >= before,
1666+
"an absence must be dated with when we observed it, got {} before {}",
1667+
evicted[0].1,
1668+
before,
1669+
);
1670+
}
1671+
16401672
#[tokio::test]
16411673
async fn initial_sync_waits_for_in_progress_sync() {
16421674
let status = Mutex::new(WalletSyncStatus::Completed);

‎src/wallet/mod.rs‎

Lines changed: 103 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -4414,4 +4414,107 @@ mod tests {
44144414
);
44154415
assert_ne!(locked_wallet.next_unused_address(KeychainKind::Internal).index, 0);
44164416
}
4417+
4418+
/// Pins down the eviction ordering our bitcoind mempool producer relies on: BDK keeps an
4419+
/// evicted transaction canonical for as long as the absence we reported predates the
4420+
/// `last_seen` we reported, and only drops it once the absence catches up. The producer has
4421+
/// to date absences on a clock that can clear a `last_seen` reloaded from the store, which
4422+
/// is why it reads the local clock rather than Bitcoin Core's mempool entry times — see
4423+
/// `chain::bitcoind::observation_time_secs`.
4424+
#[cfg(feature = "chain-bitcoind")]
4425+
#[tokio::test]
4426+
async fn mempool_eviction_needs_an_absence_dated_after_last_seen() {
4427+
// The mempool entry time our outgoing transaction was emitted with, i.e. the
4428+
// `last_seen` the wallet persists for it.
4429+
const LAST_SEEN: u64 = 200;
4430+
4431+
// A confirmed 100_000 sat input, spent by an unconfirmed 60_000 sat payment that paid a
4432+
// 1_000 sat fee and left 39_000 sats of change.
4433+
let store: Arc<DynStore> = Arc::new(DynStoreWrapper(InMemoryStore::new()));
4434+
let spend_txid = {
4435+
let wallet = new_test_wallet(Arc::clone(&store), false).await;
4436+
4437+
let (funding_tx, block_id) = {
4438+
let mut locked_wallet = wallet.inner.lock().unwrap();
4439+
let script_pubkey = locked_wallet
4440+
.reveal_next_address(KeychainKind::External)
4441+
.address
4442+
.script_pubkey();
4443+
let funding_tx = Transaction {
4444+
version: bitcoin::transaction::Version::TWO,
4445+
lock_time: LockTime::ZERO,
4446+
input: Vec::new(),
4447+
output: vec![TxOut { value: Amount::from_sat(100_000), script_pubkey }],
4448+
};
4449+
let block_id = BlockId {
4450+
height: locked_wallet.latest_checkpoint().height() + 1,
4451+
hash: bitcoin::BlockHash::from_byte_array([23; 32]),
4452+
};
4453+
(funding_tx, block_id)
4454+
};
4455+
let funding_txid = funding_tx.compute_txid();
4456+
let mut tx_update = TxUpdate::default();
4457+
tx_update.txs = vec![Arc::new(funding_tx)];
4458+
tx_update.anchors =
4459+
[(ConfirmationBlockTime { block_id, confirmation_time: 1 }, funding_txid)].into();
4460+
let chain = CheckPoint::from_block_ids([
4461+
wallet.inner.lock().unwrap().latest_checkpoint().block_id(),
4462+
block_id,
4463+
])
4464+
.unwrap();
4465+
wallet
4466+
.apply_update(Update { tx_update, chain: Some(chain), ..Default::default() })
4467+
.await
4468+
.unwrap();
4469+
assert_eq!(wallet.get_balances(0).unwrap(), (100_000, 100_000));
4470+
4471+
let change_script_pubkey = wallet
4472+
.inner
4473+
.lock()
4474+
.unwrap()
4475+
.reveal_next_address(KeychainKind::Internal)
4476+
.address
4477+
.script_pubkey();
4478+
let spend_tx = Transaction {
4479+
version: bitcoin::transaction::Version::TWO,
4480+
lock_time: LockTime::ZERO,
4481+
input: vec![bitcoin::TxIn {
4482+
previous_output: OutPoint { txid: funding_txid, vout: 0 },
4483+
..Default::default()
4484+
}],
4485+
output: vec![
4486+
TxOut {
4487+
value: Amount::from_sat(60_000),
4488+
script_pubkey: ScriptBuf::from_bytes(vec![0x51]),
4489+
},
4490+
TxOut { value: Amount::from_sat(39_000), script_pubkey: change_script_pubkey },
4491+
],
4492+
};
4493+
let spend_txid = spend_tx.compute_txid();
4494+
wallet.apply_mempool_txs(vec![(spend_tx, LAST_SEEN)], Vec::new()).await.unwrap();
4495+
assert_eq!(wallet.get_balances(0).unwrap(), (39_000, 39_000));
4496+
4497+
spend_txid
4498+
};
4499+
4500+
// Restart: the wallet reloads the transaction and its `last_seen` from the store.
4501+
let wallet = new_test_wallet(Arc::clone(&store), true).await;
4502+
assert_eq!(wallet.get_unconfirmed_txids(), vec![spend_txid]);
4503+
assert_eq!(wallet.get_balances(0).unwrap(), (39_000, 39_000));
4504+
4505+
// The transaction is gone from the mempool, with nothing confirming, replacing or
4506+
// descending from it. Absences predating `last_seen` leave it canonical: the input
4507+
// stays spent and the stale change keeps counting.
4508+
for stale in [0, LAST_SEEN - 1] {
4509+
wallet.apply_mempool_txs(Vec::new(), vec![(spend_txid, stale)]).await.unwrap();
4510+
assert_eq!(wallet.get_unconfirmed_txids(), vec![spend_txid]);
4511+
assert_eq!(wallet.get_balances(0).unwrap(), (39_000, 39_000));
4512+
}
4513+
4514+
// Once the absence reaches `last_seen`, the transaction stops being canonical: the
4515+
// confirmed input comes back and the stale change stops counting.
4516+
wallet.apply_mempool_txs(Vec::new(), vec![(spend_txid, LAST_SEEN)]).await.unwrap();
4517+
assert!(wallet.get_unconfirmed_txids().is_empty());
4518+
assert_eq!(wallet.get_balances(0).unwrap(), (100_000, 100_000));
4519+
}
44174520
}

0 commit comments

Comments
 (0)