From 2d853e59a2c8659c885124990bd3cfc59e31384e Mon Sep 17 00:00:00 2001 From: Borja Castellano Date: Fri, 11 Sep 2026 16:36:42 +0000 Subject: [PATCH 1/2] fix(key-wallet): re-apply a spend whose coin was funded after it MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Two mainnet restores of the same wallet ended at the same balance with 7112 and 7111 wallet records. Replaying each run's logged block applications offline, against the blocks it stored, reproduced both results exactly, with no divergence from the logs. The missing record was e66553f3…8c9e at height 2 185 057: it pays change to the BIP44 account and spends a CoinJoin coin funded at 2 182 877. When the spend is applied before its funding, the CoinJoin account cannot recognise it. The funding then parks the coin in `spent_before_funded` (#1001), and that only attributes the spend if its block is delivered again. In one run it was (funding at step 1494, spend at 1516); in the other it was not (spend at 1608, funding at 1627, no redelivery), so the result depended on delivery order. `WalletInfoInterface::unrecorded_spend_heights` reports, for a transaction, the heights of the blocks that spent its outputs before it arrived and that the owning account has not recorded yet. `process_block_for_wallets` returns them per wallet in `BlockProcessingResult::reapply_heights`, only heights above the block being applied, so re-applying cannot loop. `BlocksManager` re-applies those blocks from block storage straight away; every downloaded block is stored on arrival. Re-applications emit no `SyncEvent::BlockProcessed` and so never touch a batch's pending-block accounting. Offline, both orderings now end with identical per-account records (7112), with about 105 blocks re-applied from disk per restore. Co-Authored-By: Claude Opus 5 (1M context) Claude-Session: https://claude.ai/code/session_017DruChNTWXwJoWPartZwCf --- dash-spv/src/sync/blocks/manager.rs | 147 +++++++++++++----- key-wallet-manager/src/process_block.rs | 46 ++++++ .../src/test_utils/mock_wallet.rs | 13 +- key-wallet-manager/src/wallet_interface.rs | 1 + .../tests/observed_spent_outpoints_tests.rs | 22 ++- .../wallet_info_interface.rs | 19 ++- 6 files changed, 203 insertions(+), 45 deletions(-) diff --git a/dash-spv/src/sync/blocks/manager.rs b/dash-spv/src/sync/blocks/manager.rs index e4a1464d8..48b819184 100644 --- a/dash-spv/src/sync/blocks/manager.rs +++ b/dash-spv/src/sync/blocks/manager.rs @@ -3,6 +3,7 @@ //! Downloads blocks that matched wallet filters and processes them in height order. //! Subscribes to BlockNeeded events and emits BlockProcessed events. +use std::collections::{BTreeMap, BTreeSet}; use std::sync::Arc; use tokio::sync::RwLock; @@ -12,7 +13,8 @@ use crate::error::SyncResult; use crate::network::RequestSender; use crate::storage::{BlockHeaderStorage, BlockStorage}; use crate::sync::{BlocksProgress, SyncEvent, SyncManager, SyncState}; -use key_wallet_manager::WalletInterface; +use crate::types::HashedBlock; +use key_wallet_manager::{BlockProcessingResult, WalletId, WalletInterface}; /// Blocks manager for downloading and processing matching blocks. /// @@ -87,45 +89,10 @@ impl BlocksManager 0 { - tracing::info!( - "Found {} relevant transactions ({} new, {} existing) {} at height {}, new scripts: {}", - total_relevant, - result.new_txids.len(), - result.existing_txids.len(), - hash, - height, - new_scripts_total - ); - } + let result = self.apply_block(&block, height, &interested).await; // Collect confirmed txids before moving new_scripts out of result let confirmed_txids: Vec<_> = result.relevant_txids().cloned().collect(); - - // Collect new scripts for gap limit rescanning - let new_scripts = result.new_scripts; - if new_scripts_total > 0 { - tracing::debug!( - "Block {} generated {} new scripts for gap limit maintenance across {} wallets", - height, - new_scripts_total, - new_scripts.len() - ); - } - - self.progress.add_processed(1); - if total_relevant > 0 { - self.progress.add_relevant(1); - } - // Only count new transactions to avoid double-counting during rescans - self.progress.add_transactions(result.new_txids.len() as u32); self.progress.update_last_processed(height); last_applied = Some(height); @@ -133,9 +100,11 @@ impl BlocksManager BlocksManager, + ) -> BlockProcessingResult { + let hash = *block.hash(); + let mut wallet = self.wallet.write().await; + let result = wallet.process_block_for_wallets(block.block(), hash, height, wallets).await; + drop(wallet); + + let total_relevant = result.relevant_tx_count(); + let new_scripts_total: usize = result.new_scripts.values().map(|v| v.len()).sum(); + if total_relevant > 0 { + tracing::info!( + "Found {} relevant transactions ({} new, {} existing) {} at height {}, new scripts: {}", + total_relevant, + result.new_txids.len(), + result.existing_txids.len(), + hash, + height, + new_scripts_total + ); + } + if new_scripts_total > 0 { + tracing::debug!( + "Block {} generated {} new scripts for gap limit maintenance across {} wallets", + height, + new_scripts_total, + result.new_scripts.len() + ); + } + + self.progress.add_processed(1); + if total_relevant > 0 { + self.progress.add_relevant(1); + } + // Only count new transactions to avoid double-counting during rescans + self.progress.add_transactions(result.new_txids.len() as u32); + result + } + + async fn reapply_blocks(&mut self, reapply: BTreeMap>) { + let mut queue: BTreeSet<(u32, WalletId)> = reapply + .into_iter() + .flat_map(|(wallet_id, heights)| heights.into_iter().map(move |h| (h, wallet_id))) + .collect(); + while let Some((height, wallet_id)) = queue.pop_first() { + let block = match self.block_storage.read().await.load_block(height).await { + Ok(Some(block)) => block, + Ok(None) => { + tracing::warn!("Cannot re-apply block at height {}: not in storage", height); + continue; + } + Err(e) => { + tracing::warn!("Cannot re-apply block at height {}: {}", height, e); + continue; + } + }; + let result = self.apply_block(&block, height, &BTreeSet::from([wallet_id])).await; + queue.extend( + result.reapply_heights.into_iter().flat_map(|(wallet_id, heights)| { + heights.into_iter().map(move |h| (h, wallet_id)) + }), + ); + } + } } impl std::fmt::Debug @@ -353,6 +390,40 @@ mod tests { assert_eq!(processed[0].1, 100); } + #[tokio::test] + async fn test_process_buffered_blocks_reapplies_requested_stored_block() { + let storage = DiskStorageManager::with_temp_dir().await.unwrap(); + let mut wallet = MockWallet::new(); + wallet.set_reapply_heights(100, BTreeSet::from([200])); + let wallet = Arc::new(RwLock::new(wallet)); + let mut manager: TestBlocksManager = + BlocksManager::new(wallet.clone(), storage.block_headers(), storage.blocks()).await; + manager.progress.set_state(SyncState::Syncing); + + manager + .block_storage + .write() + .await + .store_block(200, HashedBlock::dummy(200, vec![])) + .await + .unwrap(); + manager.pipeline.add_from_storage( + HashedBlock::dummy(100, vec![]), + 100, + BTreeSet::from([MOCK_WALLET_ID]), + ); + + let events = manager.process_buffered_blocks().await.unwrap(); + assert_eq!( + events.iter().filter(|e| matches!(e, SyncEvent::BlockProcessed { .. })).count(), + 1 + ); + + let processed = wallet.read().await.processed_blocks(); + let heights: Vec = processed.lock().await.iter().map(|(_, h)| *h).collect(); + assert_eq!(heights, vec![100, 200]); + } + /// A wallet that is NOT in the pipeline's interested set must not be /// routed the block. Two wallets are registered, but only `wallet_in` /// appears in the routed set; the other wallet's processed log must diff --git a/key-wallet-manager/src/process_block.rs b/key-wallet-manager/src/process_block.rs index de39fad84..5d9b61092 100644 --- a/key-wallet-manager/src/process_block.rs +++ b/key-wallet-manager/src/process_block.rs @@ -50,6 +50,7 @@ impl WalletInterface for WalletM let mut per_wallet_inserted: BTreeMap> = BTreeMap::new(); let mut per_wallet_updated: BTreeMap> = BTreeMap::new(); let mut per_wallet_derived: BTreeMap> = BTreeMap::new(); + let mut relevant_positions: BTreeMap> = BTreeMap::new(); for (position, tx) in block.txdata.iter().enumerate() { // Stamp each record with its `block.vtx` index so consumers @@ -71,6 +72,9 @@ impl WalletInterface for WalletM result.existing_txids.push(tx.txid()); } } + for wallet_id in &check_result.affected_wallets { + relevant_positions.entry(*wallet_id).or_default().push(position); + } for (wallet_id, derived) in check_result.new_addresses { let scripts = @@ -118,6 +122,20 @@ impl WalletInterface for WalletM ); } + for (wallet_id, positions) in relevant_positions { + let Some(info) = self.wallet_infos.get(&wallet_id) else { + continue; + }; + let heights: BTreeSet = positions + .into_iter() + .flat_map(|position| info.unrecorded_spend_heights(&block.txdata[position])) + .filter(|spend_height| *spend_height > height) + .collect(); + if !heights.is_empty() { + result.reapply_heights.insert(wallet_id, heights); + } + } + self.finalize_block_advance( height, wallets, @@ -675,6 +693,34 @@ mod tests { assert_eq!(manager.last_processed_height(), 0); } + #[tokio::test] + async fn test_funding_after_its_spend_asks_to_reapply_the_spend_block() { + let (mut manager, wallet_id, addr) = setup_manager_with_wallet(); + let funding = create_tx_paying_to(&addr, 0xaa); + let spend = spend_first_output_of(&funding); + let wallets = BTreeSet::from([wallet_id]); + + let mut spend_block = make_block(vec![spend]); + spend_block.header.nonce = 1; + let funding_block = make_block(vec![funding]); + + manager + .process_block_for_wallets(&spend_block, spend_block.block_hash(), 200, &wallets) + .await; + let result = manager + .process_block_for_wallets(&funding_block, funding_block.block_hash(), 100, &wallets) + .await; + assert_eq!(result.reapply_heights, BTreeMap::from([(wallet_id, BTreeSet::from([200]))])); + + manager + .process_block_for_wallets(&spend_block, spend_block.block_hash(), 200, &wallets) + .await; + let again = manager + .process_block_for_wallets(&funding_block, funding_block.block_hash(), 100, &wallets) + .await; + assert!(again.reapply_heights.is_empty()); + } + #[tokio::test] async fn test_sweep_expired_reservations_fans_out_over_wallets() { let (mut manager, wallet_id, _addr) = setup_manager_with_wallet(); diff --git a/key-wallet-manager/src/test_utils/mock_wallet.rs b/key-wallet-manager/src/test_utils/mock_wallet.rs index 63e02b487..dfbcf6e29 100644 --- a/key-wallet-manager/src/test_utils/mock_wallet.rs +++ b/key-wallet-manager/src/test_utils/mock_wallet.rs @@ -6,7 +6,7 @@ use dashcore::ephemerealdata::instant_lock::InstantLock; use dashcore::prelude::CoreBlockHeight; use dashcore::{Address, Block, OutPoint, ScriptBuf, Transaction, Txid}; use key_wallet::transaction_checking::TransactionContext; -use std::collections::BTreeSet; +use std::collections::{BTreeMap, BTreeSet}; use std::sync::Arc; use tokio::sync::{broadcast, Mutex}; @@ -43,6 +43,7 @@ pub struct MockWallet { pub processed_instant_locks: InstantLockCaptures, /// Monitor revision counter for staleness detection. monitor_revision: u64, + reapply_heights: BTreeMap>, } impl Default for MockWallet { @@ -70,9 +71,14 @@ impl MockWallet { status_changes: Arc::new(Mutex::new(Vec::new())), processed_instant_locks: Arc::new(Mutex::new(Vec::new())), monitor_revision: 0, + reapply_heights: BTreeMap::new(), } } + pub fn set_reapply_heights(&mut self, height: u32, heights: BTreeSet) { + self.reapply_heights.insert(height, heights); + } + /// Override the wallet id used for per-wallet API surfaces. pub fn set_wallet_id(&mut self, wallet_id: WalletId) { self.wallet_id = wallet_id; @@ -145,6 +151,11 @@ impl WalletInterface for MockWallet { new_txids: block.txdata.iter().map(|tx| tx.txid()).collect(), existing_txids: Vec::new(), new_scripts: Default::default(), + reapply_heights: self + .reapply_heights + .get(&height) + .map(|heights| BTreeMap::from([(self.wallet_id, heights.clone())])) + .unwrap_or_default(), } } diff --git a/key-wallet-manager/src/wallet_interface.rs b/key-wallet-manager/src/wallet_interface.rs index 4e15faa55..50fbecefb 100644 --- a/key-wallet-manager/src/wallet_interface.rs +++ b/key-wallet-manager/src/wallet_interface.rs @@ -21,6 +21,7 @@ pub struct BlockProcessingResult { /// Cached scriptPubKeys of addresses freshly generated per wallet during /// gap-limit maintenance. pub new_scripts: BTreeMap>, + pub reapply_heights: BTreeMap>, } /// Result of processing a mempool transaction through the wallet diff --git a/key-wallet/src/tests/observed_spent_outpoints_tests.rs b/key-wallet/src/tests/observed_spent_outpoints_tests.rs index 4f89fd701..ee6e00bcc 100644 --- a/key-wallet/src/tests/observed_spent_outpoints_tests.rs +++ b/key-wallet/src/tests/observed_spent_outpoints_tests.rs @@ -211,9 +211,9 @@ fn adding_account_from_xpub_rewinds_sync_checkpoint() { /// dropped. The balance and UTXO set are unaffected (the output is genuinely /// spent on-chain), so this is purely about not losing the history record. /// -/// This is the spend-first / out-of-order ordering produced by the committed- -/// range rescan (`track_for_new_scripts`): the forward scan already processed -/// the spend, then the old funding block is re-applied. #649's fix keeps the +/// This is the spend-first / out-of-order ordering produced by a rescan +/// (`track_for_new_scripts`): the scan already processed the spend, then the +/// older funding block is re-applied. #649's fix keeps the /// record: `check_transaction_for_match` classifies relevance by address /// membership (never gated on spent-status, so the fully-spent funding is still /// relevant), and `ManagedCoreFundsAccount::record_transaction` unconditionally @@ -259,8 +259,8 @@ async fn born_fully_spent_funding_tx_is_recorded_in_history() { }; // Spend-first: the spend's block (height 200) is applied before the - // funding's block (height 100), exactly as the committed-range rescan - // re-applies old funding blocks after their spends. + // funding's block (height 100), exactly as a rescan re-applies older + // funding blocks after their spends. let spend_ctx = TransactionContext::InBlock(BlockInfo::new( 200, BlockHash::from_slice(&[2u8; 32]).expect("hash"), @@ -458,6 +458,18 @@ async fn a_held_output_stays_out_of_utxos_once_its_observed_spend_is_pruned() { assert_eq!(ctx.managed_wallet.balance.total(), 0); } +#[tokio::test] +async fn funding_after_its_spend_reports_the_spend_height() { + use crate::wallet::managed_wallet_info::wallet_info_interface::WalletInfoInterface; + use std::collections::BTreeSet; + + let (mut ctx, funding, spend) = spend_first_context(in_block(100, 1)).await; + assert_eq!(ctx.managed_wallet.unrecorded_spend_heights(&funding), BTreeSet::from([200])); + + ctx.check_transaction(&spend, in_block(200, 2)).await; + assert!(ctx.managed_wallet.unrecorded_spend_heights(&funding).is_empty()); +} + /// Abandoning the funding transaction takes its held output with it: the coin /// was never ours, so a spend of it must stop being recognisable. #[tokio::test] diff --git a/key-wallet/src/wallet/managed_wallet_info/wallet_info_interface.rs b/key-wallet/src/wallet/managed_wallet_info/wallet_info_interface.rs index f2c5d9a9b..1dabac23d 100644 --- a/key-wallet/src/wallet/managed_wallet_info/wallet_info_interface.rs +++ b/key-wallet/src/wallet/managed_wallet_info/wallet_info_interface.rs @@ -20,7 +20,7 @@ use dashcore::address::Payload; use dashcore::ephemerealdata::chain_lock::ChainLock; use dashcore::ephemerealdata::instant_lock::InstantLock; use dashcore::prelude::CoreBlockHeight; -use dashcore::{Address as DashAddress, ScriptBuf, Transaction, Txid}; +use dashcore::{Address as DashAddress, OutPoint, ScriptBuf, Transaction, Txid}; /// Outcome of [`WalletInfoInterface::apply_chain_lock`]. /// @@ -280,6 +280,8 @@ pub trait WalletInfoInterface: Sized + WalletTransactionChecker + ManagedAccount /// sweep removing a loser, so this is broader than "a UTXO was marked". fn mark_instant_send_utxos(&mut self, txid: &Txid, lock: &InstantLock) -> bool; + fn unrecorded_spend_heights(&self, tx: &Transaction) -> BTreeSet; + /// Return the aggregated monitor revision across all accounts. /// Increments whenever the monitored address set changes. fn monitor_revision(&self) -> u64 { @@ -612,6 +614,21 @@ impl WalletInfoInterface for ManagedWalletInfo { fn monitor_revision(&self) -> u64 { self.accounts.all_accounts().iter().map(|a| a.monitor_revision()).sum() } + + fn unrecorded_spend_heights(&self, tx: &Transaction) -> BTreeSet { + let txid = tx.txid(); + (0..tx.output.len() as u32) + .map(|vout| OutPoint::new(txid, vout)) + .filter(|outpoint| { + self.accounts + .all_accounts() + .into_iter() + .filter_map(|account| account.as_funds()) + .any(|account| account.spent_before_funded.contains_key(outpoint)) + }) + .filter_map(|outpoint| self.observed_spent_outpoints.get(&outpoint).copied()) + .collect() + } } #[cfg(test)] From a49c7298528c059503a7334c3dd5a05fdf9dca74 Mon Sep 17 00:00:00 2001 From: Borja Castellano Date: Wed, 16 Sep 2026 15:05:45 +0000 Subject: [PATCH 2/2] fix(dash-spv): don't count block re-applications as processed blocks Re-applying a stored block went through the same accounting as a first pass, so every re-application added to `processed` and `relevant`: a mainnet restore reported 11 294 processed blocks for 10 326 downloaded plus 855 loaded from storage. Only new transactions are still counted on re-application, since a transaction first recorded there was never counted. Co-Authored-By: Claude Opus 5 (1M context) --- dash-spv/src/sync/blocks/manager.rs | 9 +++++---- 1 file changed, 5 insertions(+), 4 deletions(-) diff --git a/dash-spv/src/sync/blocks/manager.rs b/dash-spv/src/sync/blocks/manager.rs index 48b819184..3c7e50d6e 100644 --- a/dash-spv/src/sync/blocks/manager.rs +++ b/dash-spv/src/sync/blocks/manager.rs @@ -91,6 +91,10 @@ impl BlocksManager 0 { + self.progress.add_relevant(1); + } // Collect confirmed txids before moving new_scripts out of result let confirmed_txids: Vec<_> = result.relevant_txids().cloned().collect(); self.progress.update_last_processed(height); @@ -172,10 +176,6 @@ impl BlocksManager 0 { - self.progress.add_relevant(1); - } // Only count new transactions to avoid double-counting during rescans self.progress.add_transactions(result.new_txids.len() as u32); result @@ -422,6 +422,7 @@ mod tests { let processed = wallet.read().await.processed_blocks(); let heights: Vec = processed.lock().await.iter().map(|(_, h)| *h).collect(); assert_eq!(heights, vec![100, 200]); + assert_eq!(manager.progress.processed(), 1); } /// A wallet that is NOT in the pipeline's interested set must not be