Skip to content
Merged
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
146 changes: 109 additions & 37 deletions dash-spv/src/sync/blocks/manager.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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.
///
Expand Down Expand Up @@ -87,55 +89,26 @@ impl<H: BlockHeaderStorage, B: BlockStorage, W: WalletInterface> BlocksManager<H

// Process the block only for the wallets whose filter matched it.
// Already-synced wallets that did not match are not touched.
let mut wallet = self.wallet.write().await;
let result =
wallet.process_block_for_wallets(block.block(), hash, height, &interested).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
);
}

// 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()
);
}
let result = self.apply_block(&block, height, &interested).await;

self.progress.add_processed(1);
if total_relevant > 0 {
if result.relevant_tx_count() > 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);
// 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);
last_applied = Some(height);

events.push(SyncEvent::BlockProcessed {
block_hash: hash,
height,
wallets: interested,
new_scripts,
new_scripts: result.new_scripts,
confirmed_txids,
});

self.reapply_blocks(result.reapply_heights).await;
}

// Blocks are drained in strict height order, so `last_applied` is the
Expand Down Expand Up @@ -169,6 +142,70 @@ impl<H: BlockHeaderStorage, B: BlockStorage, W: WalletInterface> BlocksManager<H

Ok(events)
}

async fn apply_block(
&mut self,
block: &HashedBlock,
height: u32,
wallets: &BTreeSet<WalletId>,
) -> 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()
);
}

// 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<WalletId, BTreeSet<u32>>) {
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;
}
Comment thread
ZocoLini marked this conversation as resolved.
};
let result = self.apply_block(&block, height, &BTreeSet::from([wallet_id])).await;
Comment thread
coderabbitai[bot] marked this conversation as resolved.
queue.extend(
result.reapply_heights.into_iter().flat_map(|(wallet_id, heights)| {
heights.into_iter().map(move |h| (h, wallet_id))
}),
);
Comment thread
coderabbitai[bot] marked this conversation as resolved.
}
}
}

impl<H: BlockHeaderStorage, B: BlockStorage, W: WalletInterface> std::fmt::Debug
Expand Down Expand Up @@ -353,6 +390,41 @@ 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<u32> = 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
/// routed the block. Two wallets are registered, but only `wallet_in`
/// appears in the routed set; the other wallet's processed log must
Expand Down
46 changes: 46 additions & 0 deletions key-wallet-manager/src/process_block.rs
Original file line number Diff line number Diff line change
Expand Up @@ -50,6 +50,7 @@ impl<T: WalletInfoInterface + Send + Sync + 'static> WalletInterface for WalletM
let mut per_wallet_inserted: BTreeMap<WalletId, Vec<TransactionRecord>> = BTreeMap::new();
let mut per_wallet_updated: BTreeMap<WalletId, Vec<TransactionRecord>> = BTreeMap::new();
let mut per_wallet_derived: BTreeMap<WalletId, Vec<DerivedAddressInfo>> = BTreeMap::new();
let mut relevant_positions: BTreeMap<WalletId, Vec<usize>> = BTreeMap::new();

for (position, tx) in block.txdata.iter().enumerate() {
// Stamp each record with its `block.vtx` index so consumers
Expand All @@ -71,6 +72,9 @@ impl<T: WalletInfoInterface + Send + Sync + 'static> 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 =
Expand Down Expand Up @@ -118,6 +122,20 @@ impl<T: WalletInfoInterface + Send + Sync + 'static> WalletInterface for WalletM
);
}

for (wallet_id, positions) in relevant_positions {
let Some(info) = self.wallet_infos.get(&wallet_id) else {
continue;
};
let heights: BTreeSet<CoreBlockHeight> = 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,
Expand Down Expand Up @@ -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();
Expand Down
13 changes: 12 additions & 1 deletion key-wallet-manager/src/test_utils/mock_wallet.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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};

Expand Down Expand Up @@ -43,6 +43,7 @@ pub struct MockWallet {
pub processed_instant_locks: InstantLockCaptures,
/// Monitor revision counter for staleness detection.
monitor_revision: u64,
reapply_heights: BTreeMap<u32, BTreeSet<u32>>,
}

impl Default for MockWallet {
Expand Down Expand Up @@ -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<u32>) {
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;
Expand Down Expand Up @@ -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(),
}
}

Expand Down
1 change: 1 addition & 0 deletions key-wallet-manager/src/wallet_interface.rs
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,7 @@ pub struct BlockProcessingResult {
/// Cached scriptPubKeys of addresses freshly generated per wallet during
/// gap-limit maintenance.
pub new_scripts: BTreeMap<WalletId, Vec<ScriptBuf>>,
pub reapply_heights: BTreeMap<WalletId, BTreeSet<CoreBlockHeight>>,
}

/// Result of processing a mempool transaction through the wallet
Expand Down
22 changes: 17 additions & 5 deletions key-wallet/src/tests/observed_spent_outpoints_tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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"),
Expand Down Expand Up @@ -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]
Expand Down
Loading
Loading