diff --git a/src/bin/electrs.rs b/src/bin/electrs.rs index 74d06f806..6f94d028b 100644 --- a/src/bin/electrs.rs +++ b/src/bin/electrs.rs @@ -60,7 +60,13 @@ fn run_server(config: Arc, salt_rwlock: Arc>) -> Result<( info!("starting electrs"); if let Some(zmq_addr) = config.zmq_addr.as_ref() { - zmq::start(&format!("tcp://{zmq_addr}"), block_hash_notify); + if let Err(e) = zmq::start(&format!("tcp://{zmq_addr}"), &block_hash_notify, &metrics) { + error!( + "ZMQ notifications disabled, falling back to polling: zmq_addr='{}' err='{}'", + zmq_addr, + e.display_chain() + ); + } } info!("connecting to daemon at {}", config.daemon_rpc_addr); @@ -86,7 +92,15 @@ fn run_server(config: Arc, salt_rwlock: Arc>) -> Result<( ); info!("starting initial sync"); let mut tip = indexer.update(&daemon)?; - info!("initial sync complete, tip at {}", tip); + let daemon_tip = daemon.getbestblockhash()?; + if tip == daemon_tip { + info!("initial sync complete, tip at {}", tip); + } else { + info!( + "initial sync incomplete, serving a partial index and retrying outstanding blocks in the background tip='{}' daemon_tip='{}'", + tip, daemon_tip + ); + } let chain = Arc::new(ChainQuery::new( Arc::clone(&store), @@ -108,10 +122,18 @@ fn run_server(config: Arc, salt_rwlock: Arc>) -> Result<( Arc::clone(&config), ))); + // A mempool update is retried whenever the tip moves, so index and try again. An update + // that cannot advance the tip never will here, and no listener is bound yet. while !Mempool::update(&mempool, &daemon, &tip)? { - // Mempool syncing was aborted because the chain tip moved; - // Index the new block(s) and try again. - tip = indexer.update(&daemon)?; + let new_tip = indexer.update(&daemon)?; + if new_tip == tip { + warn!( + "index could not advance, starting up with a partial index tip='{}'", + tip + ); + break; + } + tip = new_tip; } #[cfg(feature = "liquid")] diff --git a/src/daemon.rs b/src/daemon.rs index d0b207370..7e4442e73 100644 --- a/src/daemon.rs +++ b/src/daemon.rs @@ -1270,7 +1270,14 @@ impl Daemon { pub fn getblock(&self, blockhash: &BlockHash) -> Result { let block = block_from_value(self.request("getblock", json!([blockhash, /*verbose=*/ false]))?)?; - assert_eq!(block.block_hash(), *blockhash); + let returned = block.block_hash(); + if returned != *blockhash { + bail!( + "bitcoind returned wrong block for getblock requested_hash='{}' returned_hash='{}'", + blockhash, + returned + ); + } Ok(block) } @@ -1295,23 +1302,46 @@ impl Daemon { Err(e) => { let err_msg = format!("{e:?}"); if err_msg.contains("Block not found on disk") - || err_msg.contains("Block not available") + || err_msg.contains("Block not available") { // There is a small chance the node returns the header but didn't finish to index the block - log::warn!("getblocks failing with: {e:?} trying {attempts} more time") + if attempts == 0 { + bail!( + "failed to get blocks from bitcoind attempts='0' err='{}'", + err_msg + ); + } + log::warn!( + "getblocks failed, retrying err='{}' attempts_left='{}'", + err_msg, + attempts + ); } else { - panic!("failed to get blocks from bitcoind: {}", err_msg); + bail!("failed to get blocks from bitcoind err='{}'", err_msg); } } } - if attempts == 0 { - panic!("failed to get blocks from bitcoind") - } std::thread::sleep(RETRY_WAIT_DURATION); }; - let mut blocks = vec![]; - for value in values { - blocks.push(block_from_value(value)?); + if values.len() != blockhashes.len() { + bail!( + "bitcoind returned wrong number of blocks requested='{}' returned='{}'", + blockhashes.len(), + values.len() + ); + } + let mut blocks = Vec::with_capacity(values.len()); + for (value, requested) in values.into_iter().zip(blockhashes.iter()) { + let block = block_from_value(value)?; + let returned = block.block_hash(); + if returned != *requested { + bail!( + "bitcoind returned wrong block for getblocks requested_hash='{}' returned_hash='{}'", + requested, + returned + ); + } + blocks.push(block); } Ok(blocks) } @@ -1466,20 +1496,31 @@ impl Daemon { #[trace] fn get_all_headers(&self, tip: &BlockHash) -> Result> { + const CHUNK_SIZE: usize = 100_000; + let info: Value = self.request("getblockheader", json!([tip]))?; let tip_height = info .get("height") - .expect("missing height") - .as_u64() - .expect("non-numeric height") as usize; - let all_heights: Vec = (0..=tip_height).collect(); - let chunk_size = 100_000; - let mut result = vec![]; - for heights in all_heights.chunks(chunk_size) { - let mut headers = self.getblockheaders(&heights)?; - assert!(headers.len() == heights.len()); + .and_then(|v| v.as_u64()) + .ok_or_else(|| format!("bitcoind returned malformed getblockheader info='{:?}'", info))?; + let tip_height = usize::try_from(tip_height).map_err(|_| { + format!("bitcoind returned out-of-range tip_height='{tip_height}'") + })?; - result.append(&mut headers); + // Chunked so the height allocation stays bounded whatever height bitcoind reports. + let mut result = vec![]; + for start in (0..=tip_height).step_by(CHUNK_SIZE) { + let end = start.saturating_add(CHUNK_SIZE - 1).min(tip_height); + let heights: Vec = (start..=end).collect(); + let headers = self.getblockheaders(&heights)?; + if headers.len() != heights.len() { + bail!( + "bitcoind returned wrong number of headers requested='{}' returned='{}'", + heights.len(), + headers.len() + ); + } + result.extend(headers); debug!( "downloaded {}/{} block headers ({:.0}%)", @@ -1491,10 +1532,22 @@ impl Daemon { let mut blockhash = *DEFAULT_BLOCKHASH; for header in &result { - assert_eq!(header.prev_blockhash, blockhash); + if header.prev_blockhash != blockhash { + bail!( + "bitcoind returned non-contiguous headers expected_prev='{}' got_prev='{}'", + blockhash, + header.prev_blockhash + ); + } blockhash = header.block_hash(); } - assert_eq!(blockhash, *tip); + if blockhash != *tip { + bail!( + "bitcoind headers do not end at requested tip expected_tip='{}' got_tip='{}'", + tip, + blockhash + ); + } Ok(result) } diff --git a/src/electrum/server.rs b/src/electrum/server.rs index c843cb522..44d2420ed 100644 --- a/src/electrum/server.rs +++ b/src/electrum/server.rs @@ -585,11 +585,7 @@ impl Connection { #[trace(method = %method)] fn handle_command(&mut self, method: &str, params: &[Value], id: &Value) -> Result { - let timer = self - .stats - .latency - .with_label_values(&[method]) - .start_timer(); + let started = Instant::now(); let result = match method { "blockchain.block.header" => self.blockchain_block_header(¶ms), @@ -626,7 +622,11 @@ impl Connection { "server.add_peer" => self.server_add_peer(¶ms), &_ => { - warn!("rpc #{} unknown method {} {:?}", id, method, params); + debug!("rpc #{} unknown method {} {:?}", id, method, params); + self.stats + .latency + .with_label_values(&["unknown"]) + .observe(started.elapsed().as_secs_f64()); return Ok(json_rpc_error( format!("unknown method {}", method), Some(id), @@ -634,7 +634,12 @@ impl Connection { )); } }; - timer.observe_duration(); + + self.stats + .latency + .with_label_values(&[method]) + .observe(started.elapsed().as_secs_f64()); + Ok(match result { Ok(result) => json!({"jsonrpc": "2.0", "id": id, "result": result}), Err(e) => { diff --git a/src/new_index/db.rs b/src/new_index/db.rs index 180d773bb..73193fc7f 100644 --- a/src/new_index/db.rs +++ b/src/new_index/db.rs @@ -388,7 +388,7 @@ impl DB { } #[cfg(test)] - fn open_test(path: &Path) -> DB { + pub(crate) fn open_test(path: &Path) -> DB { let mut db_opts = rocksdb::Options::default(); db_opts.create_if_missing(true); db_opts.set_prefix_extractor(rocksdb::SliceTransform::create_fixed_prefix(33)); diff --git a/src/new_index/fetch.rs b/src/new_index/fetch.rs index 0dc92aeaa..24ea1fa40 100644 --- a/src/new_index/fetch.rs +++ b/src/new_index/fetch.rs @@ -54,22 +54,33 @@ type SizedBlock = (Block, u32); pub struct Fetcher { receiver: Receiver, - thread: thread::JoinHandle<()>, + thread: thread::JoinHandle>, } impl Fetcher { - fn from(receiver: Receiver, thread: thread::JoinHandle<()>) -> Self { + fn from(receiver: Receiver, thread: thread::JoinHandle>) -> Self { Fetcher { receiver, thread } } - pub fn map(self, mut func: F) + pub fn map(self, mut func: F) -> Result<()> where - F: FnMut(T) -> (), + F: FnMut(T) -> Result<()>, { + let mut consumer_error = None; for item in self.receiver { - func(item); + if let Err(e) = func(item) { + consumer_error = Some(e); + break; + } + } + let producer_result = match self.thread.join() { + Ok(result) => result, + Err(_) => bail!("fetcher thread panicked"), + }; + match consumer_error { + Some(e) => Err(e), + None => producer_result, } - self.thread.join().expect("fetcher thread panicked") } } @@ -103,10 +114,20 @@ fn bitcoind_fetcher( fetcher_count += 1; let blockhashes: Vec = entries.iter().map(|e| *e.hash()).collect(); - let blocks = daemon - .getblocks(&blockhashes) - .expect("failed to get blocks from bitcoind"); - assert_eq!(blocks.len(), entries.len()); + let blocks = daemon.getblocks(&blockhashes).chain_err(|| { + format!( + "failed to fetch block batch first_hash='{:?}' batch_size='{}'", + blockhashes.first(), + blockhashes.len() + ) + })?; + if blocks.len() != entries.len() { + bail!( + "bitcoind_fetcher block/entry count mismatch expected='{}' got='{}'", + entries.len(), + blocks.len() + ); + } let block_entries: Vec = blocks .into_iter() .zip(entries) @@ -120,12 +141,12 @@ fn bitcoind_fetcher( } }) .collect(); - assert_eq!(block_entries.len(), entries.len()); sender .send(block_entries) - .expect("failed to send fetched blocks"); + .chain_err(|| "bitcoind_fetcher receiver dropped")?; log::debug!("last fetch {:?}", entries.last()); } + Ok(()) }), )) } @@ -177,16 +198,18 @@ fn blkfiles_fetcher( }) .collect(); trace!("fetched {} blocks", block_entries.len()); - sender - .send(block_entries) - .expect("failed to send blocks entries from blk*.dat files"); - }); + sender.send(block_entries).chain_err(|| { + "failed to send block entries from blk*.dat files" + })?; + Ok(()) + })?; if !entry_map.is_empty() { - panic!( + bail!( "failed to index {} blocks from blk*.dat files", entry_map.len() ) } + Ok(()) }), )) } @@ -211,14 +234,15 @@ fn blkfiles_reader(blk_files: Vec, xor_key: Option<[u8; 8]>) -> Fetcher trace!("reading {:?}", path); let mut blob = fs::read(&path) - .unwrap_or_else(|e| panic!("failed to read {:?}: {:?}", path, e)); + .chain_err(|| format!("failed to read {:?}", path))?; if let Some(xor_key) = xor_key { blkfile_apply_xor_key(xor_key, &mut blob); } sender .send(blob) - .unwrap_or_else(|_| panic!("failed to send {:?} contents", path)); + .chain_err(|| format!("failed to send {:?} contents", path))?; } + Ok(()) }), ) } @@ -247,11 +271,13 @@ fn blkfiles_parser(blobs: Fetcher>, magic: u32) -> Fetcher, magic: u32) -> Result max_pos { + break; + } + // If Core's WriteBlockToDisk ftell fails, only the magic bytes and size will be written // and the block body won't be written to the blk*.dat file. // Since the first 4 bytes should contain the block's version, we can skip such blocks @@ -294,10 +326,73 @@ fn parse_blocks(pool: &rayon::ThreadPool, blob: Vec, magic: u32) -> Result>>() + }) +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn fetcher_surfaces_producer_error() { + let chan = SyncChannel::::new(1); + let fetcher = Fetcher::from( + chan.into_receiver(), + spawn_thread("failing_fetcher", move || bail!("producer failed")), + ); + + let error = fetcher.map(|_| Ok(())).unwrap_err(); + assert!(format!("{:?}", error).contains("producer failed")); + } + + #[test] + fn fetcher_delivers_partial_output_before_producer_error() { + let chan = SyncChannel::new(1); + let sender = chan.sender(); + let fetcher = Fetcher::from( + chan.into_receiver(), + spawn_thread("partially_failing_fetcher", move || { + sender.send(1).chain_err(|| "receiver dropped")?; + bail!("producer failed after partial output") + }), + ); + let mut items = vec![]; + + let error = fetcher + .map(|item| { + items.push(item); + Ok(()) + }) + .unwrap_err(); + + assert_eq!(items, vec![1]); + assert!(format!("{:?}", error).contains("producer failed after partial output")); + } + + #[test] + fn fetcher_surfaces_consumer_error() { + let chan = SyncChannel::new(1); + let sender = chan.sender(); + let fetcher = Fetcher::from( + chan.into_receiver(), + spawn_thread("successful_fetcher", move || { + sender.send(1).chain_err(|| "receiver dropped")?; + Ok(()) + }), + ); + + let error = fetcher + .map(|_| bail!("consumer failed")) + .unwrap_err(); + assert!(format!("{:?}", error).contains("consumer failed")); + } } diff --git a/src/new_index/query.rs b/src/new_index/query.rs index 86dada56e..14a851f5c 100644 --- a/src/new_index/query.rs +++ b/src/new_index/query.rs @@ -1,5 +1,5 @@ use std::collections::{BTreeSet, HashMap}; -use std::sync::{Arc, RwLock, RwLockReadGuard}; +use std::sync::{Arc, Mutex, RwLock, RwLockReadGuard, TryLockError}; use std::time::{Duration, Instant}; use crate::chain::{Network, OutPoint, Transaction, TxOut, Txid}; @@ -32,7 +32,9 @@ pub struct Query { daemon: Arc, config: Arc, cached_estimates: RwLock<(HashMap, Option)>, + estimates_refresh: Mutex<()>, cached_relayfee: RwLock>, + relayfee_refresh: Mutex<()>, cached_block_template: BlockTemplateCache, #[cfg(feature = "liquid")] asset_db: Option>>, @@ -52,7 +54,9 @@ impl Query { daemon, config, cached_estimates: RwLock::new((HashMap::new(), None)), + estimates_refresh: Mutex::new(()), cached_relayfee: RwLock::new(None), + relayfee_refresh: Mutex::new(()), cached_block_template: BlockTemplateCache::new(), } } @@ -233,7 +237,7 @@ impl Query { } } - self.update_fee_estimates(); + self.refresh_fee_estimates_if_stale(); self.cached_estimates .read() .unwrap() @@ -250,10 +254,25 @@ impl Query { } } - self.update_fee_estimates(); + self.refresh_fee_estimates_if_stale(); self.cached_estimates.read().unwrap().0.clone() } + #[trace] + fn refresh_fee_estimates_if_stale(&self) { + let _guard = match self.estimates_refresh.try_lock() { + Ok(guard) => guard, + Err(TryLockError::WouldBlock) => return, + Err(TryLockError::Poisoned(poisoned)) => poisoned.into_inner(), + }; + if let (_, Some(cache_time)) = *self.cached_estimates.read().unwrap() { + if cache_time.elapsed() < Duration::from_secs(FEE_ESTIMATES_TTL) { + return; + } + } + self.update_fee_estimates(); + } + #[trace] fn update_fee_estimates(&self) { match self.daemon.estimatesmartfee_batch(&CONF_TARGETS) { @@ -272,6 +291,14 @@ impl Query { return Ok(cached); } + let _guard = self + .relayfee_refresh + .lock() + .unwrap_or_else(|poisoned| poisoned.into_inner()); + if let Some(cached) = *self.cached_relayfee.read().unwrap() { + return Ok(cached); + } + let relayfee = self.daemon.get_relayfee()?; self.cached_relayfee.write().unwrap().replace(relayfee); Ok(relayfee) @@ -292,7 +319,9 @@ impl Query { config, asset_db, cached_estimates: RwLock::new((HashMap::new(), None)), + estimates_refresh: Mutex::new(()), cached_relayfee: RwLock::new(None), + relayfee_refresh: Mutex::new(()), cached_block_template: BlockTemplateCache::new(), } } diff --git a/src/new_index/schema.rs b/src/new_index/schema.rs index 6ee9b8732..2ef6ce74d 100644 --- a/src/new_index/schema.rs +++ b/src/new_index/schema.rs @@ -45,6 +45,10 @@ use bitcoin::VarInt; const MIN_HISTORY_ITEMS_TO_CACHE: usize = 100; +fn completed_prefix_len(items: &[T], mut is_complete: impl FnMut(&T) -> bool) -> usize { + items.iter().take_while(|item| is_complete(item)).count() +} + pub struct Store { // TODO: should be column families txstore_db: DB, @@ -138,6 +142,55 @@ impl Store { pub fn done_initial_sync(&self) -> bool { self.txstore_db.get(b"t").is_some() } + + // Headers that need any work: either not yet added to txstore or not yet indexed to history. + fn headers_to_process(&self, new_headers: &[HeaderEntry]) -> Vec { + let added = self.added_blockhashes.read().unwrap(); + let indexed = self.indexed_blockhashes.read().unwrap(); + new_headers + .iter() + .filter(|e| !added.contains(e.hash()) || !indexed.contains(e.hash())) + .cloned() + .collect() + } + + // The contiguous prefix of `new_headers` that was both added and indexed. Blocks + // after the first gap stay outstanding, for the next index update to re-request. + fn completed_header_prefix(&self, new_headers: &[HeaderEntry]) -> Vec { + let added = self.added_blockhashes.read().unwrap(); + let indexed = self.indexed_blockhashes.read().unwrap(); + let len = completed_prefix_len(new_headers, |e| { + added.contains(e.hash()) && indexed.contains(e.hash()) + }); + new_headers[..len].to_vec() + } + + // Append the completed prefix to the in-memory HeaderList and persist the resulting + // tip under `t`. Errors on an empty index rather than returning the all-zeroes hash + // that `HeaderList::tip()` yields there, which reads as a real tip to the caller. + fn advance_synced_tip(&self, completed_headers: Vec) -> Result { + let mut headers = self.indexed_headers.write().unwrap(); + headers.append(completed_headers); + if headers.is_empty() { + bail!("no blocks were indexed, the index is still empty"); + } + let synced_tip = *headers.tip(); + debug!("updating synced tip to synced_tip='{}'", synced_tip); + self.txstore_db.put_sync(b"t", &serialize(&synced_tip)); + Ok(synced_tip) + } + + #[cfg(test)] + fn open_test(path: &std::path::Path) -> Self { + Store { + txstore_db: DB::open_test(&path.join("txstore")), + history_db: DB::open_test(&path.join("history")), + cache_db: DB::open_test(&path.join("cache")), + added_blockhashes: RwLock::new(HashSet::new()), + indexed_blockhashes: RwLock::new(HashSet::new()), + indexed_headers: RwLock::new(HeaderList::empty()), + } + } } type UtxoMap = HashMap; @@ -269,17 +322,6 @@ impl Indexer { self.duration.with_label_values(&[name]).start_timer() } - // Headers that need any work: either not yet added to txstore or not yet indexed to history. - fn headers_to_process(&self, new_headers: &[HeaderEntry]) -> Vec { - let added = self.store.added_blockhashes.read().unwrap(); - let indexed = self.store.indexed_blockhashes.read().unwrap(); - new_headers - .iter() - .filter(|e| !added.contains(e.hash()) || !indexed.contains(e.hash())) - .cloned() - .collect() - } - fn start_auto_compactions(&self, db: &DB) { let key = b"F".to_vec(); if db.get(&key).is_none() { @@ -358,7 +400,10 @@ impl Indexer { // Fetch the reorged blocks, then undo their history index db rows. // The txstore db rows are kept for reorged blocks/transactions. start_fetcher(self.from, &daemon, reorged_headers, self.iconfig.block_batch_size, chain_tip_height)? - .map(|blocks| self.undo_index(&blocks)); + .map(|blocks| { + self.undo_index(&blocks); + Ok(()) + })?; } // Single-pass: add to txstore and index to history in the same per-batch loop. @@ -375,7 +420,7 @@ impl Indexer { // Crash safety: added_blockhashes / indexed_blockhashes are persisted via the // "D" done-marker rows. On restart, headers_to_process() re-derives which // blocks still need work, so partially-processed batches are re-processed safely. - let to_process = self.headers_to_process(&new_headers); + let to_process = self.store.headers_to_process(&new_headers); debug!( "processing {} blocks (add + index) using {:?}", to_process.len(), @@ -385,7 +430,14 @@ impl Indexer { let mut fetcher_count = 0; let to_process_total = to_process.len(); - start_fetcher(self.from, &daemon, to_process, self.iconfig.block_batch_size, chain_tip_height)?.map(|blocks| { + let fetch_result = start_fetcher( + self.from, + &daemon, + to_process, + self.iconfig.block_batch_size, + chain_tip_height, + )? + .map(|blocks| { if fetcher_count % 25 == 0 && to_process_total > 20 { let batch_height = blocks.last().map(|b| b.entry.height()).unwrap_or(0); info!( @@ -433,8 +485,29 @@ impl Indexer { self.sync_progress.set(h as f64 / chain_tip_height as f64 * 100.0); } } + Ok(()) }); + if let FetchFrom::BlkFiles = self.from { + self.from = FetchFrom::Bitcoind; + } + + if let Err(ref e) = fetch_result { + warn!( + "block fetch incomplete, advancing only completed prefix and retrying on next index update: {:?}", + e + ); + } + + let completed_headers = self.store.completed_header_prefix(&new_headers); + if fetch_result.is_ok() && completed_headers.len() != new_headers.len() { + bail!( + "block fetch completed without indexing all headers: completed={} requested={}", + completed_headers.len(), + new_headers.len() + ); + } + // Compact after all add+index work is done, not between passes. self.start_auto_compactions(&self.store.txstore_db); self.start_auto_compactions(&self.store.history_db); @@ -461,23 +534,23 @@ impl Indexer { self.flush = DBFlush::Enable; } - // Update the synced tip after all db writes are flushed - debug!("updating synced tip to {:?}", tip); - self.store.txstore_db.put_sync(b"t", &serialize(&tip)); - - // Finally, append the new headers to the in-memory HeaderList. + // Finally, append the completed headers to the in-memory HeaderList. // This will make both the headers and the history entries visible in the public APIs, consistently with each-other. - let mut headers = self.store.indexed_headers.write().unwrap(); - headers.append(new_headers); - assert_eq!(tip, *headers.tip()); - - if let FetchFrom::BlkFiles = self.from { - self.from = FetchFrom::Bitcoind; + // Done only once all the db writes above are flushed. + let synced_tip = self.store.advance_synced_tip(completed_headers)?; + + if fetch_result.is_ok() && synced_tip != tip { + bail!( + "synced tip does not match the daemon tip after a complete fetch daemon_tip='{}' synced_tip='{}'", + tip, + synced_tip + ); } - self.tip_metric.set(headers.best_height() as i64); + self.tip_metric + .set(self.store.headers().best_height() as i64); - Ok(tip) + Ok(synced_tip) } fn add(&self, blocks: &[BlockEntry]) { @@ -2001,6 +2074,156 @@ pub mod bench { mod tests { use super::*; + #[test] + fn test_completed_prefix_stops_at_first_gap() { + let items = [true, true, false, true]; + assert_eq!(completed_prefix_len(&items, |complete| *complete), 2); + } + + #[cfg(not(feature = "liquid"))] + fn test_header(prev_blockhash: BlockHash) -> BlockHeader { + BlockHeader { + version: bitcoin::block::Version::ONE, + prev_blockhash, + merkle_root: crate::chain::TxMerkleNode::all_zeros(), + time: 0, + bits: bitcoin::CompactTarget::from_consensus(0x207f_ffff), + nonce: 0, + } + } + + #[cfg(feature = "liquid")] + fn test_header(prev_blockhash: BlockHash) -> BlockHeader { + BlockHeader { + version: 1, + prev_blockhash, + merkle_root: crate::chain::TxMerkleNode::all_zeros(), + time: 0, + height: 0, + ext: elements::BlockExtData::Proof { + challenge: Script::new(), + solution: Script::new(), + }, + } + } + + /// A chain of `count` linked headers starting at height 0. + fn test_chain(count: usize) -> Vec { + let mut prev = *crate::util::DEFAULT_BLOCKHASH; + (0..count) + .map(|height| { + let header = test_header(prev); + prev = header.block_hash(); + HeaderEntry::new(height, prev, header) + }) + .collect() + } + + /// Mark a block as both added to txstore and indexed to history. + fn mark_complete(store: &Store, entry: &HeaderEntry) { + store + .added_blockhashes + .write() + .unwrap() + .insert(*entry.hash()); + store + .indexed_blockhashes + .write() + .unwrap() + .insert(*entry.hash()); + } + + fn persisted_tip(store: &Store) -> Option { + store + .txstore_db + .get(b"t") + .map(|raw| deserialize(&raw).unwrap()) + } + + #[test] + fn test_synced_tip_stops_at_the_first_unindexed_block() { + let dir = tempfile::tempdir().unwrap(); + let store = Store::open_test(dir.path()); + let headers = test_chain(4); + + // Heights 0, 1 and 3 completed; height 2 was lost to a failed fetch. The tip + // must stop at height 1, not run on to the completed block above the gap. + mark_complete(&store, &headers[0]); + mark_complete(&store, &headers[1]); + mark_complete(&store, &headers[3]); + + let completed = store.completed_header_prefix(&headers); + assert_eq!(completed.len(), 2); + + let synced_tip = store.advance_synced_tip(completed).unwrap(); + assert_eq!(synced_tip, *headers[1].hash()); + assert_eq!(persisted_tip(&store), Some(*headers[1].hash())); + assert_eq!(store.headers().len(), 2); + + // Only the block that actually failed is outstanding. Height 3 is already in the + // db and merely needs the gap below it filled before it can be appended. + let outstanding = store.headers_to_process(&headers); + assert_eq!(outstanding.len(), 1); + assert_eq!(outstanding[0].height(), 2); + + // Once the refetch of height 2 lands, the next update carries the tip past the gap. + mark_complete(&store, &headers[2]); + let completed = store.completed_header_prefix(&headers[2..]); + assert_eq!(completed.len(), 2); + + let synced_tip = store.advance_synced_tip(completed).unwrap(); + assert_eq!(synced_tip, *headers[3].hash()); + assert_eq!(persisted_tip(&store), Some(*headers[3].hash())); + } + + #[test] + fn test_completed_prefix_requires_both_added_and_indexed() { + let dir = tempfile::tempdir().unwrap(); + let store = Store::open_test(dir.path()); + let headers = test_chain(2); + + // Added to txstore but never indexed to history. + store + .added_blockhashes + .write() + .unwrap() + .insert(*headers[0].hash()); + + assert!(store.completed_header_prefix(&headers).is_empty()); + } + + #[test] + fn test_advance_synced_tip_errs_when_nothing_was_indexed() { + let dir = tempfile::tempdir().unwrap(); + let store = Store::open_test(dir.path()); + + // Every block of the first batch failed to fetch, so the prefix is empty and + // there is no tip to report. `HeaderList::tip()` would hand back the all-zeroes + // default hash here, which reads as a successful sync to the caller. + assert!(store.advance_synced_tip(vec![]).is_err()); + assert_eq!(persisted_tip(&store), None); + } + + #[test] + fn test_advance_synced_tip_keeps_the_existing_tip_when_nothing_completed() { + let dir = tempfile::tempdir().unwrap(); + let store = Store::open_test(dir.path()); + let headers = test_chain(3); + + for entry in &headers[..2] { + mark_complete(&store, entry); + } + store + .advance_synced_tip(store.completed_header_prefix(&headers)) + .unwrap(); + + // A later update whose very first block fails to fetch must hold the tip where + // it is, not error out on an index that already has blocks in it. + let synced_tip = store.advance_synced_tip(vec![]).unwrap(); + assert_eq!(synced_tip, *headers[1].hash()); + assert_eq!(persisted_tip(&store), Some(*headers[1].hash())); + } + #[test] fn test_compute_script_hash_p2pkh() { // P2PKH scriptPubKey for address 1A1zP1eP5QGefi2DMPTfTL5SLmv7DivfNa diff --git a/src/new_index/zmq.rs b/src/new_index/zmq.rs index 73d5d7964..25353b983 100644 --- a/src/new_index/zmq.rs +++ b/src/new_index/zmq.rs @@ -1,40 +1,604 @@ +//! Subscriber for bitcoind's ZMQ `hashblock` notifications. +//! +//! The notification is only ever used as a wake-up hint: `Waiter::wait` discards +//! the hash and simply lets the main loop run a tip check one cycle early. The +//! subscriber therefore treats everything the publisher sends as untrusted input +//! and bounds what it can cost us - there is no transport authentication here, so +//! anyone able to reach or impersonate the endpoint can publish to us. + +use std::sync::atomic::{AtomicUsize, Ordering}; +use std::thread; +use std::time::{Duration, Instant}; + use bitcoin::{hashes::Hash, BlockHash}; -use crossbeam_channel::Sender; +use crossbeam_channel::{Sender, TrySendError}; +use error_chain::ChainedError; +use crate::errors::*; +use crate::metrics::{CounterVec, MetricOpts, Metrics}; use crate::util::spawn_thread; -pub fn start(url: &str, block_hash_notify: Sender) { - log::debug!("Starting ZMQ thread"); +const TOPIC_HASHBLOCK: &[u8] = b"hashblock"; +const BLOCK_HASH_BYTES: usize = 32; + +/// bitcoind publishes `hashblock` as three frames - topic, 32-byte hash, 4-byte +/// sequence number - roughly 45 bytes in total. Anything materially larger is not +/// something we have a use for. +/// +/// This bound is applied in two places and needs both. `ZMQ_MAXMSGSIZE` is the only +/// one that acts before libzmq allocates: the length prefix is checked as it is +/// decoded, so a publisher declaring a huge frame is refused without the body being +/// read or allocated. A check after `recv` cannot do that - by then the allocation +/// has already happened, and a large enough declared length aborts the process +/// rather than failing. The second application, in [`frame_admissible`], covers what +/// a size limit inherently cannot: how *many* frames arrive in one message. +/// +/// The catch is that libzmq classes an oversized message as a *protocol* error, so +/// it tears the session down, does not reconnect, and reports nothing to `recv` - +/// notifications would simply stop, silently and permanently. That is what +/// [`Subscriber::reconnect_if_peer_dropped`] exists to notice. +const MAX_FRAME_BYTES: usize = 1024; + +/// Upper bound on frames buffered from one multipart message. A size limit bounds +/// each frame but not how many a publisher may send, so cap that separately - +/// otherwise unbounded small frames add up to the same thing. +const MAX_FRAMES: usize = 8; + +/// Bound on libzmq's own receive queue. SUB sockets discard messages past the +/// high-water mark, which is the behaviour we want given we coalesce anyway. +const RCVHWM: i32 = 16; + +/// Return from `recv` this often even when idle, so the thread can notice that +/// the receiving end has gone away and exit instead of blocking forever. +const RCVTIMEO_MS: i32 = 1000; + +/// Minimum spacing between forwarded notifications. The main loop polls every 5s +/// regardless, so coalescing at this granularity costs nothing in normal +/// operation while bounding how much daemon RPC work a flood of notifications can +/// force us into. +const MIN_NOTIFY_INTERVAL: Duration = Duration::from_millis(250); + +const MIN_ERROR_BACKOFF: Duration = Duration::from_millis(100); +const MAX_ERROR_BACKOFF: Duration = Duration::from_secs(5); + +/// Minimum-interval gate over a stream of events. +struct Throttle { + min_interval: Duration, + last: Option, +} + +impl Throttle { + fn new(min_interval: Duration) -> Self { + Throttle { + min_interval, + last: None, + } + } + + /// Whether an event at `now` should be let through. Anything arriving less + /// than `min_interval` after the last admitted event is rejected. + fn allow(&mut self, now: Instant) -> bool { + match self.last { + Some(last) if now.saturating_duration_since(last) < self.min_interval => false, + _ => { + self.last = Some(now); + true + } + } + } +} + +/// Extract the block hash from a well-formed `hashblock` notification. +/// +/// Returns `None` for a topic we did not subscribe to, a body that is not exactly +/// a hash, or a message missing frames. bitcoind publishes the hash in display +/// byte order, so it is reversed to get the internal order `BlockHash` expects. +fn parse_hashblock(frames: &[Vec]) -> Option { + let topic = frames.first()?; + let body = frames.get(1)?; + + if topic.as_slice() != TOPIC_HASHBLOCK || body.len() != BLOCK_HASH_BYTES { + return None; + } + + let mut reversed = body.clone(); + reversed.reverse(); + BlockHash::from_slice(&reversed).ok() +} + +/// Whether a frame is worth copying out of libzmq, given how many we already hold. +/// +/// The count is the part that matters here: `ZMQ_MAXMSGSIZE` bounds how large a +/// frame may be but says nothing about how many of them one message may contain, so +/// a publisher is otherwise free to send unlimited small parts. The size check is +/// kept alongside it as a backstop, since it also covers transports that do not go +/// through the wire decoder that enforces `ZMQ_MAXMSGSIZE`. +fn frame_admissible(frame_len: usize, already_buffered: usize) -> bool { + frame_len <= MAX_FRAME_BYTES && already_buffered < MAX_FRAMES +} + +/// Receive one complete multipart message, buffering at most [`MAX_FRAMES`] frames +/// of at most [`MAX_FRAME_BYTES`] each. +/// +/// Inadmissible frames are drained and dropped where they land rather than copied, +/// so neither an oversized part nor an unbounded number of parts can grow our +/// memory. The whole message is still read to completion - leaving parts queued +/// would desynchronise the next read. Returns `None` when anything was rejected, +/// which discards the message rather than acting on a partial one. +fn recv_message(subscriber: &zmq::Socket) -> zmq::Result>>> { + let mut frames = Vec::new(); + let mut rejected = false; + + loop { + let frame = subscriber.recv_msg(0)?; + if frame_admissible(frame.len(), frames.len()) { + frames.push(frame.to_vec()); + } else { + rejected = true; + } + + // libzmq delivers a multipart message atomically, so the remaining parts + // are already queued and these calls do not block. + if !subscriber.get_rcvmore()? { + break; + } + } + + Ok(if rejected { None } else { Some(frames) }) +} + +/// A SUB socket together with the monitor used to notice it losing its peer. +/// +/// The two are kept as a unit because they are rebuilt as a unit: a monitor is +/// bound to one specific socket, so replacing the socket means replacing the +/// monitor with it. +struct Subscriber { + ctx: zmq::Context, + url: String, + socket: zmq::Socket, + monitor: zmq::Socket, + /// Whether a rebuild is owed. Latched rather than inferred from the monitor, + /// because the monitor only reports the *transition* and a torn-down socket + /// never reports a second one - so a rebuild that fails has to leave the + /// fact behind, or nothing will ever ask for another. + dead: bool, +} + +impl Subscriber { + fn connect(ctx: &zmq::Context, url: &str) -> Result { + let socket = ctx + .socket(zmq::SUB) + .chain_err(|| format!("failed creating ZMQ subscriber for url='{url}'"))?; + + // Refuse an oversized frame while its length prefix is being decoded, + // before libzmq reads or allocates the body. A check after `recv` is too + // late to prevent the allocation. + socket + .set_maxmsgsize(MAX_FRAME_BYTES as i64) + .chain_err(|| "failed setting ZMQ maxmsgsize")?; + socket + .set_rcvhwm(RCVHWM) + .chain_err(|| "failed setting ZMQ rcvhwm")?; + socket + .set_rcvtimeo(RCVTIMEO_MS) + .chain_err(|| "failed setting ZMQ rcvtimeo")?; + + // Attached before connecting, so no event can be missed in between. + let monitor = attach_monitor(ctx, &socket)?; + + socket + .connect(url) + .chain_err(|| format!("failed connecting ZMQ subscriber to url='{url}'"))?; + socket + .set_subscribe(TOPIC_HASHBLOCK) + .chain_err(|| "failed subscribing to ZMQ hashblock")?; + + Ok(Subscriber { + ctx: ctx.clone(), + url: url.to_owned(), + socket, + monitor, + dead: false, + }) + } + + fn recv(&self) -> zmq::Result>>> { + recv_message(&self.socket) + } + + /// Rebuild the socket if libzmq has dropped the connection to the publisher. + /// + /// Reconnecting is normally libzmq's job, and for an ordinary disconnect it + /// does it. It makes an exception for a protocol error - an oversized message + /// being the one a publisher can trigger at will - where it terminates the + /// session instead and neither reconnects nor surfaces an error. `recv` just + /// returns `EAGAIN` forever, so the only evidence is the monitor event emitted + /// on the way down. Rebuilding on any disconnect covers that case without + /// having to distinguish it, and costs nothing when libzmq would have + /// recovered on its own. + /// + /// Returns whether a rebuild was attempted, successfully or not. + fn reconnect_if_peer_dropped(&mut self) -> bool { + // `|=` rather than a short-circuiting `||`: the monitor has to be drained + // on every pass, including while a rebuild is already owed, so a backlog + // cannot accumulate behind the flag. + self.dead |= self.peer_dropped(); + if !self.dead { + return false; + } + + log::warn!( + "ZMQ publisher connection dropped, reconnecting: url='{}'", + self.url + ); + match Subscriber::connect(&self.ctx, &self.url) { + Ok(replacement) => { + self.socket = replacement.socket; + self.monitor = replacement.monitor; + self.dead = false; + } + Err(e) => { + // Deliberately still dead. The socket left in place has had its + // session terminated, so it will neither deliver anything nor + // report a second disconnect - only the flag can ask for the + // retry, and the caller's backoff paces it. + log::warn!( + "failed reconnecting ZMQ subscriber, will retry: url='{}' err='{}'", + self.url, + e.display_chain() + ); + } + } + true + } + + /// Whether any disconnect event is queued on the monitor, draining it either + /// way so a backlog cannot build up. + fn peer_dropped(&self) -> bool { + let mut dropped = false; + while let Ok(event) = self.monitor.recv_multipart(zmq::DONTWAIT) { + if let Some(raw) = event.first().filter(|frame| frame.len() >= 2) { + let id = u16::from_ne_bytes([raw[0], raw[1]]); + dropped |= id == zmq::SocketEvent::DISCONNECTED.to_raw(); + } + } + dropped + } +} + +/// Attach a disconnect monitor to `socket` and return the end we read events from. +fn attach_monitor(ctx: &zmq::Context, socket: &zmq::Socket) -> Result { + // One endpoint per socket: a monitor binds its own inproc address, and sockets + // are rebuilt over the process lifetime, so the name has to be unique. + static NEXT_MONITOR_ID: AtomicUsize = AtomicUsize::new(0); + let endpoint = format!( + "inproc://zmq-subscriber-monitor-{}", + NEXT_MONITOR_ID.fetch_add(1, Ordering::Relaxed) + ); + + socket + .monitor( + &endpoint, + i32::from(zmq::SocketEvent::DISCONNECTED.to_raw()), + ) + .chain_err(|| format!("failed monitoring ZMQ subscriber on endpoint='{endpoint}'"))?; + + let monitor = ctx + .socket(zmq::PAIR) + .chain_err(|| "failed creating ZMQ monitor socket")?; + monitor + .connect(&endpoint) + .chain_err(|| format!("failed connecting ZMQ monitor to endpoint='{endpoint}'"))?; + + Ok(monitor) +} + +/// Connect a SUB socket to `url` and forward `hashblock` notifications. +/// +/// Takes the sender by reference deliberately. If setup fails the caller's sender +/// must stay alive: dropping it would disconnect the channel, and `Waiter::wait` +/// treats a disconnected channel as a fatal error, so a ZMQ misconfiguration +/// would stop the server instead of falling back to polling. +pub fn start(url: &str, block_hash_notify: &Sender, metrics: &Metrics) -> Result<()> { + log::debug!("starting ZMQ subscriber: url='{url}'"); + let ctx = zmq::Context::new(); - let subscriber: zmq::Socket = ctx.socket(zmq::SUB).expect("failed creating subscriber"); - subscriber - .connect(url) - .expect("failed connecting subscriber"); - - // subscriber.set_subscribe(b"rawtx").unwrap(); - subscriber - .set_subscribe(b"hashblock") - .expect("failed subscribing to hashblock"); - - spawn_thread("zmq", move || loop { - match subscriber.recv_multipart(0) { - Ok(data) => match (data.get(0), data.get(1)) { - (Some(topic), Some(data)) => { - if &topic[..] == &[114, 97, 119, 116, 120] { - //rawtx - } else if &topic[..] == &[104, 97, 115, 104, 98, 108, 111, 99, 107] { - //hashblock - let mut reversed = data.to_vec(); - reversed.reverse(); - if let Ok(block_hash) = BlockHash::from_slice(&reversed[..]) { - log::debug!("New block from ZMQ: {block_hash}"); - let _ = block_hash_notify.send(block_hash); - } + let subscriber = Subscriber::connect(&ctx, url)?; + + let notifications = metrics.counter_vec( + MetricOpts::new( + "electrs_zmq_notifications_total", + "ZMQ block notifications received, by disposition", + ), + &["result"], + ); + + let block_hash_notify = block_hash_notify.clone(); + spawn_thread("zmq", move || { + subscriber_loop(subscriber, block_hash_notify, notifications) + }); + + Ok(()) +} + +fn subscriber_loop( + mut subscriber: Subscriber, + block_hash_notify: Sender, + notifications: CounterVec, +) { + let mut throttle = Throttle::new(MIN_NOTIFY_INTERVAL); + let mut backoff = MIN_ERROR_BACKOFF; + + loop { + // Checked every iteration, and `recv` returns at least once per + // `RCVTIMEO_MS`, so a silent teardown is picked up within that window. + if subscriber.reconnect_if_peer_dropped() { + notifications.with_label_values(&["reconnected"]).inc(); + // Spaced out deliberately: a publisher sending oversized frames can + // force a teardown per message, and rebuilding flat out would turn + // that into a hot loop. Backing off degrades to the tip poll instead. + thread::sleep(backoff); + backoff = (backoff * 2).min(MAX_ERROR_BACKOFF); + continue; + } + + match subscriber.recv() { + Ok(received) => { + backoff = MIN_ERROR_BACKOFF; + + // Over the frame cap. Already drained, so nothing was buffered. + let Some(frames) = received else { + notifications.with_label_values(&["rejected"]).inc(); + continue; + }; + let Some(block_hash) = parse_hashblock(&frames) else { + notifications.with_label_values(&["rejected"]).inc(); + continue; + }; + if !throttle.allow(Instant::now()) { + notifications.with_label_values(&["throttled"]).inc(); + continue; + } + + match block_hash_notify.try_send(block_hash) { + Ok(()) => { + notifications.with_label_values(&["forwarded"]).inc(); + log::debug!("new block from ZMQ: block_hash='{block_hash}'"); + } + // A wake-up is already pending, so this one would add nothing. + Err(TrySendError::Full(_)) => { + notifications.with_label_values(&["coalesced"]).inc(); + } + Err(TrySendError::Disconnected(_)) => { + log::debug!("ZMQ notification receiver is gone, stopping subscriber"); + return; } } - _ => (), - }, - Err(e) => log::warn!("recv_multipart error: {e:?}"), + } + // The receive timeout expired with nothing waiting. + Err(zmq::Error::EAGAIN) => continue, + Err(e) => { + log::warn!( + "ZMQ receive failed, backing off: url='{}' err='{e}' backoff='{backoff:?}'", + subscriber.url + ); + thread::sleep(backoff); + backoff = (backoff * 2).min(MAX_ERROR_BACKOFF); + } } - }); + } +} + +#[cfg(test)] +mod tests { + use super::*; + + fn frames(topic: &[u8], body: &[u8]) -> Vec> { + vec![topic.to_vec(), body.to_vec(), 0u32.to_le_bytes().to_vec()] + } + + #[test] + fn parses_a_well_formed_hashblock() { + let body: Vec = (0..BLOCK_HASH_BYTES as u8).collect(); + let parsed = parse_hashblock(&frames(TOPIC_HASHBLOCK, &body)).unwrap(); + + // The published bytes are in display order, so they should come back out + // unchanged. Skipping the reversal would render this the other way round. + assert_eq!( + parsed.to_string(), + "000102030405060708090a0b0c0d0e0f101112131415161718191a1b1c1d1e1f" + ); + } + + #[test] + fn rejects_an_unsubscribed_topic() { + let body = [0u8; BLOCK_HASH_BYTES]; + assert!(parse_hashblock(&frames(b"rawtx", &body)).is_none()); + assert!(parse_hashblock(&frames(b"hashtx", &body)).is_none()); + assert!(parse_hashblock(&frames(b"", &body)).is_none()); + } + + #[test] + fn rejects_a_body_that_is_not_a_hash() { + assert!(parse_hashblock(&frames(TOPIC_HASHBLOCK, &[0u8; 31])).is_none()); + assert!(parse_hashblock(&frames(TOPIC_HASHBLOCK, &[0u8; 33])).is_none()); + assert!(parse_hashblock(&frames(TOPIC_HASHBLOCK, &[])).is_none()); + } + + #[test] + fn rejects_a_message_missing_frames() { + assert!(parse_hashblock(&[TOPIC_HASHBLOCK.to_vec()]).is_none()); + assert!(parse_hashblock(&[]).is_none()); + } + + #[test] + fn admits_a_frame_of_the_size_bitcoind_actually_sends() { + // topic, 32-byte hash, 4-byte sequence. + assert!(frame_admissible(TOPIC_HASHBLOCK.len(), 0)); + assert!(frame_admissible(BLOCK_HASH_BYTES, 1)); + assert!(frame_admissible(4, 2)); + } + + #[test] + fn rejects_a_frame_over_the_size_limit() { + assert!(frame_admissible(MAX_FRAME_BYTES, 0)); + assert!(!frame_admissible(MAX_FRAME_BYTES + 1, 0)); + assert!(!frame_admissible(1024 * 1024, 0)); + } + + #[test] + fn rejects_a_frame_past_the_count_limit() { + assert!(frame_admissible(32, MAX_FRAMES - 1)); + assert!(!frame_admissible(32, MAX_FRAMES)); + assert!(!frame_admissible(32, MAX_FRAMES + 1)); + } + + /// Setup failure has to be reported, not panicked, and it must leave the + /// caller's sender alive. If `start` consumed the sender, this failure would + /// drop it, disconnect the channel, and turn a misconfigured endpoint into a + /// fatal error on the next `Waiter::wait` - the opposite of falling back to + /// polling. + #[test] + fn a_failed_start_reports_an_error_and_leaves_the_sender_usable() { + let (tx, rx) = crossbeam_channel::bounded(1); + let metrics = Metrics::new("127.0.0.1:0".parse().unwrap()); + + assert!(start("nosuchtransport://endpoint", &tx, &metrics).is_err()); + + let hash = BlockHash::from_slice(&[0u8; BLOCK_HASH_BYTES]).unwrap(); + assert!(tx.try_send(hash).is_ok()); + assert_eq!(rx.try_recv().unwrap(), hash); + } + + /// Publish `hashblock` until one is received, rebuilding the socket the way + /// the loop does. Retried because PUB drops anything published before the + /// subscription has reached it, so early attempts can be lost legitimately. + fn expect_notification(subscriber: &mut Subscriber, publisher: &zmq::Socket) -> bool { + let body = [7u8; BLOCK_HASH_BYTES]; + + for _ in 0..20 { + subscriber.reconnect_if_peer_dropped(); + publisher + .send_multipart([TOPIC_HASHBLOCK, &body[..]], 0) + .unwrap(); + + // Bounded by the receive timeout, so this paces the retries too. + if let Ok(Some(frames)) = subscriber.recv() { + if parse_hashblock(&frames).is_some() { + return true; + } + } + } + false + } + + /// An oversized frame is a *protocol* error to libzmq: it terminates the + /// session, does not reconnect, and reports nothing to `recv`, which just + /// returns `EAGAIN` from then on. Nothing at the receive layer can observe + /// that, so without the monitor the second `expect_notification` never + /// succeeds - one hostile message would end block notifications for the life + /// of the process. + /// + /// Deliberately over TCP: the size limit is applied by the wire decoder, and + /// inproc does not go through it. + #[test] + fn recovers_after_an_oversized_frame_tears_the_session_down() { + let ctx = zmq::Context::new(); + let publisher = ctx.socket(zmq::PUB).unwrap(); + publisher.bind("tcp://127.0.0.1:*").unwrap(); + let url = publisher.get_last_endpoint().unwrap().unwrap(); + + let mut subscriber = Subscriber::connect(&ctx, &url).unwrap(); + assert!( + expect_notification(&mut subscriber, &publisher), + "no notification received before the oversized frame" + ); + + let oversized = vec![0u8; MAX_FRAME_BYTES + 1]; + publisher + .send_multipart([TOPIC_HASHBLOCK, &oversized[..]], 0) + .unwrap(); + + assert!( + expect_notification(&mut subscriber, &publisher), + "subscriber never recovered from the oversized frame" + ); + } + + /// The rebuild itself can fail - `EMFILE` on either socket it creates is the + /// realistic way - and the attempt consumes the disconnect event that asked + /// for it. Since the socket left in place has had its session terminated, it + /// will never report a second disconnect, so nothing would ever ask again: + /// one transient failure would end block notifications permanently. + /// + /// Failure is induced by pointing the rebuild at a transport libzmq rejects, + /// which fails at the same place fd exhaustion would - after the event has + /// been taken off the monitor. + #[test] + fn a_failed_rebuild_is_retried_rather_than_latching_off() { + let ctx = zmq::Context::new(); + let publisher = ctx.socket(zmq::PUB).unwrap(); + publisher.bind("tcp://127.0.0.1:*").unwrap(); + let url = publisher.get_last_endpoint().unwrap().unwrap(); + + let mut subscriber = Subscriber::connect(&ctx, &url).unwrap(); + assert!( + expect_notification(&mut subscriber, &publisher), + "no notification received before the oversized frame" + ); + + let oversized = vec![0u8; MAX_FRAME_BYTES + 1]; + publisher + .send_multipart([TOPIC_HASHBLOCK, &oversized[..]], 0) + .unwrap(); + + subscriber.url = "bogus://not-a-transport".to_owned(); + let mut attempted = false; + for _ in 0..20 { + if subscriber.reconnect_if_peer_dropped() { + attempted = true; + break; + } + // Bounded by the receive timeout, so this paces the retries. + let _ = subscriber.recv(); + } + assert!(attempted, "the teardown was never observed"); + + subscriber.url = url; + assert!( + expect_notification(&mut subscriber, &publisher), + "subscriber never recovered from a failed rebuild" + ); + } + + #[test] + fn throttle_admits_the_first_event() { + let mut throttle = Throttle::new(MIN_NOTIFY_INTERVAL); + assert!(throttle.allow(Instant::now())); + } + + #[test] + fn throttle_rejects_inside_the_window_and_admits_on_the_boundary() { + let mut throttle = Throttle::new(Duration::from_millis(250)); + let t0 = Instant::now(); + + assert!(throttle.allow(t0)); + assert!(!throttle.allow(t0 + Duration::from_millis(1))); + assert!(!throttle.allow(t0 + Duration::from_millis(249))); + assert!(throttle.allow(t0 + Duration::from_millis(250))); + assert!(!throttle.allow(t0 + Duration::from_millis(251))); + } + + #[test] + fn throttle_bounds_a_sustained_flood() { + let mut throttle = Throttle::new(Duration::from_millis(250)); + let t0 = Instant::now(); + + // 10k notifications arriving over one second: admit at 0, 250, 500, 750. + let admitted = (0..10_000u64) + .filter(|i| throttle.allow(t0 + Duration::from_micros(i * 100))) + .count(); + + assert_eq!(admitted, 4); + } } diff --git a/src/signal.rs b/src/signal.rs index ee319c3f6..34f97f56d 100644 --- a/src/signal.rs +++ b/src/signal.rs @@ -38,38 +38,36 @@ impl Waiter { } pub fn wait(&self, duration: Duration, accept_block_notification: bool) -> Result<()> { - let start = Instant::now(); - select! { - recv(self.receiver) -> msg => { - match msg { - Ok(sig) if sig == SIGUSR1 => { - trace!("notified via SIGUSR1"); - if accept_block_notification { - Ok(()) - } else { - let wait_more = duration.saturating_sub(start.elapsed()); - self.wait(wait_more, accept_block_notification) + let mut remaining = duration; + + loop { + let start = Instant::now(); + + let deadline_expired = select! { + recv(self.receiver) -> msg => { + match msg { + Ok(sig) if sig == SIGUSR1 => { + trace!("notified via SIGUSR1"); + false } + Ok(sig) => bail!(ErrorKind::Interrupt(sig)), + Err(_) => bail!("signal hook channel disconnected"), } - Ok(sig) => bail!(ErrorKind::Interrupt(sig)), - Err(_) => bail!("signal hook channel disconnected"), - } - }, - recv(self.zmq_receiver) -> msg => { - match msg { - Ok(_) => { - if accept_block_notification { - Ok(()) - } else { - let wait_more = duration.saturating_sub(start.elapsed()); - self.wait(wait_more, accept_block_notification) - } + }, + recv(self.zmq_receiver) -> msg => { + match msg { + Ok(_) => false, + Err(_) => bail!("signal hook channel disconnected"), } - Err(_) => bail!("signal hook channel disconnected"), - } - }, - recv(after(duration)) -> _ => Ok(()), + }, + recv(after(remaining)) -> _ => true, + }; + + if deadline_expired || accept_block_notification { + return Ok(()); + } + remaining = remaining.saturating_sub(start.elapsed()); } } } diff --git a/src/util/block.rs b/src/util/block.rs index ca2da60ea..45358ae2b 100644 --- a/src/util/block.rs +++ b/src/util/block.rs @@ -45,7 +45,7 @@ pub struct HeaderEntry { } impl HeaderEntry { - #[cfg(feature = "bench")] + #[cfg(any(test, feature = "bench"))] pub fn new(height: usize, hash: BlockHash, header: BlockHeader) -> Self { Self { height, diff --git a/tests/common.rs b/tests/common.rs index 7fddf1599..7bcd19d82 100644 --- a/tests/common.rs +++ b/tests/common.rs @@ -39,6 +39,7 @@ pub struct TestRunner { daemon: Arc, mempool: Arc>, metrics: Metrics, + metrics_addr: net::SocketAddr, salt_rwlock: Arc>, } @@ -138,7 +139,8 @@ impl TestRunner { }); let signal = Waiter::start(crossbeam_channel::never()); - let metrics = Metrics::new(rand_available_addr()); + let metrics_addr = rand_available_addr(); + let metrics = Metrics::new(metrics_addr); metrics.start(); let daemon = Arc::new(Daemon::new( @@ -205,10 +207,19 @@ impl TestRunner { daemon, mempool, metrics, + metrics_addr, salt_rwlock, }) } + pub fn query(&self) -> Arc { + Arc::clone(&self.query) + } + + pub fn metrics_addr(&self) -> net::SocketAddr { + self.metrics_addr + } + pub fn node_client(&self) -> &Client { #[cfg(not(feature = "liquid"))] return &self.node.client; diff --git a/tests/query.rs b/tests/query.rs new file mode 100644 index 000000000..39e15e171 --- /dev/null +++ b/tests/query.rs @@ -0,0 +1,105 @@ +use std::net; +use std::sync::{Arc, Barrier}; +use std::thread; + +pub mod common; + +use common::Result; + +const CONCURRENT_CALLERS: usize = 10; + +fn fetch_metrics(addr: net::SocketAddr) -> String { + ureq::get(&format!("http://{}/metrics", addr)) + .call() + .expect("failed to scrape metrics") + .into_body() + .read_to_string() + .expect("failed to read metrics body") +} + +fn metric_value(body: &str, metric: &str, label_match: &str) -> u64 { + let prefix = format!("{}{{{}}} ", metric, label_match); + body.lines() + .find_map(|line| line.strip_prefix(&prefix)) + .and_then(|rest| rest.trim().parse::().ok()) + .map(|value| value as u64) + .unwrap_or(0) +} + +fn daemon_rpc_count(metrics_addr: net::SocketAddr, method: &str) -> u64 { + metric_value( + &fetch_metrics(metrics_addr), + "daemon_rpc_count", + &format!(r#"method="{}""#, method), + ) +} + +#[test] +fn test_estimate_fee_refresh_is_single_flight() -> Result<()> { + let tester = common::TestRunner::new()?; + let query = tester.query(); + let metrics_addr = tester.metrics_addr(); + + let before = daemon_rpc_count(metrics_addr, "estimatesmartfee"); + + let barrier = Arc::new(Barrier::new(CONCURRENT_CALLERS)); + let handles: Vec<_> = (0..CONCURRENT_CALLERS) + .map(|_| { + let query = query.clone(); + let barrier = Arc::clone(&barrier); + thread::spawn(move || { + barrier.wait(); + query.estimate_fee_map() + }) + }) + .collect(); + for handle in handles { + handle.join().expect("estimate_fee_map panicked"); + } + + let after = daemon_rpc_count(metrics_addr, "estimatesmartfee"); + + assert_eq!( + after - before, + 28, + "expected exactly one fee-estimate refresh (28 RPCs) from {} concurrent callers, got {}", + CONCURRENT_CALLERS, + after - before + ); + Ok(()) +} + +#[test] +fn test_get_relayfee_refresh_is_single_flight() -> Result<()> { + let tester = common::TestRunner::new()?; + let query = tester.query(); + let metrics_addr = tester.metrics_addr(); + + let before = daemon_rpc_count(metrics_addr, "getnetworkinfo"); + + let barrier = Arc::new(Barrier::new(CONCURRENT_CALLERS)); + let handles: Vec<_> = (0..CONCURRENT_CALLERS) + .map(|_| { + let query = query.clone(); + let barrier = Arc::clone(&barrier); + thread::spawn(move || { + barrier.wait(); + query.get_relayfee() + }) + }) + .collect(); + for handle in handles { + handle.join().expect("get_relayfee panicked")?; + } + + let after = daemon_rpc_count(metrics_addr, "getnetworkinfo"); + + assert_eq!( + after - before, + 1, + "expected exactly one getnetworkinfo call from {} concurrent get_relayfee() callers, got {}", + CONCURRENT_CALLERS, + after - before + ); + Ok(()) +}