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
12 changes: 12 additions & 0 deletions src/builder.rs
Original file line number Diff line number Diff line change
Expand Up @@ -123,6 +123,18 @@ impl Builder {
self
}

/// Skip filter header and filter downloads during the initial sync. Once the node is
/// believed to be at the current chain tip, filters will be downloaded for newly mined
/// blocks. This feature is intended for new wallets that have never received a payment.
///
/// A subsequent call to
/// [`Requester::rescan`](crate::Requester::rescan) or
/// [`Requester::rescan_from`](crate::Requester::rescan_from) resumes normal filter syncing.
pub fn headers_only_sync(mut self) -> Self {
self.config.headers_only_sync = true;
self
}

/// Route network traffic through a Tor daemon using a Socks5 proxy. Currently, proxies
/// must be reachable by IP address.
pub fn socks5_proxy(mut self, proxy: impl Into<Socks5Proxy>) -> Self {
Expand Down
2 changes: 2 additions & 0 deletions src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -381,6 +381,7 @@ struct Config {
peer_timeout_config: PeerTimeoutConfig,
filter_type: FilterType,
block_type: BlockType,
headers_only_sync: bool,
}

impl Default for Config {
Expand All @@ -394,6 +395,7 @@ impl Default for Config {
peer_timeout_config: PeerTimeoutConfig::default(),
filter_type: FilterType::default(),
block_type: BlockType::default(),
headers_only_sync: Default::default(),
}
}
}
Expand Down
53 changes: 39 additions & 14 deletions src/node.rs
Original file line number Diff line number Diff line change
Expand Up @@ -59,6 +59,7 @@ pub struct Node {
block_queue: BlockQueue,
client_recv: UnboundedReceiver<ClientMessage>,
peer_recv: Receiver<PeerThreadMessage>,
headers_only_sync: bool,
}

impl Node {
Expand All @@ -72,6 +73,7 @@ impl Node {
peer_timeout_config,
filter_type,
block_type,
headers_only_sync,
} = config;
// Set up a communication channel between the node and client
let (info_tx, info_rx) = mpsc::channel::<Info>(32);
Expand Down Expand Up @@ -116,6 +118,7 @@ impl Node {
block_queue: BlockQueue::new(),
client_recv: crx,
peer_recv: mrx,
headers_only_sync,
},
client,
)
Expand Down Expand Up @@ -338,21 +341,20 @@ impl Node {
// This state is updated upon receiving new block headers
NodeState::Behind => (),
NodeState::HeadersSynced => {
if self.chain.is_cf_headers_synced() {
if self.headers_only_sync {
self.state = NodeState::FiltersSynced;
self.headers_only_sync = false;
let tip = self.chain.header_chain.height();
self.chain.header_chain.assume_checked_to(tip);
self.emit_sync_complete();
} else if self.chain.is_cf_headers_synced() {
self.state = NodeState::FilterHeadersSynced;
}
}
NodeState::FilterHeadersSynced => {
if self.chain.is_filters_synced() {
self.state = NodeState::FiltersSynced;
let update = SyncUpdate::new(
HashCheckpoint::new(
self.chain.header_chain.height(),
self.chain.header_chain.tip_hash(),
),
self.chain.last_ten(),
);
self.dialog.send_event(Event::FiltersSynced(update));
self.emit_sync_complete();
}
}
NodeState::FiltersSynced => {
Expand All @@ -366,6 +368,17 @@ impl Node {
}
}

fn emit_sync_complete(&self) {
let update = SyncUpdate::new(
HashCheckpoint::new(
self.chain.header_chain.height(),
self.chain.header_chain.tip_hash(),
),
self.chain.last_ten(),
);
self.dialog.send_event(Event::FiltersSynced(update));
}

// When syncing headers we are only interested in one peer to start
fn next_required_peers(&self) -> PeerRequirement {
match self.state {
Expand All @@ -384,7 +397,11 @@ impl Node {
stop_hash: BlockHash::all_zeros(),
};
return Some(MainThreadMessage::GetHeaders(headers));
} else if !self.chain.is_cf_headers_synced() {
}
if self.headers_only_sync {
return None;
}
if !self.chain.is_cf_headers_synced() {
return Some(MainThreadMessage::GetFilterHeaders(
self.chain.next_cf_header_message(),
));
Expand Down Expand Up @@ -638,14 +655,22 @@ impl Node {
NodeState::Behind => None,
NodeState::HeadersSynced => None,
_ => {
self.headers_only_sync = false;
self.chain.clear_filters();
if let Some(height) = height_opt {
self.chain.header_chain.assume_checked_to(height);
}
self.state = NodeState::FilterHeadersSynced;
Some(MainThreadMessage::GetFilters(
self.chain.next_filter_message(),
))
if !self.chain.is_cf_headers_synced() {
self.state = NodeState::HeadersSynced;
Some(MainThreadMessage::GetFilterHeaders(
self.chain.next_cf_header_message(),
))
} else {
self.state = NodeState::FilterHeadersSynced;
Some(MainThreadMessage::GetFilters(
self.chain.next_filter_message(),
))
}
}
}
}
Expand Down
90 changes: 90 additions & 0 deletions tests/core.rs
Original file line number Diff line number Diff line change
Expand Up @@ -733,6 +733,96 @@ async fn whitelist_only_sync() {
rpc.stop().unwrap();
}

#[tokio::test]
async fn headers_only_sync_then_rescan() {
setup_debug_output();
let (bitcoind, socket_addr) = start_bitcoind(true).unwrap();
let rpc = &bitcoind.client;
let miner = rpc.new_address().unwrap();
mine_blocks(rpc, &miner, 10, 2).await;
let best = best_hash(rpc);
let host = (IpAddr::V4(*socket_addr.ip()), Some(socket_addr.port()));
let builder = bip157::builder::Builder::new(bitcoin::Network::Regtest)
.chain_state(ChainState::Checkpoint(HashCheckpoint::from_genesis(
bitcoin::Network::Regtest,
)))
.add_peer(host)
.headers_only_sync();
let (node, client) = builder.build();
tokio::task::spawn(async move { node.run().await });
let Client {
requester,
info_rx,
warn_rx,
event_rx: mut channel,
} = client;
tokio::task::spawn(async move { print_logs(info_rx, warn_rx).await });
let mut filters_before = 0usize;
tokio::time::timeout(Duration::from_secs(60), async {
loop {
match channel.recv().await {
Some(Event::IndexedFilter(_)) => filters_before += 1,
Some(Event::FiltersSynced(update)) => {
assert_eq!(update.tip().hash, best);
break;
}
_ => {}
}
}
})
.await
.expect("headers-only sync timed out");
assert_eq!(
filters_before, 0,
"no filters should be delivered during headers-only sync"
);
// A block mined after the initial catchup should stream a filter, since the flag
// only skips the initial sync and flips off once the tip is reached.
mine_blocks(rpc, &miner, 1, 2).await;
let best = best_hash(rpc);
let mut new_block_filters = 0usize;
tokio::time::timeout(Duration::from_secs(60), async {
loop {
match channel.recv().await {
Some(Event::IndexedFilter(_)) => new_block_filters += 1,
Some(Event::FiltersSynced(update)) => {
assert_eq!(update.tip().hash, best);
break;
}
_ => {}
}
}
})
.await
.expect("new-block filter sync timed out");
assert_eq!(
new_block_filters, 1,
"one filter should be delivered for the newly mined block"
);
requester.rescan().unwrap();
let mut filters_after = 0usize;
tokio::time::timeout(Duration::from_secs(60), async {
loop {
match channel.recv().await {
Some(Event::IndexedFilter(_)) => filters_after += 1,
Some(Event::FiltersSynced(update)) => {
assert_eq!(update.tip().hash, best);
break;
}
_ => {}
}
}
})
.await
.expect("rescan sync timed out");
assert_eq!(
filters_after, 11,
"rescan should deliver filters for the synced chain"
);
requester.shutdown().unwrap();
rpc.stop().unwrap();
}

#[tokio::test]
async fn inv_fallback_after_burst_mine() {
setup_debug_output();
Expand Down
Loading