diff --git a/dash-spv/src/sync/filters/batch.rs b/dash-spv/src/sync/filters/batch.rs index 3fd609b79..67263ef3b 100644 --- a/dash-spv/src/sync/filters/batch.rs +++ b/dash-spv/src/sync/filters/batch.rs @@ -22,8 +22,6 @@ pub(super) struct FiltersBatch { scanned: bool, /// Number of blocks still being downloaded for this batch. pending_blocks: u32, - /// Whether rescan has been completed for this batch. - rescan_complete: bool, /// Wallets that were behind for this batch's height range at scan time — /// and therefore need their `synced_height` advanced when the batch /// commits — each mapped to the wallet's `account_generation` at scan @@ -33,10 +31,6 @@ pub(super) struct FiltersBatch { /// current account set (dashpay/rust-dashcore#649). Already-synced wallets /// must not be touched. scanned_wallets: BTreeMap, - /// Cached scriptPubKeys discovered during block processing that still - /// need rescan, attributed per wallet so we can rerun matching only - /// against the wallet that produced each new script. - collected_scripts: HashMap>, /// Every script already matched against this batch's filters, per wallet. tested_scripts: HashMap>, } @@ -55,9 +49,7 @@ impl FiltersBatch { verified: false, scanned: false, pending_blocks: 0, - rescan_complete: false, scanned_wallets: BTreeMap::new(), - collected_scripts: HashMap::new(), tested_scripts: HashMap::new(), } } @@ -106,22 +98,6 @@ impl FiltersBatch { self.pending_blocks = self.pending_blocks.saturating_sub(1); self.pending_blocks } - /// Returns whether rescan has been completed for this batch. - pub(super) fn rescan_complete(&self) -> bool { - self.rescan_complete - } - /// Mark rescan as complete for this batch. - pub(super) fn mark_rescan_complete(&mut self) { - self.rescan_complete = true; - } - /// Add scriptPubKeys discovered during block processing for later rescan. - pub(super) fn add_scripts_for_wallet( - &mut self, - wallet_id: WalletId, - scripts: impl IntoIterator, - ) { - self.collected_scripts.entry(wallet_id).or_default().extend(scripts); - } /// Record that `scripts` have been matched against this batch's filters. pub(super) fn mark_tested>( &mut self, @@ -141,10 +117,6 @@ impl FiltersBatch { monitored.iter().filter(move |script| tested.is_none_or(|t| !t.contains(*script))) } - /// Take collected per-wallet scripts for rescan, leaving the map empty. - pub(super) fn take_collected_scripts(&mut self) -> HashMap> { - std::mem::take(&mut self.collected_scripts) - } /// Record the wallets that were behind for this batch at scan time, each /// with its `account_generation` snapshot. pub(super) fn set_scanned_wallets(&mut self, wallets: BTreeMap) { diff --git a/dash-spv/src/sync/filters/block_match_tracker.rs b/dash-spv/src/sync/filters/block_match_tracker.rs index eb81d5344..e0d911b47 100644 --- a/dash-spv/src/sync/filters/block_match_tracker.rs +++ b/dash-spv/src/sync/filters/block_match_tracker.rs @@ -148,11 +148,6 @@ impl BlockMatchTracker { self.processed_blocks_per_wallet.split_off(&(height + 1)); } - /// True while matched blocks are still awaiting their `BlockProcessed`. - pub(super) fn has_blocks_in_flight(&self) -> bool { - !self.blocks_remaining.is_empty() - } - /// True when there is no in-flight or processed-record state. pub(super) fn is_empty(&self) -> bool { self.blocks_remaining.is_empty() && self.processed_blocks_per_wallet.is_empty() diff --git a/dash-spv/src/sync/filters/coinjoin_gap_discovery_tests.rs b/dash-spv/src/sync/filters/coinjoin_gap_discovery_tests.rs index bf064fd97..235645464 100644 --- a/dash-spv/src/sync/filters/coinjoin_gap_discovery_tests.rs +++ b/dash-spv/src/sync/filters/coinjoin_gap_discovery_tests.rs @@ -12,7 +12,7 @@ //! `rescan_batch` now re-queues already-processed blocks of ACTIVE batches //! against newly derived scripts, to a fixpoint. //! -//! These tests pin both the fixed behaviour and the remaining hole: +//! These tests pin the fixed behaviour: //! //! 1. [`coinjoin_gap_limit_dense_same_batch_recovers`] — the empirical //! stall-at-59 shape (two dense blocks in one batch). GREEN since #820; @@ -20,16 +20,6 @@ //! 2. [`coinjoin_gap_limit_inversion_within_batch_recovers`] — a gap-window //! output in an EARLIER block of the SAME (still-active) batch is //! recovered by the commit-time rescan. GREEN since #820. -//! 3. [`coinjoin_gap_limit_stall_across_committed_batch`] — the SAME shape -//! with the earlier block in an already-COMMITTED batch (#846): rescans -//! only reach `active_batches`, and committed batches are gone. GREEN -//! since `rescan_committed_range` re-tests newly derived scripts against -//! the STORED filters below the committing batch. -//! 4. [`committed_range_sweep_coalesces_across_batch_commits`] — the cost -//! side of #846's fix: N script-carrying commits share one drain-time -//! committed-range sweep instead of walking the stored history once per -//! commit, while a late-derived script still finds its -//! committed-prefix block before `FiltersSyncComplete`. //! //! Each test drives the manager exactly the way the production event loop //! does: `try_process_batch` → `BlocksNeeded` → (blocks-manager stand-in) @@ -296,8 +286,7 @@ async fn coinjoin_gap_limit_dense_same_batch_recovers() { /// only at height 20. The initial scan cannot see block A (nothing watched /// matches it), but once block B's processing extends the window past G+21 /// the commit-time rescan re-tests block A's filter and recovers it. GREEN -/// since #820 — contrast for the cross-batch test below, isolating the -/// commit boundary as the broken seam. +/// since #820. #[tokio::test] async fn coinjoin_gap_limit_inversion_within_batch_recovers() { let (mut manager, wallet, wallet_id) = setup().await; @@ -324,263 +313,3 @@ async fn coinjoin_gap_limit_inversion_within_batch_recovers() { recovered by the commit-time rescan (PR #820). used_count={used_count}" ); } - -/// Gap-window outputs in an already-COMMITTED batch (#846). -/// -/// Same funding shape as the within-batch inversion test, but the early -/// block (indices G+10..=G+21, height 10) sits in batch 0..=99 while the -/// in-window block (indices 0..=29) sits at height 110 in batch 100..=199. -/// Batch 0 scans clean (nothing watched matches) and commits. Processing the -/// height-110 block extends the window past G+21, and those scripts DO match -/// block 10's filter — but `rescan_batch` only reaches `active_batches`, and -/// committed batches are gone (`try_commit_batches` removes them; the -/// tracker prunes at-or-below the committed height). Indices G+10..=G+21 — -/// squarely inside the BIP-44/CoinJoin gap-limit recovery contract -/// (G+21 < 29 + 1 + G) — used to stay invisible forever, along with their -/// funds; a fresh re-sync from genesis hit the same wall deterministically. -/// -/// GREEN since `rescan_committed_range`: newly derived scripts are re-tested -/// against the persisted filters below the committing batch (BIP-158 filters -/// are address-independent, so re-matching needs no re-download), and hits -/// flow through the `track_for_new_scripts` re-download path to the same -/// commit-time fixpoint. `highest_used` reaches G+21. -#[tokio::test] -async fn coinjoin_gap_limit_stall_across_committed_batch() { - let (mut manager, wallet, wallet_id) = setup().await; - let addresses = coinjoin_external_addresses(&wallet, &wallet_id, (G + 22) as u32).await; - - let (block_a, filter_a, key_a) = block_paying(10, &addresses[(G + 10)..=(G + 21)]); - let (block_b, filter_b, key_b) = block_paying(110, &addresses[0..=29]); - - // Uphold the production invariant the injected batches imply: every - // height at or below `stored_height` has its header and filter - // persisted (store_and_match_batches stores a batch's filters before - // stored_height advances past it). The committed-range recovery path - // re-tests exactly this stored data, so the invariant is load-bearing - // here: batch 0 commits before block B's processing derives the missing - // scripts, and by then its in-memory filters are gone. - { - let mut header_storage = manager.header_storage.write().await; - let mut filter_storage = manager.filter_storage.write().await; - for height in 0..=99u32 { - let (header, filter_bytes) = if height == 10 { - (block_a.header, filter_a.content.clone()) - } else { - let filler = Block::dummy(height, vec![]); - let filter = BlockFilter::dummy(&filler); - (filler.header, filter.content) - }; - header_storage - .store_headers_at_height(&[header.into()], height) - .await - .expect("seed header"); - filter_storage.store_filter(height, &filter_bytes).await.expect("seed filter"); - } - } - - let blocks: HashMap = - HashMap::from([(block_a.block_hash(), block_a), (block_b.block_hash(), block_b)]); - - let mut batch_0 = FiltersBatch::new(0, 99, HashMap::from([(key_a, filter_a)])); - batch_0.mark_verified(); - manager.active_batches.insert(0, batch_0); - let mut batch_1 = FiltersBatch::new(100, 199, HashMap::from([(key_b, filter_b)])); - batch_1.mark_verified(); - manager.active_batches.insert(100, batch_1); - manager.progress.update_stored_height(199); - - let initial_events = manager.try_process_batch().await.unwrap(); - drive_to_quiescence(&mut manager, &wallet, &blocks, initial_events).await; - - let (highest_used, highest_generated, used_count) = - coinjoin_pool_state(&wallet, &wallet_id).await; - // Sanity: the in-window block was found and the gap window extended past - // index G+21, so the missed indices ARE inside the watched range by now. - assert!( - highest_generated >= Some((G + 21) as u32), - "gap maintenance must have extended the watch window past index G+21 \ - (got {highest_generated:?})" - ); - assert_eq!( - highest_used, - Some((G + 21) as u32), - "CoinJoin External indices G+10..=G+21 were funded at height 10 in a batch that \ - committed before their scripts were derived, and the new-script rescan never \ - looks below the committed boundary (rescan_batch only reaches active_batches; \ - BlockMatchTracker/commit pruning drops the range). The addresses are within \ - the gap-limit recovery contract and are watched now (highest_generated = \ - {highest_generated:?}), yet their outputs stay invisible: highest_used stalls \ - at {highest_used:?}, used_count={used_count}. Fix direction: key re-scan \ - suppression by (wallet, address/script) instead of block/commit progress, or \ - trigger a below-committed-height rescan for a wallet whose gap maintenance \ - derives scripts mid-sync." - ); -} - -/// Committed-range sweeps coalesce across batch commits. -/// -/// Four batches; the last three each contain one in-window block whose -/// processing derives new scripts (each funds the next run of CoinJoin -/// indices, so gap maintenance extends the watch window at every commit). -/// Per-commit sweeping walks the entire committed prefix once per -/// script-carrying commit — on a real mainnet restore that shape produced -/// 191 full-prefix sweeps totalling ~14.5 minutes. Coalesced, the -/// accumulated scripts cross the stored history when the forward pipeline -/// drains: one sweep, plus one follow-up round for the scripts derived from -/// the block that sweep recovers. -/// -/// The early block (height 10, beyond-window indices G+10..=G+21, unwatched -/// when its batch scans and commits) pins the #846 correctness contract at -/// the same time: the deferred sweep must still find it, and -/// `FiltersSyncComplete` must not be emitted before it has been found and -/// applied. `committed_range_sweeps` counts sweeps that reach the chunk walk -/// in `rescan_committed_range`. -#[tokio::test] -async fn committed_range_sweep_coalesces_across_batch_commits() { - let (mut manager, wallet, wallet_id) = setup().await; - let addresses = coinjoin_external_addresses(&wallet, &wallet_id, (G + 22) as u32).await; - - // Beyond-window block in the range that commits first (#846 shape). - let (block_early, filter_early, key_early) = block_paying(10, &addresses[(G + 10)..=(G + 21)]); - // One in-window block per later batch. Each extends `highest_used` by 30, - // so every one of these batches carries newly derived scripts into its - // commit. After block 110 the generated window reaches 29 + G >= G + 21, - // so the early block's scripts are among the first commit's derivations. - let (block_1, filter_1, key_1) = block_paying(110, &addresses[0..=29]); - let (block_2, filter_2, key_2) = block_paying(210, &addresses[30..=59]); - let (block_3, filter_3, key_3) = block_paying(310, &addresses[60..=89]); - - // Persist headers and filters for the committed prefix 0..=299 — the - // range the drain-time sweep reloads from storage. Real data at the - // three block heights, filler elsewhere. - { - let mut header_storage = manager.header_storage.write().await; - let mut filter_storage = manager.filter_storage.write().await; - for height in 0..=299u32 { - let (header, filter_bytes) = match height { - 10 => (block_early.header, filter_early.content.clone()), - 110 => (block_1.header, filter_1.content.clone()), - 210 => (block_2.header, filter_2.content.clone()), - _ => { - let filler = Block::dummy(height, vec![]); - let filter = BlockFilter::dummy(&filler); - (filler.header, filter.content) - } - }; - header_storage - .store_headers_at_height(&[header.into()], height) - .await - .expect("seed header"); - filter_storage.store_filter(height, &filter_bytes).await.expect("seed filter"); - } - } - - let blocks: HashMap = HashMap::from([ - (block_early.block_hash(), block_early), - (block_1.block_hash(), block_1), - (block_2.block_hash(), block_2), - (block_3.block_hash(), block_3), - ]); - - for (start, filters) in [ - (0u32, HashMap::from([(key_early, filter_early)])), - (100, HashMap::from([(key_1, filter_1)])), - (200, HashMap::from([(key_2, filter_2)])), - (300, HashMap::from([(key_3, filter_3)])), - ] { - let mut batch = FiltersBatch::new(start, start + 99, filters); - batch.mark_verified(); - manager.active_batches.insert(start, batch); - } - manager.progress.update_stored_height(399); - // Both halves of the drain gate have to be live, or the test only proves - // the `active_batches.len() == 1` half: with the tips left at their - // default 0, `end_height() >= filter_header_tip_height()` is trivially - // true and the comparison that keeps the sweep from firing while more - // batches are still to come is never exercised. Same for the target - // height, which `FiltersSyncComplete` is checked against below. - manager.progress.update_filter_header_tip_height(399); - manager.progress.update_target_height(399); - - // Drive the production event loop to quiescence like `drive_to_quiescence` - // does, additionally watching for `FiltersSyncComplete` so the completion - // contract can be asserted at the moment it is emitted. - let (tx, _rx) = unbounded_channel(); - let requests = RequestSender::new(tx); - let mut events = manager.try_process_batch().await.unwrap(); - let mut sync_complete_seen = false; - 'rounds: for _round in 0..64 { - let mut pending: BTreeMap<(u32, BlockHash), BTreeSet> = BTreeMap::new(); - for event in events.drain(..) { - match event { - SyncEvent::BlocksNeeded { - blocks: needed, - } => { - for (key, wallets) in needed { - pending.entry((key.height(), *key.hash())).or_default().extend(wallets); - } - } - SyncEvent::FiltersSyncComplete { - .. - } => { - let (highest_used, _, _) = coinjoin_pool_state(&wallet, &wallet_id).await; - assert_eq!( - highest_used, - Some((G + 21) as u32), - "FiltersSyncComplete emitted before the deferred committed-range \ - sweep recovered the height-10 block: a script derived after its \ - range committed was never tested against the committed prefix" - ); - sync_complete_seen = true; - } - _ => {} - } - } - if pending.is_empty() { - break 'rounds; - } - for ((height, block_hash), wallets) in pending { - let block = blocks.get(&block_hash).expect("BlocksNeeded for an unknown test block"); - let result = wallet - .write() - .await - .process_block_for_wallets(block, block_hash, height, &wallets) - .await; - let confirmed_txids = result.relevant_txids().cloned().collect(); - let event = SyncEvent::BlockProcessed { - block_hash, - height, - wallets, - new_scripts: result.new_scripts, - confirmed_txids, - }; - events.extend( - manager.handle_sync_event(&event, &requests).await.expect("BlockProcessed"), - ); - } - } - - assert!(sync_complete_seen, "the run must reach FiltersSyncComplete"); - - let (highest_used, _, used_count) = coinjoin_pool_state(&wallet, &wallet_id).await; - assert_eq!( - highest_used, - Some((G + 21) as u32), - "the height-10 block's beyond-window outputs must be recovered by the \ - drain-time committed-range sweep (used_count={used_count})" - ); - assert_eq!(used_count, 90 + 12, "indices 0..=89 and G+10..=G+21 must all be marked used"); - - assert!( - manager.committed_range_sweeps >= 1, - "the deferred committed-range sweep must still run before completion" - ); - assert!( - manager.committed_range_sweeps <= 2, - "three script-carrying batch commits must share the committed-range sweep \ - (one drain-time sweep plus one follow-up round for the scripts derived from \ - the recovered block); got {} sweeps — one sweep per commit means the \ - coalescing regressed", - manager.committed_range_sweeps - ); -} diff --git a/dash-spv/src/sync/filters/manager.rs b/dash-spv/src/sync/filters/manager.rs index f5e0cdf05..be77c0668 100644 --- a/dash-spv/src/sync/filters/manager.rs +++ b/dash-spv/src/sync/filters/manager.rs @@ -91,21 +91,6 @@ pub struct FiltersManager< /// `BlockProcessed` and the per-wallet record of which wallets already /// have a given processed block applied. pub(super) tracker: BlockMatchTracker, - /// Scripts already forward-rescanned but still awaiting the combined - /// backward sweep over the committed range (#846). Held at the manager - /// level and carried across batch commits: intermediate commits only - /// accumulate here, and `try_commit_batches` runs the sweep when the - /// forward pipeline drains — the committing batch is the last one active - /// and no lookahead can extend past it — so script-carrying commits share - /// one walk of the stored history instead of walking it once per commit. - /// Deliberately survives `reset_for_rescan`: the restarted scan covers - /// heights above its entry point, while these scripts still owe a pass - /// over the committed prefix below it. - backward_scripts: HashMap>, - /// Number of committed-range sweeps that reached the chunk walk in - /// `rescan_committed_range`. Diagnostic counter; the sweep-coalescing - /// regression test asserts on it. - pub(super) committed_range_sweeps: u64, } impl @@ -150,8 +135,6 @@ impl= self.progress.filter_header_tip_height() && self.progress.committed_height() >= self.progress.target_height() { - // A block delivered after its batch committed derives scripts with - // no active batch to route them to, so they land in the accumulator - // instead. It is in-memory only and nothing looks below the - // committed frontier again, so sweep here rather than wait for a - // next commit that may never come. - if !self.backward_scripts.is_empty() { - let backward_scripts = std::mem::take(&mut self.backward_scripts); - let sweep_start = self.progress.committed_height().saturating_add(1); - events.extend(self.rescan_committed_range(sweep_start, &backward_scripts).await?); - } - - // Blocks that sweep found are charged to no batch, so the commit - // gate cannot hold completion — gate on the tracker. Their - // `BlockProcessed` re-enters here, and a round deriving no new - // scripts is the fixpoint. - if self.tracker.has_blocks_in_flight() { - return Ok(events); - } - // Blocks applied after the last commit leave processed records that // no commit will prune. self.tracker.prune_at_or_below(self.progress.committed_height()); @@ -620,89 +579,11 @@ impl = self - .active_batches - .iter() - .filter(|(&start, batch)| start > batch_start && batch.scanned()) - .map(|(&start, _)| start) - .collect(); - - for later_start in later_batches { - events.extend(self.rescan_batch(later_start, &scripts_by_wallet).await?); - } - - // Newly derived scripts also have to reach ranges that - // already committed: those blocks were matched against a - // watch set that predates these scripts, and nothing else - // ever looks below `committed_height` again (#846). That - // backward sweep walks stored history — the expensive - // direction — so defer it: accumulate the scripts at the - // manager level, across batch commits, and sweep below. - for (wallet_id, scripts) in scripts_by_wallet { - self.backward_scripts.entry(wallet_id).or_default().extend(scripts); - } - if let Some(batch) = self.active_batches.get(&batch_start) { - if batch.pending_blocks() > 0 { - // Forward rescan found blocks; converge the - // forward direction first. - break; - } - } - } - - events.extend(self.reconcile_untested_scripts(batch_start).await?); - if let Some(batch) = self.active_batches.get(&batch_start) { - if batch.pending_blocks() > 0 { - // Reconciliation found blocks; converge before committing. - break; - } - } - - // The backward sweep waits for the forward pipeline to drain: - // it runs only when this batch is the last one active and no - // lookahead batch can be created past it. Commits before that - // point leave the accumulated scripts in place, so a sync's - // script-carrying commits share one walk of the stored - // history — plus one walk per follow-up round whose block - // processing derives genuinely new scripts — instead of - // walking it once per commit. Hits attribute to this batch, - // so scripts their processing derives re-enter through - // `collected_scripts` above and only genuinely new scripts - // get a follow-up sweep. - let forward_drained = self.active_batches.len() == 1 - && self.active_batches.get(&batch_start).is_some_and(|b| { - b.end_height() >= self.progress.filter_header_tip_height() - }); - if forward_drained && !self.backward_scripts.is_empty() { - let backward_scripts = std::mem::take(&mut self.backward_scripts); - events - .extend(self.rescan_committed_range(batch_start, &backward_scripts).await?); - - // Check if the backward sweep found more blocks - if let Some(batch) = self.active_batches.get(&batch_start) { - if batch.pending_blocks() > 0 { - // Found more blocks, can't commit yet - break; - } - } - } - // Mark rescan as complete - if let Some(batch) = self.active_batches.get_mut(&batch_start) { - batch.mark_rescan_complete(); + events.extend(self.reconcile_untested_scripts(batch_start).await?); + if let Some(batch) = self.active_batches.get(&batch_start) { + if batch.pending_blocks() > 0 { + // Reconciliation found blocks; converge before committing. + break; } } @@ -856,49 +737,9 @@ impl>, - ) -> usize { - let target = self - .active_batches - .range(..=height) - .next_back() - .filter(|(_, batch)| batch.end_height() >= height) - .map(|(&start, _)| start); - - let mut routed = 0; - for (wallet_id, scripts) in new_scripts { - if scripts.is_empty() { - continue; - } - routed += scripts.len(); - match target.and_then(|start| self.active_batches.get_mut(&start)) { - Some(batch) => batch.add_scripts_for_wallet(*wallet_id, scripts.iter().cloned()), - // Its range has committed: only the sweep can still test these. - None => self - .backward_scripts - .entry(*wallet_id) - .or_default() - .extend(scripts.iter().cloned()), - } - } - routed - } - /// Re-test anything the wallet watches that this batch has never been - /// matched against, before letting it commit. Closing the loop on wallet - /// state does not depend on a `new_scripts` notification arriving; a - /// script the wallet has dropped from `scan_script_pubkeys_for` is still - /// not re-tested. + /// matched against, before letting it commit. A script the wallet has + /// dropped from `scan_script_pubkeys_for` is not re-tested. async fn reconcile_untested_scripts(&mut self, batch_start: u32) -> SyncResult> { let Some(batch) = self.active_batches.get(&batch_start) else { return Ok(vec![]); @@ -933,11 +774,7 @@ impl(), ); // Same accounting as any other newly derived script. - let events = self.rescan_batch(batch_start, &untested).await?; - for (wallet_id, scripts) in untested { - self.backward_scripts.entry(wallet_id).or_default().extend(scripts); - } - Ok(events) + self.rescan_batch(batch_start, &untested).await } /// Rescan a specific batch for newly discovered scriptPubKeys, attributed @@ -1012,26 +849,6 @@ impl>, - context: &str, - ) -> Vec { let mut events = Vec::new(); let mut blocks_needed: BTreeMap> = BTreeMap::new(); let mut new_blocks_count = 0; @@ -1062,7 +879,7 @@ impl>, - ) -> SyncResult> { - let Some(range_end) = batch_start.checked_sub(1) else { - return Ok(vec![]); - }; - - let wallet_queries: Vec<(WalletId, Vec)> = new_scripts - .iter() - .filter(|(_, scripts)| !scripts.is_empty()) - .map(|(id, scripts)| (*id, scripts.iter().cloned().collect())) - .collect(); - if wallet_queries.is_empty() { - return Ok(vec![]); - } - - // Nothing relevant can precede the earliest wallet birth height, and - // nothing is loadable below the first stored filter. - let wallet_base = self.wallet.read().await.earliest_required_height().await; - let Some(filter_base) = self.filter_storage.read().await.filter_start_height().await else { - return Ok(vec![]); - }; - let range_start = wallet_base.max(filter_base); - if range_start > range_end { - return Ok(vec![]); - } - - self.committed_range_sweeps += 1; - tracing::info!( - "Rescan committed filters ({}-{}) for new scripts across {} wallets (sweep #{})", - range_start, - range_end, - wallet_queries.len(), - self.committed_range_sweeps - ); - - let mut block_to_wallets: BTreeMap> = BTreeMap::new(); - let mut chunk_start = range_start; - while chunk_start <= range_end { - let chunk_end = (chunk_start + BATCH_PROCESSING_SIZE - 1).min(range_end); - // A chunk the storage cannot serve (e.g. filters pruned or never - // stored for a sub-range) is skipped rather than failing the - // commit: the sweep is best-effort recovery over whatever - // history is locally available. - let filters = match self.load_filters(chunk_start, chunk_end).await { - Ok(filters) => filters, - Err(e) => { - tracing::warn!( - "Committed-range rescan skipping {}-{}: {}", - chunk_start, - chunk_end, - e - ); - chunk_start = chunk_end + 1; - continue; - } - }; - for (wallet_id, scripts) in &wallet_queries { - let matches = check_compact_filters_for_elements(&filters, scripts, &[], 0); - for key in matches { - block_to_wallets.entry(key).or_default().insert(*wallet_id); - } - } - chunk_start = chunk_end + 1; - } - - Ok(self.queue_new_script_matches(batch_start, block_to_wallets, "Committed-range rescan")) - } - /// Handle notification that new filter headers are available. /// Used by both FilterHeadersSyncComplete and FilterHeadersStored events. pub(super) async fn handle_new_filter_headers( @@ -1831,7 +1552,6 @@ mod tests { let mut batch1 = FiltersBatch::new(0, 4999, HashMap::new()); batch1.set_pending_blocks(0); batch1.mark_scanned(); - batch1.mark_rescan_complete(); manager.active_batches.insert(0, batch1); @@ -1858,7 +1578,6 @@ mod tests { let mut batch1 = FiltersBatch::new(0, 4999, HashMap::new()); batch1.set_pending_blocks(0); batch1.mark_scanned(); - batch1.mark_rescan_complete(); batch1.set_scanned_wallets(BTreeMap::from([(MOCK_WALLET_ID, 0)])); manager.active_batches.insert(0, batch1); @@ -1873,7 +1592,6 @@ mod tests { let mut batch2 = FiltersBatch::new(5000, 9999, HashMap::new()); batch2.set_pending_blocks(0); batch2.mark_scanned(); - batch2.mark_rescan_complete(); manager.active_batches.insert(5000, batch2); manager.try_commit_batches().await.unwrap(); @@ -1902,7 +1620,6 @@ mod tests { let mut batch = FiltersBatch::new(0, 4999, HashMap::new()); batch.set_pending_blocks(0); batch.mark_scanned(); - batch.mark_rescan_complete(); batch.set_scanned_wallets(BTreeMap::from([(wallet_a, 0)])); manager.active_batches.insert(0, batch); @@ -1940,7 +1657,6 @@ mod tests { let mut batch = FiltersBatch::new(5000, 9999, HashMap::new()); batch.set_pending_blocks(0); batch.mark_scanned(); - batch.mark_rescan_complete(); batch.set_scanned_wallets(BTreeMap::from([(wallet_a, 0)])); manager.active_batches.insert(5000, batch); @@ -1981,7 +1697,6 @@ mod tests { let mut batch = FiltersBatch::new(205_000, 209_999, HashMap::new()); batch.set_pending_blocks(0); batch.mark_scanned(); - batch.mark_rescan_complete(); batch.set_scanned_wallets(BTreeMap::from([(wallet_b, 0)])); manager.active_batches.insert(205_000, batch); @@ -2010,7 +1725,6 @@ mod tests { let mut batch1 = FiltersBatch::new(0, 4999, HashMap::new()); batch1.set_pending_blocks(0); batch1.mark_scanned(); - batch1.mark_rescan_complete(); batch1.set_scanned_wallets(BTreeMap::from([(wallet_a, 0)])); manager.active_batches.insert(0, batch1); @@ -2022,7 +1736,6 @@ mod tests { let mut batch2 = FiltersBatch::new(5000, 9999, HashMap::new()); batch2.set_pending_blocks(0); batch2.mark_scanned(); - batch2.mark_rescan_complete(); batch2.set_scanned_wallets(BTreeMap::from([(wallet_a, 0)])); manager.active_batches.insert(5000, batch2); @@ -2068,7 +1781,6 @@ mod tests { let mut batch = FiltersBatch::new(5000, 9999, HashMap::new()); batch.set_pending_blocks(0); batch.mark_scanned(); - batch.mark_rescan_complete(); batch.set_scanned_wallets(BTreeMap::from([(wallet_a, 0), (wallet_b, 0)])); manager.active_batches.insert(5000, batch); @@ -2118,7 +1830,6 @@ mod tests { let mut batch = FiltersBatch::new(5000, 9999, HashMap::new()); batch.set_pending_blocks(0); batch.mark_scanned(); - batch.mark_rescan_complete(); batch.set_scanned_wallets(BTreeMap::from([(wallet_a, 0)])); manager.active_batches.insert(5000, batch); @@ -2175,7 +1886,6 @@ mod tests { let mut batch = FiltersBatch::new(1000, 5999, HashMap::new()); batch.set_pending_blocks(0); batch.mark_scanned(); - batch.mark_rescan_complete(); batch.set_scanned_wallets(BTreeMap::from([(wallet_a, 0)])); manager.active_batches.insert(1000, batch); @@ -2205,7 +1915,6 @@ mod tests { let mut batch = FiltersBatch::new(200_000, 204_999, HashMap::new()); batch.set_pending_blocks(0); batch.mark_scanned(); - batch.mark_rescan_complete(); batch.set_scanned_wallets(BTreeMap::from([(wallet_b, 0)])); manager.active_batches.insert(200_000, batch); @@ -2243,7 +1952,6 @@ mod tests { let mut batch = FiltersBatch::new(0, 4999, HashMap::new()); batch.set_pending_blocks(0); batch.mark_scanned(); - batch.mark_rescan_complete(); batch.set_scanned_wallets(BTreeMap::from([(wallet_a, 3)])); manager.active_batches.insert(0, batch); @@ -2612,7 +2320,6 @@ mod tests { // Mark batch ready so commit can run, then commit. if let Some(b) = manager.active_batches.get_mut(&0) { b.set_pending_blocks(0); - b.mark_rescan_complete(); } manager.try_commit_batches().await.unwrap(); @@ -2806,7 +2513,6 @@ mod tests { let mut batch = FiltersBatch::new(0, 4999, HashMap::new()); batch.set_pending_blocks(0); batch.mark_scanned(); - batch.mark_rescan_complete(); manager.active_batches.insert(0, batch); manager.try_commit_batches().await.unwrap(); @@ -2932,12 +2638,10 @@ mod tests { let mut batch1 = FiltersBatch::new(0, 4999, HashMap::new()); batch1.set_pending_blocks(0); batch1.mark_scanned(); - batch1.mark_rescan_complete(); let mut batch2 = FiltersBatch::new(5000, 9999, HashMap::new()); batch2.set_pending_blocks(0); batch2.mark_scanned(); - batch2.mark_rescan_complete(); manager.active_batches.insert(5000, batch2); // Insert higher one first manager.active_batches.insert(0, batch1); @@ -3063,73 +2767,6 @@ mod tests { assert!(manager.is_idle()); } - #[tokio::test] - async fn test_batch_collects_scripts() { - use crate::sync::filters::batch::FiltersBatch; - use dashcore::Network; - - let mut batch = FiltersBatch::new(0, 4999, HashMap::new()); - - // Initially empty - assert!(batch.take_collected_scripts().is_empty()); - - // Add scripts using test utility - let script1 = dashcore::Address::dummy(Network::Testnet, 1).script_pubkey(); - let script2 = dashcore::Address::dummy(Network::Testnet, 2).script_pubkey(); - let wallet_id: WalletId = [7; 32]; - - batch.add_scripts_for_wallet(wallet_id, [script1.clone(), script2.clone()]); - - let collected = batch.take_collected_scripts(); - let for_wallet = collected.get(&wallet_id).expect("wallet entry"); - assert_eq!(for_wallet.len(), 2); - assert!(for_wallet.contains(&script1)); - assert!(for_wallet.contains(&script2)); - - // After take, should be empty - assert!(batch.take_collected_scripts().is_empty()); - } - - /// Proves that scripts from a block with no in-flight record — every - /// delivery after the first — reach the batch covering their height, and - /// the backward accumulator when no active batch does. - #[tokio::test] - async fn test_new_scripts_route_by_height_with_no_in_flight_record() { - use crate::sync::filters::batch::FiltersBatch; - use dashcore::Network; - - let mut manager = create_test_manager().await; - manager.active_batches.insert(0, FiltersBatch::new(0, 4999, HashMap::new())); - manager.active_batches.insert(5000, FiltersBatch::new(5000, 9999, HashMap::new())); - - let wallet_id: WalletId = [7; 32]; - let covered = Address::dummy(Network::Testnet, 1).script_pubkey(); - let below = Address::dummy(Network::Testnet, 2).script_pubkey(); - - let (tx, _rx) = unbounded_channel(); - let requests = RequestSender::new(tx); - let block_processed = |height: u32, script: &ScriptBuf| SyncEvent::BlockProcessed { - block_hash: Header::dummy(height).block_hash(), - height, - wallets: BTreeSet::from([wallet_id]), - new_scripts: BTreeMap::from([(wallet_id, vec![script.clone()])]), - confirmed_txids: vec![], - }; - - manager.handle_sync_event(&block_processed(6000, &covered), &requests).await.unwrap(); - - let lower = manager.active_batches.get_mut(&0).unwrap().take_collected_scripts(); - assert!(lower.is_empty(), "scripts must not land in a batch that cannot match them"); - let covering = manager.active_batches.get_mut(&5000).unwrap().take_collected_scripts(); - assert!(covering.get(&wallet_id).is_some_and(|s| s.contains(&covered))); - assert!(manager.backward_scripts.is_empty()); - - manager.active_batches.remove(&0); - manager.handle_sync_event(&block_processed(120, &below), &requests).await.unwrap(); - assert!(manager.active_batches.get_mut(&5000).unwrap().take_collected_scripts().is_empty()); - assert!(manager.backward_scripts.get(&wallet_id).is_some_and(|s| s.contains(&below))); - } - /// Proves a batch with no filters still marks the scripts tested, so /// reconciliation does not hand it the same set on every commit attempt. #[tokio::test] @@ -3208,89 +2845,6 @@ mod tests { ); } - /// Proves the tip is not reported synced while scripts sit in the backward - /// accumulator with no batch left to sweep them: the sweep runs here, and - /// completion waits for the block it finds. - #[tokio::test] - async fn test_tip_sweeps_stranded_backward_scripts_before_completing() { - let wallet_id: WalletId = [3; 32]; - let watched = Address::dummy(Network::Testnet, 31); - - let mut multi = MultiMockWallet::new(); - multi.insert_wallet( - wallet_id, - MockWalletState { - addresses: vec![watched.clone()], - synced_height: 9, - last_processed_height: 9, - account_generation: 0, - }, - ); - let mut manager = create_multi_test_manager(Arc::new(RwLock::new(multi))).await; - manager.set_state(SyncState::Syncing); - - // Every height at or below the committed frontier has its header and - // filter persisted; only height 4 pays `watched`. - let paying = Block::dummy(4, vec![Transaction::dummy(&watched, 0..0, &[4])]); - let paying_filter = BlockFilter::dummy(&paying); - let key = FilterMatchKey::new(4, paying.block_hash()); - { - let mut header_storage = manager.header_storage.write().await; - let mut filter_storage = manager.filter_storage.write().await; - for height in 0..=9u32 { - let (header, bytes) = if height == 4 { - (paying.header, paying_filter.content.clone()) - } else { - let filler = Block::dummy(height, vec![]); - (filler.header, BlockFilter::dummy(&filler).content) - }; - header_storage.store_headers_at_height(&[header.into()], height).await.unwrap(); - filter_storage.store_filter(height, &bytes).await.unwrap(); - } - } - - // At the tip with nothing active: the last commit left - // `processing_height` past the tip, so no lookahead batch is created - // and the completion branch is reached. - manager.processing_height = 10; - manager.progress.update_stored_height(9); - manager.progress.update_committed_height(9); - manager.progress.update_filter_header_tip_height(9); - manager.progress.update_target_height(9); - manager.backward_scripts.entry(wallet_id).or_default().insert(watched.script_pubkey()); - - let events = manager.try_process_batch().await.unwrap(); - assert!( - events.iter().any(|e| match e { - SyncEvent::BlocksNeeded { - blocks, - } => blocks.contains_key(&key), - _ => false, - }), - "the tip sweep must find the block paying the stranded script" - ); - assert!( - !events.iter().any(|e| matches!(e, SyncEvent::FiltersSyncComplete { .. })), - "completion must wait for that block to be applied" - ); - assert!(manager.backward_scripts.is_empty(), "the sweep consumes the accumulator"); - - let (tx, _rx) = unbounded_channel(); - let requests = RequestSender::new(tx); - let processed = SyncEvent::BlockProcessed { - block_hash: paying.block_hash(), - height: 4, - wallets: BTreeSet::from([wallet_id]), - new_scripts: BTreeMap::new(), - confirmed_txids: vec![], - }; - let events = manager.handle_sync_event(&processed, &requests).await.unwrap(); - assert!( - events.iter().any(|e| matches!(e, SyncEvent::FiltersSyncComplete { .. })), - "completion follows once the sweep's block is applied" - ); - } - #[tokio::test] async fn test_start_download_waits_when_filter_headers_insufficient() { let mut manager = create_test_manager().await; diff --git a/dash-spv/src/sync/filters/sync_manager.rs b/dash-spv/src/sync/filters/sync_manager.rs index 8e763b8b3..3a3b9cc69 100644 --- a/dash-spv/src/sync/filters/sync_manager.rs +++ b/dash-spv/src/sync/filters/sync_manager.rs @@ -174,7 +174,6 @@ impl< block_hash, height, wallets, - new_scripts, .. } => { // Record per-wallet processing so a future scan can give a @@ -183,8 +182,7 @@ impl< self.tracker.record_processed(*height, *block_hash, wallets); // Check if this block is part of our tracked blocks - let in_flight = self.tracker.finish_in_flight(block_hash); - if let Some((_, batch_start)) = in_flight { + if let Some((_, batch_start)) = self.tracker.finish_in_flight(block_hash) { if let Some(batch) = self.active_batches.get_mut(&batch_start) { batch.decrement_pending_blocks(); tracing::debug!( @@ -195,13 +193,6 @@ impl< batch.pending_blocks() ); } - } - - // Outside the in-flight arm on purpose: that record is consumed - // by the first delivery, and a block is delivered more than once. - let derived = self.collect_new_scripts(*height, new_scripts); - - if in_flight.is_some() || derived > 0 { return self.try_process_batch().await; } }