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
3 changes: 2 additions & 1 deletion dash-spv-ffi/include/dash_spv_ffi.h
Original file line number Diff line number Diff line change
Expand Up @@ -85,7 +85,8 @@ typedef struct FFIFilterHeadersProgress {
*/
typedef struct FFIFiltersProgress {
enum FFISyncState state;
uint32_t current_height;
uint32_t committed_height;
uint32_t stored_height;
uint32_t target_height;
uint32_t filter_header_tip_height;
uint32_t downloaded;
Expand Down
2 changes: 1 addition & 1 deletion dash-spv-ffi/src/bin/ffi_cli.rs
Original file line number Diff line number Diff line change
Expand Up @@ -221,7 +221,7 @@ extern "C" fn on_progress_update(progress: *const FFISyncProgress, _user_data: *
}
if !p.filters.is_null() {
let f = unsafe { &*p.filters };
print!("filters:{}/{} ", f.current_height, f.target_height);
print!("filters:{}/{} stored: {} ", f.committed_height, f.target_height, f.stored_height);
}
if !p.blocks.is_null() {
let f = unsafe { &*p.blocks };
Expand Down
6 changes: 4 additions & 2 deletions dash-spv-ffi/src/types.rs
Original file line number Diff line number Diff line change
Expand Up @@ -131,7 +131,8 @@ impl From<&FilterHeadersProgress> for FFIFilterHeadersProgress {
#[derive(Debug, Clone, Default)]
pub struct FFIFiltersProgress {
pub state: FFISyncState,
pub current_height: u32,
pub committed_height: u32,
pub stored_height: u32,
pub target_height: u32,
pub filter_header_tip_height: u32,
pub downloaded: u32,
Expand All @@ -145,7 +146,8 @@ impl From<&FiltersProgress> for FFIFiltersProgress {
fn from(progress: &FiltersProgress) -> Self {
FFIFiltersProgress {
state: progress.state().into(),
current_height: progress.current_height(),
committed_height: progress.committed_height(),
stored_height: progress.stored_height(),
target_height: progress.target_height(),
filter_header_tip_height: progress.filter_header_tip_height(),
downloaded: progress.downloaded(),
Expand Down
6 changes: 4 additions & 2 deletions dash-spv-ffi/tests/test_types.rs
Original file line number Diff line number Diff line change
Expand Up @@ -65,7 +65,8 @@ mod tests {

let mut filters = FiltersProgress::default();
filters.set_state(SyncState::WaitForEvents);
filters.update_current_height(120);
filters.update_stored_height(150);
filters.update_committed_height(120);
filters.update_target_height(200);
filters.update_filter_header_tip_height(150);
filters.add_downloaded(40);
Expand Down Expand Up @@ -138,7 +139,8 @@ mod tests {
unsafe {
let filters = &*ffi_progress.filters;
assert_eq!(filters.state, FFISyncState::WaitForEvents);
assert_eq!(filters.current_height, 120);
assert_eq!(filters.stored_height, 150);
assert_eq!(filters.committed_height, 120);
assert_eq!(filters.target_height, 200);
assert_eq!(filters.filter_header_tip_height, 150);
assert_eq!(filters.downloaded, 40);
Expand Down
59 changes: 28 additions & 31 deletions dash-spv/src/sync/filters/manager.rs
Original file line number Diff line number Diff line change
Expand Up @@ -64,8 +64,6 @@ pub struct FiltersManager<
// === Multi-batch processing state ===
/// Active batches being processed (keyed by start_height).
pub(super) active_batches: BTreeMap<u32, FiltersBatch>,
/// Height that has been committed to wallet (all blocks up to this height processed).
pub(super) committed_height: u32,
/// Current block height being processed (for progress tracking).
processing_height: u32,
/// Blocks remaining that need to be processed.
Expand Down Expand Up @@ -96,7 +94,6 @@ impl<H: BlockHeaderStorage, FH: FilterHeaderStorage, F: FilterStorage, W: Wallet
next_batch_to_store: 0,
// Multi-batch processing
active_batches: BTreeMap::new(),
committed_height: 0,
processing_height: 0,
blocks_remaining: BTreeMap::new(),
filters_matched: HashSet::new(),
Expand Down Expand Up @@ -159,11 +156,11 @@ impl<H: BlockHeaderStorage, FH: FilterHeaderStorage, F: FilterStorage, W: Wallet
if scan_start > self.progress.filter_header_tip_height() {
// Only emit FiltersSyncComplete if we've also reached the chain tip
// This prevents premature sync complete while filter headers are still syncing
if self.committed_height >= self.progress.target_height() {
if self.progress.committed_height() >= self.progress.target_height() {
self.set_state(SyncState::Synced);
tracing::info!("Filters already synced to {}", self.progress.target_height());
return Ok(vec![SyncEvent::FiltersSyncComplete {
tip_height: self.committed_height,
tip_height: self.progress.committed_height(),
}]);
}
// Caught up to available filter headers but chain tip not reached yet
Expand Down Expand Up @@ -245,8 +242,8 @@ impl<H: BlockHeaderStorage, FH: FilterHeaderStorage, F: FilterStorage, W: Wallet
scan_start,
end_height
);
// Update current_height to reflect stored filters are available
self.progress.update_current_height(stored_filters_tip);
// Update stored_height to reflect stored filters are available
self.progress.update_stored_height(stored_filters_tip);
self.load_filters(scan_start, end_height).await?
} else {
HashMap::new()
Expand All @@ -257,17 +254,17 @@ impl<H: BlockHeaderStorage, FH: FilterHeaderStorage, F: FilterStorage, W: Wallet
batch.mark_verified();
}
self.active_batches.insert(scan_start, batch);
self.committed_height = scan_start.saturating_sub(1);
self.progress.update_committed_height(scan_start.saturating_sub(1));

// Only scan if all filters for the batch are already loaded
if self.progress.current_height() >= batch_end {
if self.progress.stored_height() >= batch_end {
self.scan_batch(scan_start).await
} else {
tracing::debug!(
"Initial batch {}-{}: waiting for filters (current_height={})",
"Initial batch {}-{}: waiting for filters (stored_height={})",
scan_start,
batch_end,
self.progress.current_height()
self.progress.stored_height()
);
Ok(vec![])
}
Expand Down Expand Up @@ -386,16 +383,16 @@ impl<H: BlockHeaderStorage, FH: FilterHeaderStorage, F: FilterStorage, W: Wallet
}

self.progress.add_processed(batch.end_height() - batch.start_height() + 1);
self.progress.update_current_height(batch.end_height());
self.progress.update_stored_height(batch.end_height());
self.next_batch_to_store = batch.end_height() + 1;
}

// If we stored any batches, try to process the batch containing the current processing height.
// This is called only when batches complete, not on every filter
if !events.is_empty() {
tracing::debug!(
"Calling try_process_batch after storing batches (current_height={}, target_height={})",
self.progress.current_height(),
"Calling try_process_batch after storing batches (stored_height={}, target_height={})",
self.progress.stored_height(),
self.progress.target_height()
);
events.extend(self.try_process_batch().await?);
Expand Down Expand Up @@ -423,15 +420,15 @@ impl<H: BlockHeaderStorage, FH: FilterHeaderStorage, F: FilterStorage, W: Wallet
// updates (already Synced, signal BlocksManager that no more blocks are coming).
if self.active_batches.is_empty()
&& matches!(self.state(), SyncState::Syncing | SyncState::Synced)
&& self.committed_height >= self.progress.filter_header_tip_height()
&& self.committed_height >= self.progress.target_height()
&& self.progress.committed_height() >= self.progress.filter_header_tip_height()
&& self.progress.committed_height() >= self.progress.target_height()
{
if self.state() == SyncState::Syncing {
self.set_state(SyncState::Synced);
}
tracing::info!("Filter sync complete at height {}", self.committed_height);
tracing::info!("Filter sync complete at height {}", self.progress.committed_height());
events.push(SyncEvent::FiltersSyncComplete {
tip_height: self.committed_height,
tip_height: self.progress.committed_height(),
});
}

Expand Down Expand Up @@ -500,8 +497,8 @@ impl<H: BlockHeaderStorage, FH: FilterHeaderStorage, F: FilterStorage, W: Wallet
// Commit this batch
let batch = self.active_batches.remove(&batch_start).unwrap();
let end = batch.end_height();
if end > self.committed_height {
self.committed_height = end;
if end > self.progress.committed_height() {
self.progress.update_committed_height(end);
self.wallet.write().await.update_filter_committed_height(end);
}
self.processing_height = end + 1;
Expand All @@ -510,7 +507,7 @@ impl<H: BlockHeaderStorage, FH: FilterHeaderStorage, F: FilterStorage, W: Wallet
"Committed batch {}-{}, committed_height now {}",
batch.start_height(),
batch.end_height(),
self.committed_height
self.progress.committed_height()
);
}

Expand All @@ -526,7 +523,7 @@ impl<H: BlockHeaderStorage, FH: FilterHeaderStorage, F: FilterStorage, W: Wallet
.active_batches
.iter()
.filter(|(_, batch)| {
!batch.scanned() && self.progress.current_height() >= batch.end_height()
!batch.scanned() && self.progress.stored_height() >= batch.end_height()
})
.map(|(&start, _)| start)
.collect();
Expand Down Expand Up @@ -566,21 +563,21 @@ impl<H: BlockHeaderStorage, FH: FilterHeaderStorage, F: FilterStorage, W: Wallet
);

// Load available filters into the new batch
let available_end = self.progress.current_height().min(next_end);
let available_end = self.progress.stored_height().min(next_end);
let filters = if next_start <= available_end {
self.load_filters(next_start, available_end).await?
} else {
HashMap::new()
};

let mut batch = FiltersBatch::new(next_start, next_end, filters);
if self.progress.current_height() >= next_end {
if self.progress.stored_height() >= next_end {
batch.mark_verified();
}
self.active_batches.insert(next_start, batch);

// Scan immediately if filters are available
if self.progress.current_height() >= next_end {
if self.progress.stored_height() >= next_end {
events.extend(self.scan_batch(next_start).await?);
}
}
Expand Down Expand Up @@ -736,7 +733,7 @@ impl<H: BlockHeaderStorage, FH: FilterHeaderStorage, F: FilterStorage, W: Wallet

match self.state() {
SyncState::Syncing | SyncState::Synced
if self.progress.current_height() < self.progress.filter_header_tip_height() =>
if self.progress.stored_height() < self.progress.filter_header_tip_height() =>
{
self.filter_pipeline.extend_target(tip_height);
{
Expand All @@ -750,7 +747,7 @@ impl<H: BlockHeaderStorage, FH: FilterHeaderStorage, F: FilterStorage, W: Wallet
}
}
SyncState::WaitingForConnections | SyncState::WaitForEvents
if self.progress.current_height() < self.progress.filter_header_tip_height() =>
if self.progress.stored_height() < self.progress.filter_header_tip_height() =>
{
return self.start_download(requests).await;
}
Expand Down Expand Up @@ -809,7 +806,7 @@ mod tests {
async fn test_filters_manager_progress() {
let mut manager = create_test_manager().await;
manager.set_state(SyncState::Syncing);
manager.progress.update_current_height(500);
manager.progress.update_stored_height(500);
manager.progress.update_target_height(1000);
manager.progress.add_processed(350);
manager.progress.add_downloaded(250);
Expand All @@ -819,7 +816,7 @@ mod tests {
let progress = manager_ref.progress();
if let SyncManagerProgress::Filters(progress) = progress {
assert_eq!(progress.state(), SyncState::Syncing);
assert_eq!(progress.current_height(), 500);
assert_eq!(progress.stored_height(), 500);
assert_eq!(progress.target_height(), 1000);
assert_eq!(progress.processed(), 350);
assert_eq!(progress.downloaded(), 250);
Expand Down Expand Up @@ -874,7 +871,7 @@ mod tests {
// Commit should work
manager.try_commit_batches().await.unwrap();
assert_eq!(manager.active_batches.len(), 0);
assert_eq!(manager.committed_height, 4999);
assert_eq!(manager.progress.committed_height(), 4999);
}

#[tokio::test]
Expand All @@ -899,7 +896,7 @@ mod tests {
// Commit should commit both in order
manager.try_commit_batches().await.unwrap();
assert_eq!(manager.active_batches.len(), 0);
assert_eq!(manager.committed_height, 9999); // Both committed
assert_eq!(manager.progress.committed_height(), 9999); // Both committed
}

#[tokio::test]
Expand Down
38 changes: 29 additions & 9 deletions dash-spv/src/sync/filters/progress.rs
Original file line number Diff line number Diff line change
Expand Up @@ -8,8 +8,10 @@ use std::time::Instant;
pub struct FiltersProgress {
/// Current sync state.
state: SyncState,
/// Tip height of the filter storage.
current_height: u32,
/// Height up to which all filter batches have been fully committed to wallet.
committed_height: u32,
/// Tip height of the filter storage (how far filters have been downloaded/stored).
stored_height: u32,
/// Target height (peer's best height). Used for progress display.
target_height: u32,
/// The tip height of the filter header storage (the download limit for filters).
Expand All @@ -29,7 +31,8 @@ impl Default for FiltersProgress {
fn default() -> Self {
Self {
state: Default::default(),
current_height: 0,
committed_height: 0,
stored_height: 0,
target_height: 0,
filter_header_tip_height: 0,
downloaded: 0,
Expand All @@ -46,6 +49,16 @@ impl FiltersProgress {
self.state
}

/// Get the committed height (all batches up to this height fully processed).
pub fn committed_height(&self) -> u32 {
self.committed_height
}

/// Get the stored height (tip of filter storage).
pub fn stored_height(&self) -> u32 {
self.stored_height
}

/// Get the filter header tip height (the download limit for filters).
pub fn filter_header_tip_height(&self) -> u32 {
self.filter_header_tip_height
Expand Down Expand Up @@ -77,9 +90,15 @@ impl FiltersProgress {
self.bump_last_activity();
}

/// Update the current height (last successfully processed height).
pub fn update_current_height(&mut self, height: u32) {
self.current_height = height;
/// Update the committed height (all batches up to this height fully processed).
pub fn update_committed_height(&mut self, height: u32) {
self.committed_height = height;
self.bump_last_activity();
}

/// Update the stored height (tip of filter storage).
pub fn update_stored_height(&mut self, height: u32) {
self.stored_height = height;
self.bump_last_activity();
}

Expand Down Expand Up @@ -127,11 +146,12 @@ impl fmt::Display for FiltersProgress {
let pct = self.percentage() * 100.0;
write!(
f,
"{:?} {}/{} ({:.1}%) downloaded: {}, processed: {}, matched: {}, last_activity: {}s",
"{:?} {}/{} ({:.1}%) stored:{}, downloaded: {}, processed: {}, matched: {}, last_activity: {}s",
self.state,
self.current_height,
self.committed_height,
self.target_height,
pct,
self.stored_height,
self.downloaded,
self.processed,
self.matched,
Expand All @@ -145,6 +165,6 @@ impl ProgressPercentage for FiltersProgress {
self.target_height
}
fn current_height(&self) -> u32 {
self.current_height
self.committed_height
}
}
Loading
Loading