From 77dc674db0656f68eaea75433971db9b0bb27df3 Mon Sep 17 00:00:00 2001 From: Oleksandr Fedorenko Date: Mon, 14 Sep 2026 05:31:29 +0300 Subject: [PATCH 01/13] fix: single-flight fee-estimate and relayfee cache refreshes --- src/new_index/query.rs | 34 +++++++++++-- tests/common.rs | 13 ++++- tests/query.rs | 105 +++++++++++++++++++++++++++++++++++++++++ 3 files changed, 148 insertions(+), 4 deletions(-) create mode 100644 tests/query.rs diff --git a/src/new_index/query.rs b/src/new_index/query.rs index 86dada56e..165cb8d0f 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}; 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,24 @@ 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 = self + .estimates_refresh + .lock() + .unwrap_or_else(|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 +290,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 +318,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/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(()) +} From aee7c21ecadd5ef60f0997eba1402d6509278d18 Mon Sep 17 00:00:00 2001 From: Oleksandr Fedorenko Date: Mon, 14 Sep 2026 05:31:29 +0300 Subject: [PATCH 02/13] fix: use try_lock for fee-estimate refresh to avoid retry herd --- src/new_index/query.rs | 11 ++++++----- 1 file changed, 6 insertions(+), 5 deletions(-) diff --git a/src/new_index/query.rs b/src/new_index/query.rs index 165cb8d0f..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, Mutex, RwLock, RwLockReadGuard}; +use std::sync::{Arc, Mutex, RwLock, RwLockReadGuard, TryLockError}; use std::time::{Duration, Instant}; use crate::chain::{Network, OutPoint, Transaction, TxOut, Txid}; @@ -260,10 +260,11 @@ impl Query { #[trace] fn refresh_fee_estimates_if_stale(&self) { - let _guard = self - .estimates_refresh - .lock() - .unwrap_or_else(|poisoned| poisoned.into_inner()); + 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; From 9a5a743bae6eeb190fc9de8799e6e2278d499c86 Mon Sep 17 00:00:00 2001 From: Edward Houston Date: Thu, 3 Sep 2026 16:25:30 +0200 Subject: [PATCH 03/13] Harden the ZMQ block notification subscriber The hashblock subscriber trusted its publisher completely. There is no transport authentication on the socket, so anyone able to reach or impersonate the endpoint could inject arbitrary 32-byte hashes, and every accepted one woke the main loop into an early tip check plus mempool update against bitcoind. Notifications are now coalesced to at most one per 250ms, which bounds that amplification and costs nothing in normal operation, since the main loop polls every 5 seconds regardless. Message handling is bounded. ZMQ_MAXMSGSIZE caps each frame at 1 KiB, ZMQ_RCVHWM caps the receive queue, and recv_multipart is replaced by a manual frame loop that drains anything past 8 frames rather than buffering it - the size cap limits how large a frame may be but not how many parts a publisher may send in a single message. The error arm now backs off from 100ms to 5 seconds instead of retrying immediately, which previously turned a persistent socket error into a busy loop and a log flood. A receive timeout lets the thread notice that the receiving end has gone away and exit. Startup returns Result instead of calling expect, so a bad address logs an error and leaves the server running on its poll interval rather than aborting the process. The sender is borrowed rather than moved, so a failed startup cannot disconnect the notification channel and take the server down with it. Waiter::wait was recursive: on the path where a caller does not accept block notifications, each one added a stack frame that only unwound when the deadline expired, so a flood could exhaust the stack. It is now a loop with identical semantics. Adds an electrs_zmq_notifications_total counter labelled by disposition, and unit tests for the frame parsing and the throttle. --- src/bin/electrs.rs | 11 +- src/new_index/zmq.rs | 313 +++++++++++++++++++++++++++++++++++++++---- src/signal.rs | 60 +++++---- 3 files changed, 327 insertions(+), 57 deletions(-) diff --git a/src/bin/electrs.rs b/src/bin/electrs.rs index 74d06f806..0b83a436c 100644 --- a/src/bin/electrs.rs +++ b/src/bin/electrs.rs @@ -60,7 +60,16 @@ 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); + // Block notifications only save the main loop from waiting out its poll + // interval, so a subscriber that will not start is worth an error but not + // an exit - we keep serving, just without the early wake-ups. + 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); diff --git a/src/new_index/zmq.rs b/src/new_index/zmq.rs index 73d5d7964..a4aa784c0 100644 --- a/src/new_index/zmq.rs +++ b/src/new_index/zmq.rs @@ -1,40 +1,297 @@ +//! 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::thread; +use std::time::{Duration, Instant}; + use bitcoin::{hashes::Hash, BlockHash}; -use crossbeam_channel::Sender; +use crossbeam_channel::{Sender, TrySendError}; +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, so let libzmq discard it before it reaches us. +const MAX_FRAME_BYTES: i64 = 1024; + +/// Upper bound on frames buffered from one multipart message. `maxmsgsize` caps +/// the size of each frame but not how many frames a publisher may send, so cap +/// that separately. +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() +} + +/// Receive one complete multipart message, buffering at most [`MAX_FRAMES`]. +/// +/// Frames past the cap are drained and discarded rather than accumulated, so a +/// publisher sending an unbounded number of parts cannot grow our memory. Returns +/// `None` when the message exceeded the cap and was therefore dropped. +fn recv_message(subscriber: &zmq::Socket) -> zmq::Result>>> { + let mut frames = Vec::new(); + let mut over_cap = false; + + loop { + let frame = subscriber.recv_msg(0)?; + if frames.len() < MAX_FRAMES { + frames.push(frame.to_vec()); + } else { + over_cap = 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 over_cap { None } else { Some(frames) }) +} + +/// 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"); + let subscriber = ctx + .socket(zmq::SUB) + .chain_err(|| format!("failed creating ZMQ subscriber for url='{url}'"))?; + subscriber - .connect(url) - .expect("failed connecting subscriber"); + .set_maxmsgsize(MAX_FRAME_BYTES) + .chain_err(|| "failed setting ZMQ maxmsgsize")?; + subscriber + .set_rcvhwm(RCVHWM) + .chain_err(|| "failed setting ZMQ rcvhwm")?; + subscriber + .set_rcvtimeo(RCVTIMEO_MS) + .chain_err(|| "failed setting ZMQ rcvtimeo")?; - // 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); - } + .connect(url) + .chain_err(|| format!("failed connecting ZMQ subscriber to url='{url}'"))?; + subscriber + .set_subscribe(TOPIC_HASHBLOCK) + .chain_err(|| "failed subscribing to ZMQ hashblock")?; + + let notifications = metrics.counter_vec( + MetricOpts::new( + "electrs_zmq_notifications_total", + "ZMQ block notifications received, by disposition", + ), + &["result"], + ); + + let url = url.to_owned(); + let block_hash_notify = block_hash_notify.clone(); + spawn_thread("zmq", move || { + subscriber_loop(&url, subscriber, block_hash_notify, notifications) + }); + + Ok(()) +} + +fn subscriber_loop( + url: &str, + subscriber: zmq::Socket, + block_hash_notify: Sender, + notifications: CounterVec, +) { + let mut throttle = Throttle::new(MIN_NOTIFY_INTERVAL); + let mut backoff = MIN_ERROR_BACKOFF; + + loop { + match recv_message(&subscriber) { + 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='{url}' err='{e}' backoff='{backoff:?}'" + ); + 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 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..753ab0d0c 100644 --- a/src/signal.rs +++ b/src/signal.rs @@ -38,38 +38,42 @@ 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) + // Iterative rather than recursive. A caller that is not accepting block + // notifications keeps waiting out the rest of its budget after each one, + // and recursing there meant one stack frame per notification - a flood of + // them could exhaust the stack before the deadline ever expired. + let mut remaining = duration; + + loop { + let start = Instant::now(); + + // `false` means we were woken by a notification rather than by the + // deadline, so the loop goes round again unless we accept those. + 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()); } } } From 6e5fe7dbe389c29a6d305bf3d984dd5b068040a1 Mon Sep 17 00:00:00 2001 From: Edward Houston Date: Thu, 3 Sep 2026 17:30:11 +0200 Subject: [PATCH 04/13] Enforce the ZMQ frame size limit ourselves, not via maxmsgsize libzmq treats a ZMQ_MAXMSGSIZE violation as a protocol error: it destroys the session and does not reconnect, and no error is returned to recv. A single oversized message would therefore disable block notifications for the lifetime of the process, silently. Check the frame length in recv_message instead, alongside the existing frame-count cap, so an oversized message is dropped without being copied out and the connection stays usable. Verified against a live regtest: one 1 MiB message previously ended all notifications permanently, and now costs one rejected counter increment. --- src/new_index/zmq.rs | 89 +++++++++++++++++++++++++++++++++++--------- 1 file changed, 72 insertions(+), 17 deletions(-) diff --git a/src/new_index/zmq.rs b/src/new_index/zmq.rs index a4aa784c0..da0b21da4 100644 --- a/src/new_index/zmq.rs +++ b/src/new_index/zmq.rs @@ -21,12 +21,19 @@ 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, so let libzmq discard it before it reaches us. -const MAX_FRAME_BYTES: i64 = 1024; - -/// Upper bound on frames buffered from one multipart message. `maxmsgsize` caps -/// the size of each frame but not how many frames a publisher may send, so cap -/// that separately. +/// something we have a use for, so it is dropped instead of being copied out. +/// +/// Enforced here rather than with `ZMQ_MAXMSGSIZE`, deliberately. libzmq treats an +/// oversized message as a *protocol* error: it tears down the session and does not +/// reconnect afterwards. A single hostile message would therefore disable block +/// notifications for the lifetime of the process, silently - no error is returned +/// to `recv`, so nothing logs and nothing retries. Checking the frame ourselves +/// costs one comparison and leaves the connection usable. +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 @@ -91,21 +98,33 @@ fn parse_hashblock(frames: &[Vec]) -> Option { BlockHash::from_slice(&reversed).ok() } -/// Receive one complete multipart message, buffering at most [`MAX_FRAMES`]. +/// Whether a frame is worth copying out of libzmq, given how many we already hold. +/// +/// Rejects a frame that is too large or that is past the count we are willing to +/// buffer. Both bounds matter: the size limit alone leaves a publisher free to send +/// unlimited small parts, and the count limit alone leaves each part unbounded. +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. /// -/// Frames past the cap are drained and discarded rather than accumulated, so a -/// publisher sending an unbounded number of parts cannot grow our memory. Returns -/// `None` when the message exceeded the cap and was therefore dropped. +/// 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 over_cap = false; + let mut rejected = false; loop { let frame = subscriber.recv_msg(0)?; - if frames.len() < MAX_FRAMES { + if frame_admissible(frame.len(), frames.len()) { frames.push(frame.to_vec()); } else { - over_cap = true; + rejected = true; } // libzmq delivers a multipart message atomically, so the remaining parts @@ -115,7 +134,7 @@ fn recv_message(subscriber: &zmq::Socket) -> zmq::Result>>> { } } - Ok(if over_cap { None } else { Some(frames) }) + Ok(if rejected { None } else { Some(frames) }) } /// Connect a SUB socket to `url` and forward `hashblock` notifications. @@ -132,9 +151,6 @@ pub fn start(url: &str, block_hash_notify: &Sender, metrics: &Metrics .socket(zmq::SUB) .chain_err(|| format!("failed creating ZMQ subscriber for url='{url}'"))?; - subscriber - .set_maxmsgsize(MAX_FRAME_BYTES) - .chain_err(|| "failed setting ZMQ maxmsgsize")?; subscriber .set_rcvhwm(RCVHWM) .chain_err(|| "failed setting ZMQ rcvhwm")?; @@ -264,6 +280,45 @@ mod tests { 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); + } + #[test] fn throttle_admits_the_first_event() { let mut throttle = Throttle::new(MIN_NOTIFY_INTERVAL); From 1d1c73e01c6d877232fdf79c0f050ad7be2a0f17 Mon Sep 17 00:00:00 2001 From: Edward Houston Date: Fri, 4 Sep 2026 14:05:42 +0200 Subject: [PATCH 05/13] Rebuild the ZMQ subscriber when libzmq drops its publisher An oversized message is a protocol error to libzmq, not a connection error: the session is terminated rather than reconnected, and nothing is reported to recv, which returns EAGAIN from then on. A single hostile message could therefore end block notifications for the lifetime of the process, silently. That teardown is the price of ZMQ_MAXMSGSIZE, and ZMQ_MAXMSGSIZE is the only bound that acts before libzmq allocates: the length prefix is checked as it is decoded, so a declared multi-gigabyte frame is refused without the body being read. Measured on libzmq 4.3.5, one 256 MiB frame grows a subscriber with no size limit by 206 MB, and the body is never handed to user code at all, so a check after recv cannot prevent it. A large enough declared length aborts the process outright. So set the size limit, and attach a socket monitor to notice the teardown it makes possible. On a disconnect the socket and its monitor are rebuilt and resubscribed. Rebuilds are spaced by the existing backoff, so a publisher sending oversized frames in a loop degrades to the tip poll instead of spinning. Adds a socket-level test that sends an oversized frame mid-stream and asserts the subscriber still receives afterwards. With the reconnect disabled it fails, which is the behaviour it is there to catch. --- src/new_index/zmq.rs | 257 +++++++++++++++++++++++++++++++++++++------ 1 file changed, 224 insertions(+), 33 deletions(-) diff --git a/src/new_index/zmq.rs b/src/new_index/zmq.rs index da0b21da4..864aeca9c 100644 --- a/src/new_index/zmq.rs +++ b/src/new_index/zmq.rs @@ -6,11 +6,13 @@ //! 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, TrySendError}; +use error_chain::ChainedError; use crate::errors::*; use crate::metrics::{CounterVec, MetricOpts, Metrics}; @@ -21,14 +23,20 @@ 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, so it is dropped instead of being copied out. +/// something we have a use for. /// -/// Enforced here rather than with `ZMQ_MAXMSGSIZE`, deliberately. libzmq treats an -/// oversized message as a *protocol* error: it tears down the session and does not -/// reconnect afterwards. A single hostile message would therefore disable block -/// notifications for the lifetime of the process, silently - no error is returned -/// to `recv`, so nothing logs and nothing retries. Checking the frame ourselves -/// costs one comparison and leaves the connection usable. +/// 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 @@ -100,9 +108,11 @@ fn parse_hashblock(frames: &[Vec]) -> Option { /// Whether a frame is worth copying out of libzmq, given how many we already hold. /// -/// Rejects a frame that is too large or that is past the count we are willing to -/// buffer. Both bounds matter: the size limit alone leaves a publisher free to send -/// unlimited small parts, and the count limit alone leaves each part unbounded. +/// 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 } @@ -137,6 +147,137 @@ fn recv_message(subscriber: &zmq::Socket) -> zmq::Result>>> { 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, +} + +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, + }) + } + + 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. + fn reconnect_if_peer_dropped(&mut self) -> bool { + if !self.peer_dropped() { + 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; + } + Err(e) => { + log::warn!( + "failed reconnecting ZMQ subscriber: 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 @@ -147,23 +288,7 @@ pub fn start(url: &str, block_hash_notify: &Sender, metrics: &Metrics log::debug!("starting ZMQ subscriber: url='{url}'"); let ctx = zmq::Context::new(); - let subscriber = ctx - .socket(zmq::SUB) - .chain_err(|| format!("failed creating ZMQ subscriber for url='{url}'"))?; - - subscriber - .set_rcvhwm(RCVHWM) - .chain_err(|| "failed setting ZMQ rcvhwm")?; - subscriber - .set_rcvtimeo(RCVTIMEO_MS) - .chain_err(|| "failed setting ZMQ rcvtimeo")?; - - subscriber - .connect(url) - .chain_err(|| format!("failed connecting ZMQ subscriber to url='{url}'"))?; - subscriber - .set_subscribe(TOPIC_HASHBLOCK) - .chain_err(|| "failed subscribing to ZMQ hashblock")?; + let subscriber = Subscriber::connect(&ctx, url)?; let notifications = metrics.counter_vec( MetricOpts::new( @@ -173,18 +298,16 @@ pub fn start(url: &str, block_hash_notify: &Sender, metrics: &Metrics &["result"], ); - let url = url.to_owned(); let block_hash_notify = block_hash_notify.clone(); spawn_thread("zmq", move || { - subscriber_loop(&url, subscriber, block_hash_notify, notifications) + subscriber_loop(subscriber, block_hash_notify, notifications) }); Ok(()) } fn subscriber_loop( - url: &str, - subscriber: zmq::Socket, + mut subscriber: Subscriber, block_hash_notify: Sender, notifications: CounterVec, ) { @@ -192,7 +315,19 @@ fn subscriber_loop( let mut backoff = MIN_ERROR_BACKOFF; loop { - match recv_message(&subscriber) { + // 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; @@ -229,7 +364,8 @@ fn subscriber_loop( Err(zmq::Error::EAGAIN) => continue, Err(e) => { log::warn!( - "ZMQ receive failed, backing off: url='{url}' err='{e}' backoff='{backoff:?}'" + "ZMQ receive failed, backing off: url='{}' err='{e}' backoff='{backoff:?}'", + subscriber.url ); thread::sleep(backoff); backoff = (backoff * 2).min(MAX_ERROR_BACKOFF); @@ -319,6 +455,61 @@ mod tests { 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" + ); + } + #[test] fn throttle_admits_the_first_event() { let mut throttle = Throttle::new(MIN_NOTIFY_INTERVAL); From 62d9a69826fd1f355cffc8199fdf99e33284ec95 Mon Sep 17 00:00:00 2001 From: Edward Houston Date: Fri, 4 Sep 2026 14:52:44 +0200 Subject: [PATCH 06/13] Retry a failed ZMQ subscriber rebuild instead of latching off The disconnect monitor is edge-triggered, and a socket whose session libzmq has terminated never emits a second DISCONNECTED. The rebuild attempt consumed that one event, so a rebuild that failed left the dead socket in place with nothing able to ask for another attempt, ending block notifications for the life of the process. Track the owed rebuild in a latched flag cleared only on success, so a transient failure - fd exhaustion on either socket the rebuild creates is the realistic one - is retried on the loop backoff. Add a test that induces a failing rebuild and asserts recovery. The existing oversized-frame test passes with and without the fix, since it only covers rebuilds that succeed. --- src/new_index/zmq.rs | 67 ++++++++++++++++++++++++++++++++++++++++++-- 1 file changed, 64 insertions(+), 3 deletions(-) diff --git a/src/new_index/zmq.rs b/src/new_index/zmq.rs index 864aeca9c..25353b983 100644 --- a/src/new_index/zmq.rs +++ b/src/new_index/zmq.rs @@ -157,6 +157,11 @@ struct Subscriber { 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 { @@ -193,6 +198,7 @@ impl Subscriber { url: url.to_owned(), socket, monitor, + dead: false, }) } @@ -211,9 +217,13 @@ impl Subscriber { /// having to distinguish it, and costs nothing when libzmq would have /// recovered on its own. /// - /// Returns whether a rebuild was attempted. + /// Returns whether a rebuild was attempted, successfully or not. fn reconnect_if_peer_dropped(&mut self) -> bool { - if !self.peer_dropped() { + // `|=` 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; } @@ -225,10 +235,15 @@ impl Subscriber { 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: url='{}' err='{}'", + "failed reconnecting ZMQ subscriber, will retry: url='{}' err='{}'", self.url, e.display_chain() ); @@ -510,6 +525,52 @@ mod tests { ); } + /// 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); From 122984d383600743a25e8d1ad596df02f6688e62 Mon Sep 17 00:00:00 2001 From: Edward Houston Date: Mon, 14 Sep 2026 15:33:52 +0200 Subject: [PATCH 07/13] chore: drop redundant implementation comments --- src/bin/electrs.rs | 3 --- src/signal.rs | 6 ------ 2 files changed, 9 deletions(-) diff --git a/src/bin/electrs.rs b/src/bin/electrs.rs index 0b83a436c..870fc52bd 100644 --- a/src/bin/electrs.rs +++ b/src/bin/electrs.rs @@ -60,9 +60,6 @@ fn run_server(config: Arc, salt_rwlock: Arc>) -> Result<( info!("starting electrs"); if let Some(zmq_addr) = config.zmq_addr.as_ref() { - // Block notifications only save the main loop from waiting out its poll - // interval, so a subscriber that will not start is worth an error but not - // an exit - we keep serving, just without the early wake-ups. 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='{}'", diff --git a/src/signal.rs b/src/signal.rs index 753ab0d0c..34f97f56d 100644 --- a/src/signal.rs +++ b/src/signal.rs @@ -38,17 +38,11 @@ impl Waiter { } pub fn wait(&self, duration: Duration, accept_block_notification: bool) -> Result<()> { - // Iterative rather than recursive. A caller that is not accepting block - // notifications keeps waiting out the rest of its budget after each one, - // and recursing there meant one stack frame per notification - a flood of - // them could exhaust the stack before the deadline ever expired. let mut remaining = duration; loop { let start = Instant::now(); - // `false` means we were woken by a notification rather than by the - // deadline, so the loop goes round again unless we accept those. let deadline_expired = select! { recv(self.receiver) -> msg => { match msg { From 304635cae8ed3b187321ae4e97294b5f659b86a6 Mon Sep 17 00:00:00 2001 From: Philippe McLean Date: Sat, 12 Sep 2026 14:27:14 -0700 Subject: [PATCH 08/13] capture unknown methods as distinct metrics label - avoid creating new labels in prometheus for unknown method calls - add separate bucket for unknown in histogram - downgrade log from warn to debug for unknown method calls --- src/electrum/server.rs | 19 ++++++++++++------- 1 file changed, 12 insertions(+), 7 deletions(-) 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) => { From eba47e99825c7eca0116becf4c083e2a3d00358c Mon Sep 17 00:00:00 2001 From: Edward Houston Date: Mon, 24 Aug 2026 13:06:25 +0200 Subject: [PATCH 09/13] fix(daemon): verify RPC responses; propagate errors instead of aborting Addresses three related findings on the electrs to bitcoind RPC path where daemon-supplied data was accepted by assert/expect/panic and, under panic=abort, terminated the process. src/daemon.rs - getblock now returns Err on hash mismatch instead of asserting. - getblocks verifies each returned block against the requested hash by position, with a count check first. Non-whitelisted error and retry-exhaustion arms now propagate Err instead of panicking, so the caller can decide policy rather than the process dying. - get_all_headers replaces expect/assert/assert_eq on daemon-supplied data with propagated errors, and caps tip_height at 100_000_000 before allocating the heights vector so a hostile daemon cannot drive a huge allocation before any network check. src/new_index/fetch.rs - bitcoind_fetcher logs and returns from the fetcher thread when getblocks fails instead of panicking. Under panic=abort a panic here killed electrs; on restart the same fetch was retried and, for a pruned deep-reorg or a persistent daemon fault, the same panic fired again, producing a crash loop only manual intervention could break. Log-and-return turns that into a stall the next Indexer.update tick can retry. - parse_blocks adds a bounds check on the block-body slice end position so a truncated tail block from a bitcoind crash no longer panics inside the rayon pool. Deserialize errors are propagated via ? into a collected Result rather than expect. - blkfiles_parser logs and skips the current blob on parse failure so one corrupted blk file cannot abort initial sync; anything not found on the BlkFiles path is picked up on the Bitcoind fetch switchover. --- src/daemon.rs | 82 +++++++++++++++++++++++++++++++++--------- src/new_index/fetch.rs | 69 ++++++++++++++++++++++++++--------- 2 files changed, 119 insertions(+), 32 deletions(-) diff --git a/src/daemon.rs b/src/daemon.rs index d0b207370..5aafcf2da 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,39 @@ 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") } 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") + bail!("failed to get blocks from bitcoind attempts='0'"); } 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 +1489,35 @@ impl Daemon { #[trace] fn get_all_headers(&self, tip: &BlockHash) -> Result> { + const MAX_TIP_HEIGHT: u64 = 100_000_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; + .and_then(|v| v.as_u64()) + .ok_or_else(|| format!("bitcoind returned malformed getblockheader info='{:?}'", info))? + as u64; + if tip_height > MAX_TIP_HEIGHT { + bail!( + "bitcoind returned implausible tip_height='{}' cap='{}'", + tip_height, + MAX_TIP_HEIGHT + ); + } + let tip_height = tip_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()); - - result.append(&mut headers); + 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 +1529,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/new_index/fetch.rs b/src/new_index/fetch.rs index 0dc92aeaa..8e1806b4f 100644 --- a/src/new_index/fetch.rs +++ b/src/new_index/fetch.rs @@ -103,10 +103,26 @@ 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 = match daemon.getblocks(&blockhashes) { + Ok(blocks) => blocks, + Err(e) => { + log::error!( + "bitcoind_fetcher stopping, will retry on next index update err='{:?}' first_hash='{:?}' batch_size='{}'", + e, + blockhashes.first(), + blockhashes.len() + ); + return; + } + }; + if blocks.len() != entries.len() { + log::error!( + "bitcoind_fetcher block/entry count mismatch expected='{}' got='{}'", + entries.len(), + blocks.len() + ); + return; + } let block_entries: Vec = blocks .into_iter() .zip(entries) @@ -120,10 +136,10 @@ fn bitcoind_fetcher( } }) .collect(); - assert_eq!(block_entries.len(), entries.len()); - sender - .send(block_entries) - .expect("failed to send fetched blocks"); + if sender.send(block_entries).is_err() { + log::warn!("bitcoind_fetcher receiver dropped, stopping"); + return; + } log::debug!("last fetch {:?}", entries.last()); } }), @@ -247,10 +263,21 @@ fn blkfiles_parser(blobs: Fetcher>, magic: u32) -> Fetcher blocks, + Err(e) => { + log::warn!( + "blkfiles_parser skipping blob err='{:?}'", + e + ); + return; + } + }; + if sender.send(blocks).is_err() { + log::warn!("blkfiles_parser receiver dropped, stopping"); + } }); }), ) @@ -277,6 +304,12 @@ fn parse_blocks(pool: &rayon::ThreadPool, blob: Vec, 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 +327,14 @@ fn parse_blocks(pool: &rayon::ThreadPool, blob: Vec, magic: u32) -> Result>>() + }) } From 86f005210cd347d7d77db08fcb2283dc2ea67e43 Mon Sep 17 00:00:00 2001 From: Edward Houston Date: Mon, 7 Sep 2026 12:06:38 +0200 Subject: [PATCH 10/13] fix(index): preserve synced prefix on fetch failure Propagate fetcher errors through the pipeline while retaining successfully processed batches. Advance the persisted and in-memory chain tips only through the contiguous added-and-indexed prefix, allowing incomplete forward syncs to retry safely. Fall back from blk files to bitcoind after an incomplete fetch, while keeping reorg fetch failures fatal. --- src/new_index/fetch.rs | 152 +++++++++++++++++++++++++++------------- src/new_index/schema.rs | 83 ++++++++++++++++++---- 2 files changed, 175 insertions(+), 60 deletions(-) diff --git a/src/new_index/fetch.rs b/src/new_index/fetch.rs index 8e1806b4f..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,25 +114,19 @@ fn bitcoind_fetcher( fetcher_count += 1; let blockhashes: Vec = entries.iter().map(|e| *e.hash()).collect(); - let blocks = match daemon.getblocks(&blockhashes) { - Ok(blocks) => blocks, - Err(e) => { - log::error!( - "bitcoind_fetcher stopping, will retry on next index update err='{:?}' first_hash='{:?}' batch_size='{}'", - e, - blockhashes.first(), - blockhashes.len() - ); - return; - } - }; + 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() { - log::error!( + bail!( "bitcoind_fetcher block/entry count mismatch expected='{}' got='{}'", entries.len(), blocks.len() ); - return; } let block_entries: Vec = blocks .into_iter() @@ -136,12 +141,12 @@ fn bitcoind_fetcher( } }) .collect(); - if sender.send(block_entries).is_err() { - log::warn!("bitcoind_fetcher receiver dropped, stopping"); - return; - } + sender + .send(block_entries) + .chain_err(|| "bitcoind_fetcher receiver dropped")?; log::debug!("last fetch {:?}", entries.last()); } + Ok(()) }), )) } @@ -193,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(()) }), )) } @@ -227,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(()) }), ) } @@ -263,22 +271,13 @@ fn blkfiles_parser(blobs: Fetcher>, magic: u32) -> Fetcher blocks, - Err(e) => { - log::warn!( - "blkfiles_parser skipping blob err='{:?}'", - e - ); - return; - } - }; - if sender.send(blocks).is_err() { - log::warn!("blkfiles_parser receiver dropped, stopping"); - } - }); + let blocks = parse_blocks(&pool, blob, magic)?; + sender + .send(blocks) + .chain_err(|| "blkfiles_parser receiver dropped")?; + Ok(()) + })?; + Ok(()) }), ) } @@ -338,3 +337,62 @@ 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/schema.rs b/src/new_index/schema.rs index 6ee9b8732..31e8cdd8f 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, @@ -280,6 +284,15 @@ impl Indexer { .collect() } + fn completed_header_prefix(&self, new_headers: &[HeaderEntry]) -> Vec { + let added = self.store.added_blockhashes.read().unwrap(); + let indexed = self.store.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() + } + fn start_auto_compactions(&self, db: &DB) { let key = b"F".to_vec(); if db.get(&key).is_none() { @@ -358,7 +371,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. @@ -385,7 +401,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 +456,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.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 +505,30 @@ 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()); + headers.append(completed_headers); + let synced_tip = *headers.tip(); - if let FetchFrom::BlkFiles = self.from { - self.from = FetchFrom::Bitcoind; + // Update the synced tip only after all db writes are flushed and only as far + // as the contiguous prefix that was both added and indexed. + debug!("updating synced tip to {:?}", synced_tip); + if !headers.is_empty() { + self.store + .txstore_db + .put_sync(b"t", &serialize(&synced_tip)); + } + + if fetch_result.is_ok() { + assert_eq!(tip, synced_tip); } - self.tip_metric.set(headers.best_height() as i64); + if !headers.is_empty() { + self.tip_metric.set(headers.best_height() as i64); + } - Ok(tip) + Ok(synced_tip) } fn add(&self, blocks: &[BlockEntry]) { @@ -2001,6 +2052,12 @@ 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); + } + #[test] fn test_compute_script_hash_p2pkh() { // P2PKH scriptPubKey for address 1A1zP1eP5QGefi2DMPTfTL5SLmv7DivfNa From 37b33205936b9149a964f2b44e36ada6346d51d6 Mon Sep 17 00:00:00 2001 From: Edward Houston Date: Tue, 8 Sep 2026 13:20:49 +0200 Subject: [PATCH 11/13] fix(index): bound startup retry and fail loudly on an empty index Startup no longer completes when a block cannot be fetched. Now that indexer.update returns Ok at a short prefix, the pre-listener mempool catch-up loop retried forever: the tip stayed behind, Mempool::update kept returning false, and neither the REST nor the Electrum listener bound. Break out when an index update makes no progress and serve the prefix that is indexed; the 5 second main loop keeps retrying. A total fetch failure on an empty database reported success at the all-zeroes hash, because HeaderList::tip returns the default hash for an empty list. Move the append, the empty check and the t marker write into Store::advance_synced_tip, which errors instead. Also replace the assert_eq on the daemon tip with a bail so it is not a process abort under panic = abort, carry the underlying error into the attempts exhausted bail in getblocks, put the retry warning in the var= form, and drop MAX_TIP_HEIGHT in favour of chunking heights with step_by so get_all_headers no longer allocates one usize per block. Move headers_to_process and completed_header_prefix onto Store, since both read only Store state, and cover them with unit tests over a temporary RocksDB. Reverting the prefix logic to new_headers.to_vec fails three of the four new tests; dropping the empty index guard fails the fourth. --- src/bin/electrs.rs | 16 ++- src/daemon.rs | 40 ++++--- src/new_index/db.rs | 2 +- src/new_index/schema.rs | 244 +++++++++++++++++++++++++++++++++------- src/util/block.rs | 2 +- 5 files changed, 242 insertions(+), 62 deletions(-) diff --git a/src/bin/electrs.rs b/src/bin/electrs.rs index 870fc52bd..93b59f186 100644 --- a/src/bin/electrs.rs +++ b/src/bin/electrs.rs @@ -114,10 +114,20 @@ fn run_server(config: Arc, salt_rwlock: Arc>) -> Result<( Arc::clone(&config), ))); + // Mempool syncing is aborted whenever the chain tip moves, so index the new block(s) + // and try again. An update that cannot advance the tip never will here, so start up + // on the prefix that is indexed rather than spin with no listeners bound - the main + // loop below keeps retrying the outstanding blocks every 5 seconds. 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 5aafcf2da..0ba6a5865 100644 --- a/src/daemon.rs +++ b/src/daemon.rs @@ -1305,15 +1305,22 @@ impl Daemon { || 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 { bail!("failed to get blocks from bitcoind err='{}'", err_msg); } } } - if attempts == 0 { - bail!("failed to get blocks from bitcoind attempts='0'"); - } std::thread::sleep(RETRY_WAIT_DURATION); }; if values.len() != blockhashes.len() { @@ -1489,26 +1496,23 @@ impl Daemon { #[trace] fn get_all_headers(&self, tip: &BlockHash) -> Result> { - const MAX_TIP_HEIGHT: u64 = 100_000_000; + const CHUNK_SIZE: usize = 100_000; let info: Value = self.request("getblockheader", json!([tip]))?; let tip_height = info .get("height") .and_then(|v| v.as_u64()) - .ok_or_else(|| format!("bitcoind returned malformed getblockheader info='{:?}'", info))? - as u64; - if tip_height > MAX_TIP_HEIGHT { - bail!( - "bitcoind returned implausible tip_height='{}' cap='{}'", - tip_height, - MAX_TIP_HEIGHT - ); - } - let tip_height = tip_height as usize; - let all_heights: Vec = (0..=tip_height).collect(); - let chunk_size = 100_000; + .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}'") + })?; + + // Materialise one chunk of heights at a time, rather than every height up front, + // so the allocation stays bounded regardless of the height bitcoind reports. let mut result = vec![]; - for heights in all_heights.chunks(chunk_size) { + 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!( 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/schema.rs b/src/new_index/schema.rs index 31e8cdd8f..2ef6ce74d 100644 --- a/src/new_index/schema.rs +++ b/src/new_index/schema.rs @@ -142,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; @@ -273,26 +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 completed_header_prefix(&self, new_headers: &[HeaderEntry]) -> Vec { - let added = self.store.added_blockhashes.read().unwrap(); - let indexed = self.store.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() - } - fn start_auto_compactions(&self, db: &DB) { let key = b"F".to_vec(); if db.get(&key).is_none() { @@ -391,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(), @@ -470,7 +499,7 @@ impl Indexer { ); } - let completed_headers = self.completed_header_prefix(&new_headers); + 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={}", @@ -507,26 +536,19 @@ impl Indexer { // 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(completed_headers); - let synced_tip = *headers.tip(); + // Done only once all the db writes above are flushed. + let synced_tip = self.store.advance_synced_tip(completed_headers)?; - // Update the synced tip only after all db writes are flushed and only as far - // as the contiguous prefix that was both added and indexed. - debug!("updating synced tip to {:?}", synced_tip); - if !headers.is_empty() { - self.store - .txstore_db - .put_sync(b"t", &serialize(&synced_tip)); - } - - if fetch_result.is_ok() { - assert_eq!(tip, synced_tip); + 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 + ); } - if !headers.is_empty() { - self.tip_metric.set(headers.best_height() as i64); - } + self.tip_metric + .set(self.store.headers().best_height() as i64); Ok(synced_tip) } @@ -2058,6 +2080,150 @@ mod tests { 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/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, From 36069c60d70417fd3a01f4ea523fa02cf4de853e Mon Sep 17 00:00:00 2001 From: Edward Houston Date: Tue, 8 Sep 2026 13:38:36 +0200 Subject: [PATCH 12/13] style: trim inline comments to the non-obvious intent Reviewer feedback on this file set was to keep inline comments out unless the intent or a gotcha is not apparent from the code. Two comments added with the chunked header fetch and the startup retry bound restated what the code already shows. Cut each to the part a reader cannot infer: that the chunking is an allocation bound rather than a batching optimisation, and that the startup loop runs before any listener is bound. --- src/bin/electrs.rs | 6 ++---- src/daemon.rs | 3 +-- 2 files changed, 3 insertions(+), 6 deletions(-) diff --git a/src/bin/electrs.rs b/src/bin/electrs.rs index 93b59f186..25e4e50c7 100644 --- a/src/bin/electrs.rs +++ b/src/bin/electrs.rs @@ -114,10 +114,8 @@ fn run_server(config: Arc, salt_rwlock: Arc>) -> Result<( Arc::clone(&config), ))); - // Mempool syncing is aborted whenever the chain tip moves, so index the new block(s) - // and try again. An update that cannot advance the tip never will here, so start up - // on the prefix that is indexed rather than spin with no listeners bound - the main - // loop below keeps retrying the outstanding blocks every 5 seconds. + // 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)? { let new_tip = indexer.update(&daemon)?; if new_tip == tip { diff --git a/src/daemon.rs b/src/daemon.rs index 0ba6a5865..7e4442e73 100644 --- a/src/daemon.rs +++ b/src/daemon.rs @@ -1507,8 +1507,7 @@ impl Daemon { format!("bitcoind returned out-of-range tip_height='{tip_height}'") })?; - // Materialise one chunk of heights at a time, rather than every height up front, - // so the allocation stays bounded regardless of the height bitcoind reports. + // 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); From 39529a25e29baa41a838d27a795670711b79c621 Mon Sep 17 00:00:00 2001 From: Edward Houston Date: Tue, 8 Sep 2026 17:55:03 +0200 Subject: [PATCH 13/13] fix(sync): stop reporting initial sync complete at a partial tip Indexer::update returns Ok at the completed prefix when a fetch fails partway, so the startup log asserted completion for a chain that never reached the daemon tip. Compare against the daemon tip and log the partial case for what it is, naming both tips. --- src/bin/electrs.rs | 10 +++++++++- 1 file changed, 9 insertions(+), 1 deletion(-) diff --git a/src/bin/electrs.rs b/src/bin/electrs.rs index 25e4e50c7..6f94d028b 100644 --- a/src/bin/electrs.rs +++ b/src/bin/electrs.rs @@ -92,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),