From 97f8988aa2806a1e293e7525c2932ae5ebf50681 Mon Sep 17 00:00:00 2001 From: Benjamen Keroack Date: Mon, 28 Sep 2026 11:40:05 -0600 Subject: [PATCH 01/11] node: --health-listen serves /healthz from the start of run_p2p RPC, Electrum, and Esplora bind only after the store opens, catch-up finishes, and the scripthash index materializes. On a synced mainnet node those took 42 minutes (a schema backfill), 9 hours (IBD), and 33 minutes. A liveness probe on any of those listeners would restart the process in the middle of that work. --health-listen [ADDR] (conf health_listen, bare 127.0.0.1:9332) binds before run_node opens the store and answers GET /healthz with 200 in every phase. Other methods are 405 and other paths 404. Concurrency, body, and timeout limits are always on, as on Esplora. A bind failure stops the node before it touches the datadir. Co-Authored-By: Claude Opus 5.5 --- Cargo.lock | 3 + crates/rbitcoin-node/Cargo.toml | 3 + crates/rbitcoin-node/src/cli.rs | 13 +++- crates/rbitcoin-node/src/config.rs | 16 ++++- crates/rbitcoin-node/src/health.rs | 67 +++++++++++++++++++++ crates/rbitcoin-node/src/lib.rs | 1 + crates/rbitcoin-node/src/run.rs | 8 +++ crates/rbitcoin-primitives/src/lib.rs | 2 + crates/rbitcoin-test/tests/cross_surface.rs | 33 +++++++++- crates/rbitcoin-test/tests/scenarios.rs | 1 + 10 files changed, 143 insertions(+), 4 deletions(-) create mode 100644 crates/rbitcoin-node/src/health.rs diff --git a/Cargo.lock b/Cargo.lock index d2f0bbee3..16c449680 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -764,6 +764,7 @@ dependencies = [ name = "rbitcoin-node" version = "0.7.99" dependencies = [ + "axum", "bitcoin", "bitcoin_hashes 0.14.101", "getrandom 0.4.3", @@ -781,6 +782,8 @@ dependencies = [ "rbitcoin-store", "serde_json", "tokio", + "tower", + "tower-http", ] [[package]] diff --git a/crates/rbitcoin-node/Cargo.toml b/crates/rbitcoin-node/Cargo.toml index a299e2537..aa2089979 100644 --- a/crates/rbitcoin-node/Cargo.toml +++ b/crates/rbitcoin-node/Cargo.toml @@ -25,6 +25,9 @@ rbitcoin-esplora = { workspace = true } bitcoin = { workspace = true } bitcoin_hashes = { workspace = true } serde_json = { workspace = true } +axum = { workspace = true } +tower = { workspace = true } +tower-http = { workspace = true } tokio = { workspace = true } mimalloc = { workspace = true } getrandom = "0.4" diff --git a/crates/rbitcoin-node/src/cli.rs b/crates/rbitcoin-node/src/cli.rs index fc67de1b4..b8ff7d088 100644 --- a/crates/rbitcoin-node/src/cli.rs +++ b/crates/rbitcoin-node/src/cli.rs @@ -302,7 +302,7 @@ fn operator_usage() -> String { [--listen ADDR] [--no-listen] [--connect ADDR]... [--seed-node HOST]... [--proxy HOST:PORT] [--onion HOST:PORT] [--proxy-randomize[=0|1]] [--only-net NET]... \\\n\ [--tor-control [HOST:PORT]] [--tor-control-cookie PATH] [--tor-control-password PASS] \\\n\ [--i2p-sam [HOST:PORT]] [--i2p-accept-incoming] \\\n\ - [--electrum-listen ADDR] [--esplora-listen ADDR] [--esplora-onion[=0|1]] \\\n\ + [--electrum-listen ADDR] [--esplora-listen ADDR] [--esplora-onion[=0|1]] [--health-listen [ADDR]] \\\n\ [--sh-index] [--block-filter-index] [--prune-seqsigwit] [--prune-seqsigwit-ram-threshold-bytes N] [--sp-tweaks] [--sp-tweaks-dust SATS] [--max-sh-creates N] [--esplora-block-template] \\\n\ [--rpc] [--rpc-listen [ADDR]] [--rpc-socket PATH] [--rpc-token-file PATH] [--rpc-cookie-file PATH] [--rpc-work-queue N] \\\n\ [--milestone HEIGHT] \\\n\ @@ -353,6 +353,8 @@ Block filters: --block-filter-index (default off) builds BIP158 basic filters. I Silent payments: --sp-tweaks (default off) writes/serves the thin BIP-352 tweak index.\n\ Not with --prune-seqsigwit (tweaks read scriptSig and witness).\n\ --sp-tweaks-dust SATS omits served P2TR outs with value <= SATS (default 1000; 0 = all; 546 = Cake electrs).\n\ +Health: --health-listen [ADDR] serves GET /healthz from the first second of startup\n\ + (default 127.0.0.1:9332). Unauthenticated; keep it on loopback or a probe-only network.\n\ RPC: --rpc unix socket {{datadir}}/rpc.sock; --rpc-listen [ADDR] adds TCP (default 127.0.0.1 and Core-matching port). Token {{datadir}}/rpc.token (Bearer); --rpc-cookie-file opts TCP into Core cookie HTTP Basic. No --rpcuser.\n\ Cold files: --datadir-cold PATH puts Class A seqsigwit.body/idx under PATH/store (HDD).\n\ Default (flag omitted): hot and cold files both live under --datadir.\n\ @@ -418,7 +420,12 @@ fn is_bool_key(key: &str) -> bool { fn is_optional_addr_key(key: &str) -> bool { matches!( key, - "rpc_listen" | "electrum_listen" | "esplora_listen" | "tor_control" | "i2p_sam" + "rpc_listen" + | "electrum_listen" + | "esplora_listen" + | "health_listen" + | "tor_control" + | "i2p_sam" ) } @@ -809,6 +816,7 @@ mod tests { "--rpc-listen", "--electrum-listen", "--esplora-listen", + "--health-listen", ]); assert!(cfg.shindex && cfg.sptweaks); assert_eq!(cfg.sptweaks_dust, 546); @@ -820,6 +828,7 @@ mod tests { assert!(cfg.rpc.socket); assert_eq!(cfg.rpc.listen, Some("127.0.0.1:18443".parse().unwrap())); assert_eq!(cfg.listen.electrum.unwrap().port(), 50001); + assert_eq!(cfg.listen.health, Some("127.0.0.1:9332".parse().unwrap())); match cfg.listen.esplora.unwrap() { rbitcoin_esplora::EsploraListen::Tcp(a) => assert_eq!(a.port(), 3000), #[cfg(unix)] diff --git a/crates/rbitcoin-node/src/config.rs b/crates/rbitcoin-node/src/config.rs index bdadad0ec..25a0b97f5 100644 --- a/crates/rbitcoin-node/src/config.rs +++ b/crates/rbitcoin-node/src/config.rs @@ -3,7 +3,9 @@ use bitcoin::hex::FromHex; use bitcoin::ScriptBuf; use rbitcoin_consensus::{mainnet_milestone_anchor, ChainParams, Milestone}; use rbitcoin_esplora::EsploraListen; -use rbitcoin_primitives::{Network, DEFAULT_ELECTRUM_PORT, DEFAULT_ESPLORA_PORT}; +use rbitcoin_primitives::{ + Network, DEFAULT_ELECTRUM_PORT, DEFAULT_ESPLORA_PORT, DEFAULT_HEALTH_PORT, +}; use rbitcoin_store::HeadScale; use std::net::SocketAddr; use std::path::{Path, PathBuf}; @@ -85,6 +87,8 @@ pub struct ListenOpts { pub p2p_extra: Vec, pub electrum: Option, pub esplora: Option, + /// `--health-listen`: `/healthz` from the start of `run_p2p`. + pub health: Option, pub connect: Vec, /// `--connect` names that are not a `NetAddr` (clearnet DNS, Warnet tanks). pub connect_dns: Vec, @@ -122,6 +126,7 @@ impl Default for ListenOpts { p2p_extra: Vec::new(), electrum: None, esplora: None, + health: None, connect: Vec::new(), connect_dns: Vec::new(), seednodes: Vec::new(), @@ -1007,6 +1012,14 @@ impl NodeConfig { .map_err(|e| NodeError::Config(format!("conf esplora_listen: {e}")))? }); } + "health_listen" => { + self.listen.health = Some(if val.is_empty() { + SocketAddr::from(([127, 0, 0, 1], DEFAULT_HEALTH_PORT)) + } else { + val.parse() + .map_err(|e| NodeError::Config(format!("conf health_listen: {e}")))? + }); + } "sh_index" => { self.shindex = parse_conf_bool(val) .map_err(|e| NodeError::Config(format!("conf sh_index: {e}")))?; @@ -2113,6 +2126,7 @@ mod tests { ("listen=not-an-addr\n", "listen"), ("electrum_listen=bad\n", "electrum"), ("esplora_listen=bad\n", "esplora"), + ("health_listen=bad\n", "health"), ("mempool_size_mb=0\n", "mempool"), ("log_level=\n", "log_level"), ("network=notanet\n", "network"), diff --git a/crates/rbitcoin-node/src/health.rs b/crates/rbitcoin-node/src/health.rs new file mode 100644 index 000000000..af9bcd8ad --- /dev/null +++ b/crates/rbitcoin-node/src/health.rs @@ -0,0 +1,67 @@ +//! Health listener (`--health-listen`): unauthenticated `GET /healthz` for +//! process probes. +//! +//! It binds at the top of [`crate::run_p2p`], before the store opens, so it +//! answers through a schema migration, catch-up, and index materialize. The +//! RPC, Electrum, and Esplora listeners bind only after those. + +use axum::http::StatusCode; +use axum::routing::get; +use axum::Router; +use rbitcoin_log::{info, warn}; +use std::net::SocketAddr; +use std::time::Duration; +use tokio::net::TcpListener; +use tokio::task::JoinHandle; +use tower::limit::ConcurrencyLimitLayer; +use tower_http::limit::RequestBodyLimitLayer; +use tower_http::timeout::TimeoutLayer; + +/// Probes in flight at once. Probers send one request per period. +const MAX_CONCURRENT: usize = 16; +/// Per-request wall. Probe timeouts are usually 1s. +const REQUEST_TIMEOUT: Duration = Duration::from_secs(5); + +/// Bound health listener. Dropping it stops serving. +pub(crate) struct HealthHandle { + task: JoinHandle<()>, +} + +impl Drop for HealthHandle { + fn drop(&mut self) { + self.task.abort(); + } +} + +/// Bind `addr` and serve the health routes until the handle drops. +pub(crate) async fn run_health(addr: SocketAddr) -> std::io::Result { + let listener = TcpListener::bind(addr).await?; + let local_addr = listener.local_addr()?; + if !local_addr.ip().is_loopback() { + warn!("health: {local_addr} is not loopback; /healthz is unauthenticated"); + } + let app = router(); + let task = tokio::spawn(async move { + if let Err(e) = axum::serve(listener, app).await { + warn!("health: serve ended: {e}"); + } + }); + info!("health HTTP on {local_addr}"); + Ok(HealthHandle { task }) +} + +fn router() -> Router { + Router::new() + .route("/healthz", get(healthz)) + // Outer → inner: concurrency → body (GET only, so none) → timeout. + .layer(TimeoutLayer::with_status_code( + StatusCode::REQUEST_TIMEOUT, + REQUEST_TIMEOUT, + )) + .layer(RequestBodyLimitLayer::new(0)) + .layer(ConcurrencyLimitLayer::new(MAX_CONCURRENT)) +} + +async fn healthz() -> &'static str { + "ok\n" +} diff --git a/crates/rbitcoin-node/src/lib.rs b/crates/rbitcoin-node/src/lib.rs index fafb39c74..0b47b704c 100644 --- a/crates/rbitcoin-node/src/lib.rs +++ b/crates/rbitcoin-node/src/lib.rs @@ -3,6 +3,7 @@ mod cli; mod config; mod error; +mod health; mod inhibit; mod lock; mod regtest_rpc; diff --git a/crates/rbitcoin-node/src/run.rs b/crates/rbitcoin-node/src/run.rs index 73903d5f9..f5b9d6df9 100644 --- a/crates/rbitcoin-node/src/run.rs +++ b/crates/rbitcoin-node/src/run.rs @@ -217,6 +217,14 @@ pub fn run_node(config: NodeConfig) -> Result { /// hold the process). #[allow(clippy::cognitive_complexity)] // node bring-up / P2P follow loop pub async fn run_p2p(config: NodeConfig) -> Result<(), NodeError> { + let _health = match config.listen.health { + Some(addr) => Some( + crate::health::run_health(addr) + .await + .map_err(|e| NodeError::Config(format!("health listen {addr}: {e}")))?, + ), + None => None, + }; let handle = run_node(config.clone())?; let params = config.chain_params()?; let milestone = config.milestone(); diff --git a/crates/rbitcoin-primitives/src/lib.rs b/crates/rbitcoin-primitives/src/lib.rs index 129bb828f..f2edc3e57 100644 --- a/crates/rbitcoin-primitives/src/lib.rs +++ b/crates/rbitcoin-primitives/src/lib.rs @@ -72,6 +72,8 @@ pub fn rbitcoin_subversion( pub const DEFAULT_ELECTRUM_PORT: u16 = 50001; /// Default Esplora HTTP port when `--esplora-listen` omits ADDR. pub const DEFAULT_ESPLORA_PORT: u16 = 3000; +/// Default health HTTP port when `--health-listen` omits ADDR. +pub const DEFAULT_HEALTH_PORT: u16 = 9332; pub const STORE_MAGIC: [u8; 4] = *b"RBT1"; diff --git a/crates/rbitcoin-test/tests/cross_surface.rs b/crates/rbitcoin-test/tests/cross_surface.rs index a950fa92b..fefeff2fd 100644 --- a/crates/rbitcoin-test/tests/cross_surface.rs +++ b/crates/rbitcoin-test/tests/cross_surface.rs @@ -158,6 +158,15 @@ async fn jsonrpc_unix(path: &std::path::Path, method: &str, params: Value) -> Va serde_json::from_str(json).unwrap_or_else(|e| panic!("unix rpc {method} json: {e} body={text}")) } +/// `--health-listen` answers `GET /healthz` in every phase; nothing else is a route. +async fn pin_healthz(health_addr: SocketAddr) { + assert_eq!(http_get(health_addr, "/healthz").await, (200, "ok".into())); + let (st, body) = http_post(health_addr, "/healthz", "").await; + assert_eq!(st, 405, "POST /healthz: {body}"); + let (st, body) = http_get(health_addr, "/nope").await; + assert_eq!(st, 404, "unknown health path: {body}"); +} + async fn pin_address_prefix_404(esplora_addr: SocketAddr) { let (st, body) = http_get(esplora_addr, "/address-prefix/bc1").await; assert_eq!(st, 404, "address-prefix stays 404: {body}"); @@ -1033,6 +1042,25 @@ async fn fee_history_backfills_from_the_chain_when_relay_starts() { } } +/// The health listener binds before the store opens: a taken port stops the +/// node with nothing written to the datadir. +#[tokio::test(flavor = "multi_thread")] +async fn health_listen_bind_failure_stops_before_store_open() { + let td = TestDatadir::new().unwrap(); + let taken = std::net::TcpListener::bind("127.0.0.1:0").unwrap(); + let mut cfg = NodeConfig::default() + .with_datadir(td.path()) + .with_network(Network::Regtest) + .with_tiny_heads(); + cfg.listen.health = Some(taken.local_addr().unwrap()); + let err = run_p2p(cfg).await.unwrap_err().to_string(); + assert!(err.contains("health listen"), "{err}"); + assert!( + !td.store_path().exists(), + "store opened before the health bind" + ); +} + #[tokio::test(flavor = "multi_thread")] async fn esplora_broadcast_visible_in_rpc_and_electrum() { let td = TestDatadir::new().unwrap(); @@ -1059,6 +1087,7 @@ async fn esplora_broadcast_visible_in_rpc_and_electrum() { let electrum_addr = ephemeral_addr(); let esplora_addr = ephemeral_addr(); let rpc_addr = ephemeral_addr(); + let health_addr = ephemeral_addr(); let mut cfg = NodeConfig::default() .with_datadir(td.path()) @@ -1072,6 +1101,7 @@ async fn esplora_broadcast_visible_in_rpc_and_electrum() { cfg.listen.electrum = Some(electrum_addr); cfg.listen.esplora = Some(rbitcoin_esplora::EsploraListen::Tcp(esplora_addr)); cfg.rpc.listen = Some(rpc_addr); + cfg.listen.health = Some(health_addr); // mempool's CORE_RPC.SOCKET_PATH reaches the node from another user. let rpc_sock = td.path().join("run").join("rpc.sock"); cfg.apply_kv("rpc_socket", rpc_sock.to_str().unwrap()) @@ -1080,7 +1110,8 @@ async fn esplora_broadcast_visible_in_rpc_and_electrum() { cfg.max_run_secs = Some(90); let node = tokio::spawn(run_p2p(cfg)); - wait_listeners(&[electrum_addr, esplora_addr, rpc_addr]).await; + wait_listeners(&[electrum_addr, esplora_addr, rpc_addr, health_addr]).await; + pin_healthz(health_addr).await; pin_address_prefix_404(esplora_addr).await; #[cfg(unix)] { diff --git a/crates/rbitcoin-test/tests/scenarios.rs b/crates/rbitcoin-test/tests/scenarios.rs index 4f242ab44..c6e8a2834 100644 --- a/crates/rbitcoin-test/tests/scenarios.rs +++ b/crates/rbitcoin-test/tests/scenarios.rs @@ -124,6 +124,7 @@ fn pin_argv_usage_errors() { &["--log-level"], &["--log-level", "loud"], &["--electrum-listen", "bad"], + &["--health-listen", "bad"], &["--api-log"], &["--asmap"], &["--conf"], From e6f4831437308b9c5ff8686d27fb19060e2f4823 Mon Sep 17 00:00:00 2001 From: Benjamen Keroack Date: Mon, 28 Sep 2026 11:47:29 -0600 Subject: [PATCH 02/11] node: /readyz follows the node's phase, IBD, and index lag /readyz is 503 "not ready: " until the node can serve, then 200. The first failing check names the reason: - the bring-up phase (opening, starting, catch-up, indexing, stopping), written by run_p2p at each transition; - a configured RPC, Electrum, or Esplora listener that did not bind (their start only warns, so the node would otherwise follow the tip while a probe routes clients to a dead port); - initial block download, read from the same ChainHub::in_ibd that getblockchaininfo.initialblockdownload reports; - the tip more than 6 blocks behind the best header; - with --sh-index, the scripthash index more than 6 blocks behind the tip. The chain reads can touch the store, so /readyz takes them on the blocking pool. /healthz stays 200 in every phase: liveness must not fail during a migration or IBD. Co-Authored-By: Claude Opus 5.5 --- crates/rbitcoin-node/src/health.rs | 297 +++++++++++++++++++- crates/rbitcoin-node/src/run.rs | 33 ++- crates/rbitcoin-test/tests/cross_surface.rs | 65 +++++ 3 files changed, 388 insertions(+), 7 deletions(-) diff --git a/crates/rbitcoin-node/src/health.rs b/crates/rbitcoin-node/src/health.rs index af9bcd8ad..01163b8ee 100644 --- a/crates/rbitcoin-node/src/health.rs +++ b/crates/rbitcoin-node/src/health.rs @@ -1,15 +1,19 @@ -//! Health listener (`--health-listen`): unauthenticated `GET /healthz` for -//! process probes. +//! Health listener (`--health-listen`): unauthenticated `GET /healthz` and +//! `GET /readyz` for process probes. //! //! It binds at the top of [`crate::run_p2p`], before the store opens, so it //! answers through a schema migration, catch-up, and index materialize. The //! RPC, Electrum, and Esplora listeners bind only after those. +use axum::extract::State; use axum::http::StatusCode; use axum::routing::get; use axum::Router; use rbitcoin_log::{info, warn}; +use rbitcoin_net::{BlockingRegion, ChainHub}; use std::net::SocketAddr; +use std::sync::atomic::{AtomicU8, Ordering}; +use std::sync::{Arc, OnceLock}; use std::time::Duration; use tokio::net::TcpListener; use tokio::task::JoinHandle; @@ -21,6 +25,157 @@ use tower_http::timeout::TimeoutLayer; const MAX_CONCURRENT: usize = 16; /// Per-request wall. Probe timeouts are usually 1s. const REQUEST_TIMEOUT: Duration = Duration::from_secs(5); +/// Tip and scripthash index lag that `/readyz` still calls ready. +pub(crate) const READY_LAG_BLOCKS: u32 = 6; + +/// Where [`crate::run_p2p`] is in bring-up. `/readyz` is 503 until `Following`. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub(crate) enum Phase { + /// Store open: schema migration, backfill, spend replay. + Opening, + /// P2P, mempool, and proxy bring-up before catch-up. + Starting, + /// Initial block download. + CatchUp, + /// Tip-mode entry (scripthash materialize), follow peers, listeners. + Indexing, + /// Tip-follow loop. + Following, + /// Shutdown flush. + Stopping, +} + +impl Phase { + #[cfg(test)] + const ALL: [Self; 6] = [ + Self::Opening, + Self::Starting, + Self::CatchUp, + Self::Indexing, + Self::Following, + Self::Stopping, + ]; + + pub(crate) fn as_str(self) -> &'static str { + match self { + Self::Opening => "opening", + Self::Starting => "starting", + Self::CatchUp => "catch-up", + Self::Indexing => "indexing", + Self::Following => "following", + Self::Stopping => "stopping", + } + } + + fn from_u8(v: u8) -> Self { + match v { + 0 => Self::Opening, + 1 => Self::Starting, + 2 => Self::CatchUp, + 3 => Self::Indexing, + 4 => Self::Following, + _ => Self::Stopping, + } + } +} + +/// Bring-up state the health routes read. `run_p2p` writes the phase at each +/// transition; every read is an atomic load or a `OnceLock` get. +pub(crate) struct NodeStatus { + phase: AtomicU8, + sh_index: bool, + chain: OnceLock>, + /// Configured listeners that failed to bind (RPC, Electrum, Esplora only warn). + unbound: OnceLock>, +} + +impl NodeStatus { + pub(crate) fn new(sh_index: bool) -> Arc { + Arc::new(Self { + phase: AtomicU8::new(Phase::Opening as u8), + sh_index, + chain: OnceLock::new(), + unbound: OnceLock::new(), + }) + } + + pub(crate) fn enter(&self, phase: Phase) { + self.phase.store(phase as u8, Ordering::Release); + } + + pub(crate) fn attach_chain(&self, chain: &Arc) { + let _ = self.chain.set(Arc::clone(chain)); + } + + /// Enter [`Phase::Following`] with the configured listeners that did not bind. + pub(crate) fn follow(&self, unbound: Vec<&'static str>) { + let _ = self.unbound.set(unbound); + self.enter(Phase::Following); + } + + fn phase(&self) -> Phase { + Phase::from_u8(self.phase.load(Ordering::Acquire)) + } + + /// Chain reads may touch the store; call from the blocking pool. + fn ready_snapshot(&self) -> ReadySnapshot { + let phase = self.phase(); + let mut snap = ReadySnapshot { + phase, + unbound: self.unbound.get().cloned().unwrap_or_default(), + in_ibd: false, + blocks: 0, + headers: 0, + sh_lag: None, + }; + if phase != Phase::Following { + return snap; + } + if let Some(chain) = self.chain.get() { + snap.in_ibd = chain.in_ibd(); + snap.blocks = chain.query.tip_height().map_or(0, |h| h.0); + snap.headers = chain.best_header_height(); + snap.sh_lag = self.sh_index.then(|| chain.query.sh_lag_heights()); + } + snap + } +} + +/// The inputs of one `/readyz` answer. +#[derive(Clone, Debug, PartialEq, Eq)] +struct ReadySnapshot { + phase: Phase, + unbound: Vec<&'static str>, + /// Same value as RPC `getblockchaininfo.initialblockdownload`. + in_ibd: bool, + blocks: u32, + headers: u32, + /// `None` without `--sh-index`. + sh_lag: Option, +} + +/// `Ok` when ready; otherwise the first failing gate, in check order. +fn readiness(s: &ReadySnapshot) -> Result<(), String> { + if s.phase != Phase::Following { + return Err(s.phase.as_str().to_string()); + } + if !s.unbound.is_empty() { + return Err(format!("{} not listening", s.unbound.join(", "))); + } + if s.in_ibd { + return Err("initial block download".into()); + } + let behind = s.headers.saturating_sub(s.blocks); + if behind > READY_LAG_BLOCKS { + return Err(format!("tip {behind} blocks behind headers")); + } + match s.sh_lag { + Some(lag) if lag > READY_LAG_BLOCKS => { + Err(format!("scripthash index {lag} blocks behind tip")) + } + _ => Ok(()), + } +} /// Bound health listener. Dropping it stops serving. pub(crate) struct HealthHandle { @@ -34,13 +189,16 @@ impl Drop for HealthHandle { } /// Bind `addr` and serve the health routes until the handle drops. -pub(crate) async fn run_health(addr: SocketAddr) -> std::io::Result { +pub(crate) async fn run_health( + addr: SocketAddr, + status: Arc, +) -> std::io::Result { let listener = TcpListener::bind(addr).await?; let local_addr = listener.local_addr()?; if !local_addr.ip().is_loopback() { - warn!("health: {local_addr} is not loopback; /healthz is unauthenticated"); + warn!("health: {local_addr} is not loopback; /healthz and /readyz are unauthenticated"); } - let app = router(); + let app = router(status); let task = tokio::spawn(async move { if let Err(e) = axum::serve(listener, app).await { warn!("health: serve ended: {e}"); @@ -50,9 +208,10 @@ pub(crate) async fn run_health(addr: SocketAddr) -> std::io::Result Router { +fn router(status: Arc) -> Router { Router::new() .route("/healthz", get(healthz)) + .route("/readyz", get(readyz)) // Outer → inner: concurrency → body (GET only, so none) → timeout. .layer(TimeoutLayer::with_status_code( StatusCode::REQUEST_TIMEOUT, @@ -60,8 +219,134 @@ fn router() -> Router { )) .layer(RequestBodyLimitLayer::new(0)) .layer(ConcurrencyLimitLayer::new(MAX_CONCURRENT)) + .with_state(status) } async fn healthz() -> &'static str { "ok\n" } + +async fn readyz(State(status): State>) -> (StatusCode, String) { + let snap = tokio::task::spawn_blocking(move || { + let _g = BlockingRegion::enter(); + status.ready_snapshot() + }) + .await; + match snap.as_ref().map(readiness) { + Ok(Ok(())) => (StatusCode::OK, "ok\n".into()), + Ok(Err(reason)) => ( + StatusCode::SERVICE_UNAVAILABLE, + format!("not ready: {reason}\n"), + ), + Err(e) => (StatusCode::SERVICE_UNAVAILABLE, format!("not ready: {e}\n")), + } +} + +#[cfg(test)] +mod tests { + use super::*; + + fn following() -> ReadySnapshot { + ReadySnapshot { + phase: Phase::Following, + unbound: Vec::new(), + in_ibd: false, + blocks: 100, + headers: 100, + sh_lag: None, + } + } + + #[test] + fn readiness_reports_the_first_failing_gate() { + let mut cases = Vec::new(); + for phase in Phase::ALL { + let want = match phase { + Phase::Following => Ok(()), + other => Err(other.as_str().to_string()), + }; + cases.push(( + ReadySnapshot { + phase, + in_ibd: phase != Phase::Following, + ..following() + }, + want, + )); + } + let unbound = |names: &[&'static str]| ReadySnapshot { + unbound: names.to_vec(), + in_ibd: true, + ..following() + }; + cases.push((unbound(&["rpc"]), Err("rpc not listening".into()))); + cases.push(( + unbound(&["electrum", "esplora"]), + Err("electrum, esplora not listening".into()), + )); + cases.push(( + ReadySnapshot { + in_ibd: true, + headers: 200, + ..following() + }, + Err("initial block download".into()), + )); + let behind = |n: u32, sh_lag: Option| ReadySnapshot { + headers: 100 + n, + sh_lag, + ..following() + }; + cases.push((behind(READY_LAG_BLOCKS, None), Ok(()))); + cases.push(( + behind(READY_LAG_BLOCKS + 1, Some(0)), + Err(format!( + "tip {} blocks behind headers", + READY_LAG_BLOCKS + 1 + )), + )); + cases.push((behind(0, Some(READY_LAG_BLOCKS)), Ok(()))); + cases.push(( + behind(0, Some(READY_LAG_BLOCKS + 1)), + Err(format!( + "scripthash index {} blocks behind tip", + READY_LAG_BLOCKS + 1 + )), + )); + cases.push(( + ReadySnapshot { + headers: 90, + ..following() + }, + Ok(()), + )); + for (snap, want) in cases { + assert_eq!(readiness(&snap), want, "{snap:?}"); + } + } + + #[tokio::test(flavor = "multi_thread")] + async fn opening_node_is_live_but_not_ready() { + let addr = { + let l = std::net::TcpListener::bind("127.0.0.1:0").unwrap(); + l.local_addr().unwrap() + }; + let _h = run_health(addr, NodeStatus::new(false)).await.unwrap(); + assert_eq!(get(addr, "/healthz").await, "HTTP/1.1 200 OK|ok\n"); + assert_eq!( + get(addr, "/readyz").await, + "HTTP/1.1 503 Service Unavailable|not ready: opening\n" + ); + } + + async fn get(addr: SocketAddr, path: &str) -> String { + use tokio::io::{AsyncReadExt, AsyncWriteExt}; + let mut s = tokio::net::TcpStream::connect(addr).await.unwrap(); + let req = format!("GET {path} HTTP/1.1\r\nHost: x\r\nConnection: close\r\n\r\n"); + s.write_all(req.as_bytes()).await.unwrap(); + let mut buf = String::new(); + s.read_to_string(&mut buf).await.unwrap(); + let (head, body) = buf.split_once("\r\n\r\n").unwrap(); + format!("{}|{body}", head.lines().next().unwrap()) + } +} diff --git a/crates/rbitcoin-node/src/run.rs b/crates/rbitcoin-node/src/run.rs index f5b9d6df9..017b3c594 100644 --- a/crates/rbitcoin-node/src/run.rs +++ b/crates/rbitcoin-node/src/run.rs @@ -1,5 +1,6 @@ use crate::config::{parse_btc_to_sat, NodeConfig}; use crate::error::NodeError; +use crate::health::{run_health, NodeStatus, Phase}; use crate::regtest_rpc::HubRegtest; use bitcoin::consensus::Encodable; use rbitcoin_electrum::{run_electrum, ElectrumConfig, ElectrumHandle, TipNotify}; @@ -217,15 +218,17 @@ pub fn run_node(config: NodeConfig) -> Result { /// hold the process). #[allow(clippy::cognitive_complexity)] // node bring-up / P2P follow loop pub async fn run_p2p(config: NodeConfig) -> Result<(), NodeError> { + let status = NodeStatus::new(config.shindex); let _health = match config.listen.health { Some(addr) => Some( - crate::health::run_health(addr) + run_health(addr, Arc::clone(&status)) .await .map_err(|e| NodeError::Config(format!("health listen {addr}: {e}")))?, ), None => None, }; let handle = run_node(config.clone())?; + status.enter(Phase::Starting); let params = config.chain_params()?; let milestone = config.milestone(); if let Some(anchor) = milestone.anchor { @@ -290,6 +293,7 @@ pub async fn run_p2p(config: NodeConfig) -> Result<(), NodeError> { } } .map_err(|e| NodeError::Config(format!("p2p start: {e}")))?; + status.attach_chain(&node.hub); for extra in &config.listen.p2p_extra { let bound = node .add_listen(*extra) @@ -588,6 +592,7 @@ pub async fn run_p2p(config: NodeConfig) -> Result<(), NodeError> { pinned.extend(dns_resolved); let targets = follow_dial_targets(&pinned, &addrman, max_out, &occupied); let ibd_targets = follow_dial_targets(&pinned, &addrman, candidate_n, &occupied); + status.enter(Phase::CatchUp); let catch_up = run_ibd_or_skip( &node, &ibd_targets, @@ -617,6 +622,7 @@ pub async fn run_p2p(config: NodeConfig) -> Result<(), NodeError> { let mut sh_tip_ready = false; let mut index_writebehind = None; if catch_up.is_complete() && !shutdown.requested() { + status.enter(Phase::Indexing); let gates = enter_tip_mode( &node.hub.query, Some(Arc::clone(&shutdown.flag)), @@ -919,6 +925,30 @@ pub async fn run_p2p(config: NodeConfig) -> Result<(), NodeError> { } if tip_follow_ready && config.max_run_secs != Some(0) && !shutdown.requested() { + let unbound = [ + ( + "rpc", + config.rpc.socket || config.rpc.listen.is_some(), + rpc_handle.is_some(), + ), + ( + "electrum", + config.listen.electrum.is_some(), + !electrum_handles.is_empty(), + ), + ( + "esplora", + config.listen.esplora.is_some(), + !esplora_handles.is_empty(), + ), + ]; + status.follow( + unbound + .into_iter() + .filter(|&(_, configured, bound)| configured && !bound) + .map(|(name, ..)| name) + .collect(), + ); let deadline = config .max_run_secs .map(|s| Instant::now() + Duration::from_secs(s)); @@ -1172,6 +1202,7 @@ pub async fn run_p2p(config: NodeConfig) -> Result<(), NodeError> { } } + status.enter(Phase::Stopping); { let end_tip = node.tip_height().unwrap_or(0); let blocks_this_run = end_tip.saturating_sub(start_tip); diff --git a/crates/rbitcoin-test/tests/cross_surface.rs b/crates/rbitcoin-test/tests/cross_surface.rs index fefeff2fd..92d84cd7d 100644 --- a/crates/rbitcoin-test/tests/cross_surface.rs +++ b/crates/rbitcoin-test/tests/cross_surface.rs @@ -167,6 +167,21 @@ async fn pin_healthz(health_addr: SocketAddr) { assert_eq!(st, 404, "unknown health path: {body}"); } +/// `/readyz` once `run_p2p` is past bring-up (opening … indexing). +async fn readyz_after_startup(health_addr: SocketAddr) -> (u16, String) { + let deadline = Instant::now() + Duration::from_secs(20); + loop { + let (st, body) = http_get(health_addr, "/readyz").await; + let starting = ["opening", "starting", "catch-up", "indexing"] + .iter() + .any(|p| body == format!("not ready: {p}")); + if !starting || Instant::now() >= deadline { + return (st, body); + } + tokio::time::sleep(Duration::from_millis(50)).await; + } +} + async fn pin_address_prefix_404(esplora_addr: SocketAddr) { let (st, body) = http_get(esplora_addr, "/address-prefix/bc1").await; assert_eq!(st, 404, "address-prefix stays 404: {body}"); @@ -1042,6 +1057,46 @@ async fn fee_history_backfills_from_the_chain_when_relay_starts() { } } +/// A configured listener that did not bind keeps `/readyz` at 503 while the +/// node follows the tip. RPC only warns on a bind failure. +#[tokio::test(flavor = "multi_thread")] +async fn readyz_names_a_listener_that_did_not_bind() { + let td = TestDatadir::new().unwrap(); + { + let q = Query::open_or_create_tiny(td.store_path()).unwrap(); + let genesis = bitcoin::blockdata::constants::genesis_block(bitcoin::Network::Regtest); + accept_and_connect_block( + &q, + &ChainParams::regtest(), + Height::GENESIS, + &genesis, + Milestone::NONE, + ) + .unwrap(); + q.flush().unwrap(); + } + let taken = std::net::TcpListener::bind("127.0.0.1:0").unwrap(); + let health_addr = ephemeral_addr(); + let mut cfg = NodeConfig::default() + .with_datadir(td.path()) + .with_network(Network::Regtest) + .with_tiny_heads() + .with_p2p_listen("127.0.0.1:0".parse().unwrap()); + cfg.listen.use_seeds = false; + cfg.listen.connect.clear(); + cfg.rpc.listen = Some(taken.local_addr().unwrap()); + cfg.listen.health = Some(health_addr); + cfg.max_run_secs = Some(5); + let node = tokio::spawn(run_p2p(cfg)); + wait_listeners(&[health_addr]).await; + assert_eq!( + readyz_after_startup(health_addr).await, + (503, "not ready: rpc not listening".into()) + ); + let stopped = tokio::time::timeout(Duration::from_secs(20), node).await; + assert!(matches!(stopped, Ok(Ok(Ok(())))), "{stopped:?}"); +} + /// The health listener binds before the store opens: a taken port stops the /// node with nothing written to the datadir. #[tokio::test(flavor = "multi_thread")] @@ -1137,6 +1192,11 @@ async fn esplora_broadcast_visible_in_rpc_and_electrum() { assert_eq!(count["result"], 106, "{count}"); let chain = jsonrpc(rpc_addr, "getblockchaininfo", json!([])).await; assert_eq!(chain["result"]["initialblockdownload"], true, "{chain}"); + assert_eq!( + readyz_after_startup(health_addr).await, + (503, "not ready: initial block download".into()), + "/readyz agrees with RPC initialblockdownload" + ); let mpinfo = jsonrpc(rpc_addr, "getmempoolinfo", json!([])).await; assert_eq!(mpinfo["result"]["relay_enabled"], false, "{mpinfo}"); let tips = jsonrpc(rpc_addr, "getchaintips", json!([])).await; @@ -1830,6 +1890,11 @@ async fn esplora_broadcast_visible_in_rpc_and_electrum() { } let chain = jsonrpc(rpc_addr, "getblockchaininfo", json!([])).await; assert_eq!(chain["result"]["initialblockdownload"], false, "{chain}"); + assert_eq!( + chain["result"]["headers"], chain["result"]["blocks"], + "{chain}" + ); + assert_eq!(http_get(health_addr, "/readyz").await, (200, "ok".into())); let relay_parent = acs_spend( relay_cb, From 10330e49ed00f49d1ac7b569dc059edd51db7e31 Mon Sep 17 00:00:00 2001 From: Benjamen Keroack Date: Mon, 28 Sep 2026 11:55:26 -0600 Subject: [PATCH 03/11] node: --metrics serves Prometheus gauges that equal their RPC fields --metrics (conf metrics, refused without --health-listen) adds GET /metrics on the health listener in the Prometheus text format. It does not grow a second metrics dialect: each gauge is a value the node already publishes, named after its source. - rbitcoin_blocks, _headers, _tip_time_seconds, _initial_block_download: getblockchaininfo blocks, headers, time, initialblockdownload - rbitcoin_connections{direction}: getnetworkinfo connections_in/_out, now counted by one rbitcoin_net::connection_counts both call - rbitcoin_mempool_transactions, _bytes: getmempoolinfo size, bytes - rbitcoin_scripthash_lag_blocks: tip: accept sh_lag= (with --sh-index) - process_resident_memory_bytes: ibd: sizes rss=; process_start_time_seconds - rbitcoin_build_info, rbitcoin_phase, rbitcoin_ready: /readyz itself The cross-surface journey scrapes once during IBD and once after and checks each gauge against RPC at that moment. A scrape takes one peer snapshot and folds the mempool once (the getmempoolinfo fold) on the blocking pool. Co-Authored-By: Claude Opus 5.5 --- crates/rbitcoin-net/src/lib.rs | 5 +- crates/rbitcoin-net/src/peers.rs | 7 + crates/rbitcoin-node/src/cli.rs | 8 +- crates/rbitcoin-node/src/config.rs | 10 ++ crates/rbitcoin-node/src/health.rs | 64 +++++++-- crates/rbitcoin-node/src/health/metrics.rs | 144 ++++++++++++++++++++ crates/rbitcoin-node/src/run.rs | 7 +- crates/rbitcoin-rpc/src/methods/net.rs | 3 +- crates/rbitcoin-test/tests/cross_surface.rs | 91 +++++++++++++ crates/rbitcoin-test/tests/scenarios.rs | 1 + 10 files changed, 318 insertions(+), 22 deletions(-) create mode 100644 crates/rbitcoin-node/src/health/metrics.rs diff --git a/crates/rbitcoin-net/src/lib.rs b/crates/rbitcoin-net/src/lib.rs index c4554fe42..9d40dd884 100644 --- a/crates/rbitcoin-net/src/lib.rs +++ b/crates/rbitcoin-net/src/lib.rs @@ -59,8 +59,9 @@ pub use peer::{ }; pub use peer_dos::DEFAULT_MAX_INBOUND; pub use peers::{ - parse_peer_addr, parse_peer_addr_with_port, parse_peer_net, pick_stale_follow_evict, - DialRequest, DialTarget, LivePeer, PeerConnType, PeerHub, PeerInfo, PeerOut, PingAction, + connection_counts, parse_peer_addr, parse_peer_addr_with_port, parse_peer_net, + pick_stale_follow_evict, DialRequest, DialTarget, LivePeer, PeerConnType, PeerHub, PeerInfo, + PeerOut, PingAction, }; pub(crate) use rbitcoin_mempool::MempoolGraphStats; pub use rbitcoin_mempool::{AcceptError, Selected}; diff --git a/crates/rbitcoin-net/src/peers.rs b/crates/rbitcoin-net/src/peers.rs index 7afc22d34..f36d3ede8 100644 --- a/crates/rbitcoin-net/src/peers.rs +++ b/crates/rbitcoin-net/src/peers.rs @@ -1181,6 +1181,13 @@ fn acct_bytes(cmd: &str, payload: u64) -> u64 { } } +/// `(inbound, outbound)` among `peers`: `getnetworkinfo.connections_in` / +/// `connections_out` over a [`PeerHub::snapshot`]. +pub fn connection_counts(peers: &[PeerInfo]) -> (u64, u64) { + let inbound = peers.iter().filter(|p| p.inbound).count(); + (inbound as u64, (peers.len() - inbound) as u64) +} + /// RPC-facing snapshot. #[derive(Clone, Debug)] pub struct PeerInfo { diff --git a/crates/rbitcoin-node/src/cli.rs b/crates/rbitcoin-node/src/cli.rs index b8ff7d088..0a7e534a1 100644 --- a/crates/rbitcoin-node/src/cli.rs +++ b/crates/rbitcoin-node/src/cli.rs @@ -302,7 +302,7 @@ fn operator_usage() -> String { [--listen ADDR] [--no-listen] [--connect ADDR]... [--seed-node HOST]... [--proxy HOST:PORT] [--onion HOST:PORT] [--proxy-randomize[=0|1]] [--only-net NET]... \\\n\ [--tor-control [HOST:PORT]] [--tor-control-cookie PATH] [--tor-control-password PASS] \\\n\ [--i2p-sam [HOST:PORT]] [--i2p-accept-incoming] \\\n\ - [--electrum-listen ADDR] [--esplora-listen ADDR] [--esplora-onion[=0|1]] [--health-listen [ADDR]] \\\n\ + [--electrum-listen ADDR] [--esplora-listen ADDR] [--esplora-onion[=0|1]] [--health-listen [ADDR]] [--metrics] \\\n\ [--sh-index] [--block-filter-index] [--prune-seqsigwit] [--prune-seqsigwit-ram-threshold-bytes N] [--sp-tweaks] [--sp-tweaks-dust SATS] [--max-sh-creates N] [--esplora-block-template] \\\n\ [--rpc] [--rpc-listen [ADDR]] [--rpc-socket PATH] [--rpc-token-file PATH] [--rpc-cookie-file PATH] [--rpc-work-queue N] \\\n\ [--milestone HEIGHT] \\\n\ @@ -354,7 +354,8 @@ Silent payments: --sp-tweaks (default off) writes/serves the thin BIP-352 tweak Not with --prune-seqsigwit (tweaks read scriptSig and witness).\n\ --sp-tweaks-dust SATS omits served P2TR outs with value <= SATS (default 1000; 0 = all; 546 = Cake electrs).\n\ Health: --health-listen [ADDR] serves GET /healthz from the first second of startup\n\ - (default 127.0.0.1:9332). Unauthenticated; keep it on loopback or a probe-only network.\n\ + and GET /readyz (default 127.0.0.1:9332). Unauthenticated; keep it on loopback or a\n\ + probe-only network. --metrics adds Prometheus GET /metrics there (needs --health-listen).\n\ RPC: --rpc unix socket {{datadir}}/rpc.sock; --rpc-listen [ADDR] adds TCP (default 127.0.0.1 and Core-matching port). Token {{datadir}}/rpc.token (Bearer); --rpc-cookie-file opts TCP into Core cookie HTTP Basic. No --rpcuser.\n\ Cold files: --datadir-cold PATH puts Class A seqsigwit.body/idx under PATH/store (HDD).\n\ Default (flag omitted): hot and cold files both live under --datadir.\n\ @@ -410,6 +411,7 @@ fn is_bool_key(key: &str) -> bool { | "proxy_randomize" | "i2p_accept_incoming" | "inhibit_suspend" + | "metrics" | "trusted" | "always_relay" | "relay" @@ -817,6 +819,7 @@ mod tests { "--electrum-listen", "--esplora-listen", "--health-listen", + "--metrics", ]); assert!(cfg.shindex && cfg.sptweaks); assert_eq!(cfg.sptweaks_dust, 546); @@ -829,6 +832,7 @@ mod tests { assert_eq!(cfg.rpc.listen, Some("127.0.0.1:18443".parse().unwrap())); assert_eq!(cfg.listen.electrum.unwrap().port(), 50001); assert_eq!(cfg.listen.health, Some("127.0.0.1:9332".parse().unwrap())); + assert!(cfg.metrics); match cfg.listen.esplora.unwrap() { rbitcoin_esplora::EsploraListen::Tcp(a) => assert_eq!(a.port(), 3000), #[cfg(unix)] diff --git a/crates/rbitcoin-node/src/config.rs b/crates/rbitcoin-node/src/config.rs index 25a0b97f5..84c45579e 100644 --- a/crates/rbitcoin-node/src/config.rs +++ b/crates/rbitcoin-node/src/config.rs @@ -312,6 +312,8 @@ pub struct NodeConfig { pub max_sh_creates: u32, /// Opt-in Esplora `GET /block-template` (GBT template JSON). Default off. pub esplora_block_template: bool, + /// Prometheus `GET /metrics` on the health listener. Default off. + pub metrics: bool, /// ADD_ONION for `--esplora-listen` when `--tor-control` is set. Default on. pub esplora_onion: bool, /// Skip script/prevout checks for blocks at or below this height (0 = off). @@ -389,6 +391,7 @@ impl Default for NodeConfig { sptweaks_dust: rbitcoin_electrum::DEFAULT_TWEAKS_MIN_DUST, max_sh_creates: rbitcoin_query::DEFAULT_MAX_SH_CREATES, esplora_block_template: false, + metrics: false, esplora_onion: true, milestone_height: 0, milestone_explicit: false, @@ -607,6 +610,9 @@ impl NodeConfig { .into(), )); } + if self.metrics && self.listen.health.is_none() { + return Err(NodeError::Config("--metrics needs --health-listen".into())); + } self.validate_only_net()?; self.validate_hidden_inbound()?; self.validate_rpc_cookie() @@ -1300,6 +1306,10 @@ impl NodeConfig { .map_err(|e| NodeError::Config(format!("conf max_run_secs: {e}")))?, ); } + "metrics" => { + self.metrics = parse_conf_bool(val) + .map_err(|e| NodeError::Config(format!("conf metrics: {e}")))?; + } "inhibit_suspend" => { self.inhibit_suspend = parse_conf_bool(val) .map_err(|e| NodeError::Config(format!("conf inhibit_suspend: {e}")))?; diff --git a/crates/rbitcoin-node/src/health.rs b/crates/rbitcoin-node/src/health.rs index 01163b8ee..9292224a7 100644 --- a/crates/rbitcoin-node/src/health.rs +++ b/crates/rbitcoin-node/src/health.rs @@ -1,20 +1,24 @@ //! Health listener (`--health-listen`): unauthenticated `GET /healthz` and -//! `GET /readyz` for process probes. +//! `GET /readyz` for process probes, and `GET /metrics` with `--metrics`. //! //! It binds at the top of [`crate::run_p2p`], before the store opens, so it //! answers through a schema migration, catch-up, and index materialize. The //! RPC, Electrum, and Esplora listeners bind only after those. +mod metrics; + use axum::extract::State; -use axum::http::StatusCode; +use axum::http::{header, StatusCode}; +use axum::response::{IntoResponse, Response}; use axum::routing::get; use axum::Router; use rbitcoin_log::{info, warn}; -use rbitcoin_net::{BlockingRegion, ChainHub}; +use rbitcoin_net::{BlockingRegion, ChainHub, MempoolHub, PeerHub}; +use rbitcoin_primitives::Network; use std::net::SocketAddr; use std::sync::atomic::{AtomicU8, Ordering}; use std::sync::{Arc, OnceLock}; -use std::time::Duration; +use std::time::{Duration, SystemTime}; use tokio::net::TcpListener; use tokio::task::JoinHandle; use tower::limit::ConcurrencyLimitLayer; @@ -46,7 +50,6 @@ pub(crate) enum Phase { } impl Phase { - #[cfg(test)] const ALL: [Self; 6] = [ Self::Opening, Self::Starting, @@ -83,18 +86,26 @@ impl Phase { /// transition; every read is an atomic load or a `OnceLock` get. pub(crate) struct NodeStatus { phase: AtomicU8, + network: Network, sh_index: bool, + started: SystemTime, chain: OnceLock>, + peers: OnceLock>, + mempool: OnceLock>, /// Configured listeners that failed to bind (RPC, Electrum, Esplora only warn). unbound: OnceLock>, } impl NodeStatus { - pub(crate) fn new(sh_index: bool) -> Arc { + pub(crate) fn new(network: Network, sh_index: bool) -> Arc { Arc::new(Self { phase: AtomicU8::new(Phase::Opening as u8), + network, sh_index, + started: SystemTime::now(), chain: OnceLock::new(), + peers: OnceLock::new(), + mempool: OnceLock::new(), unbound: OnceLock::new(), }) } @@ -103,8 +114,13 @@ impl NodeStatus { self.phase.store(phase as u8, Ordering::Release); } - pub(crate) fn attach_chain(&self, chain: &Arc) { + pub(crate) fn attach_p2p(&self, chain: &Arc, peers: &Arc) { let _ = self.chain.set(Arc::clone(chain)); + let _ = self.peers.set(Arc::clone(peers)); + } + + pub(crate) fn attach_mempool(&self, mempool: &Arc) { + let _ = self.mempool.set(Arc::clone(mempool)); } /// Enter [`Phase::Following`] with the configured listeners that did not bind. @@ -188,17 +204,19 @@ impl Drop for HealthHandle { } } -/// Bind `addr` and serve the health routes until the handle drops. +/// Bind `addr` and serve the health routes (and `/metrics` when `metrics`) +/// until the handle drops. pub(crate) async fn run_health( addr: SocketAddr, status: Arc, + metrics: bool, ) -> std::io::Result { let listener = TcpListener::bind(addr).await?; let local_addr = listener.local_addr()?; if !local_addr.ip().is_loopback() { warn!("health: {local_addr} is not loopback; /healthz and /readyz are unauthenticated"); } - let app = router(status); + let app = router(status, metrics); let task = tokio::spawn(async move { if let Err(e) = axum::serve(listener, app).await { warn!("health: serve ended: {e}"); @@ -208,10 +226,16 @@ pub(crate) async fn run_health( Ok(HealthHandle { task }) } -fn router(status: Arc) -> Router { - Router::new() +fn router(status: Arc, metrics: bool) -> Router { + let routes = Router::new() .route("/healthz", get(healthz)) - .route("/readyz", get(readyz)) + .route("/readyz", get(readyz)); + let routes = if metrics { + routes.route("/metrics", get(scrape)) + } else { + routes + }; + routes // Outer → inner: concurrency → body (GET only, so none) → timeout. .layer(TimeoutLayer::with_status_code( StatusCode::REQUEST_TIMEOUT, @@ -226,6 +250,18 @@ async fn healthz() -> &'static str { "ok\n" } +async fn scrape(State(status): State>) -> Response { + let body = tokio::task::spawn_blocking(move || { + let _g = BlockingRegion::enter(); + metrics::render(&status) + }) + .await; + match body { + Ok(body) => ([(header::CONTENT_TYPE, metrics::CONTENT_TYPE)], body).into_response(), + Err(e) => (StatusCode::INTERNAL_SERVER_ERROR, format!("metrics: {e}\n")).into_response(), + } +} + async fn readyz(State(status): State>) -> (StatusCode, String) { let snap = tokio::task::spawn_blocking(move || { let _g = BlockingRegion::enter(); @@ -331,7 +367,9 @@ mod tests { let l = std::net::TcpListener::bind("127.0.0.1:0").unwrap(); l.local_addr().unwrap() }; - let _h = run_health(addr, NodeStatus::new(false)).await.unwrap(); + let _h = run_health(addr, NodeStatus::new(Network::Regtest, false), false) + .await + .unwrap(); assert_eq!(get(addr, "/healthz").await, "HTTP/1.1 200 OK|ok\n"); assert_eq!( get(addr, "/readyz").await, diff --git a/crates/rbitcoin-node/src/health/metrics.rs b/crates/rbitcoin-node/src/health/metrics.rs new file mode 100644 index 000000000..67f008313 --- /dev/null +++ b/crates/rbitcoin-node/src/health/metrics.rs @@ -0,0 +1,144 @@ +//! Prometheus text exposition for `--metrics`. Each gauge is a value RPC or a +//! log line already publishes, named after that source; `phase` and `ready` +//! are `/readyz` itself. + +use super::{readiness, NodeStatus, Phase}; +use rbitcoin_net::ChainHub; +use std::fmt::{Display, Write}; +use std::time::UNIX_EPOCH; + +pub(super) const CONTENT_TYPE: &str = "text/plain; version=0.0.4; charset=utf-8"; + +/// One scrape. Chain reads may touch the store and the mempool totals take +/// its lock, so call from the blocking pool. Cost: one peer snapshot and one +/// fold over the mempool (the same fold as `getmempoolinfo`). +pub(super) fn render(status: &NodeStatus) -> String { + let mut out = Exposition::default(); + out.family( + "rbitcoin_build_info", + "gauge", + "Version and network of this node.", + ); + out.sample( + "rbitcoin_build_info", + &format!( + "{{version=\"{}\",network=\"{}\"}}", + env!("CARGO_PKG_VERSION"), + status.network.as_str() + ), + 1, + ); + let phase = status.phase(); + out.family( + "rbitcoin_phase", + "gauge", + "Bring-up phase, as /readyz names it.", + ); + for p in Phase::ALL { + out.sample( + "rbitcoin_phase", + &format!("{{phase=\"{}\"}}", p.as_str()), + u8::from(p == phase), + ); + } + out.gauge( + "rbitcoin_ready", + "1 when /readyz answers 200.", + u8::from(readiness(&status.ready_snapshot()).is_ok()), + ); + out.gauge( + "process_start_time_seconds", + "Start time of the process since the Unix epoch in seconds.", + status + .started + .duration_since(UNIX_EPOCH) + .map_or(0, |d| d.as_secs()), + ); + let rss_kb = rbitcoin_net::read_platform_rss().rss_kb; + if rss_kb > 0 { + out.gauge( + "process_resident_memory_bytes", + "Resident memory size in bytes (ibd: sizes rss=).", + rss_kb * 1024, + ); + } + if let Some(chain) = status.chain.get() { + chain_gauges(&mut out, chain, status.sh_index); + } + if let Some(peers) = status.peers.get() { + let (inbound, outbound) = rbitcoin_net::connection_counts(&peers.snapshot()); + out.family( + "rbitcoin_connections", + "gauge", + "Peer connections (getnetworkinfo.connections_in / connections_out).", + ); + out.sample("rbitcoin_connections", "{direction=\"in\"}", inbound); + out.sample("rbitcoin_connections", "{direction=\"out\"}", outbound); + } + if let Some(mempool) = status.mempool.get() { + let (size, vbytes, _fee) = mempool.live_adjusted_totals(); + out.gauge( + "rbitcoin_mempool_transactions", + "Mempool transactions (getmempoolinfo.size).", + size, + ); + out.gauge( + "rbitcoin_mempool_bytes", + "Sum of mempool virtual sizes (getmempoolinfo.bytes).", + vbytes, + ); + } + out.0 +} + +fn chain_gauges(out: &mut Exposition, chain: &ChainHub, sh_index: bool) { + let tip = chain.query.tip_height(); + out.gauge( + "rbitcoin_blocks", + "Active chain height (getblockchaininfo.blocks).", + tip.map_or(0, |h| h.0), + ); + out.gauge( + "rbitcoin_headers", + "Best header height (getblockchaininfo.headers).", + chain.best_header_height(), + ); + let time = tip + .and_then(|h| chain.query.header_at_height(h).ok().flatten()) + .map_or(0, |(_, rec)| rec.timestamp); + out.gauge( + "rbitcoin_tip_time_seconds", + "Tip block time (getblockchaininfo.time).", + time, + ); + out.gauge( + "rbitcoin_initial_block_download", + "1 during initial block download (getblockchaininfo.initialblockdownload).", + u8::from(chain.in_ibd()), + ); + if sh_index { + out.gauge( + "rbitcoin_scripthash_lag_blocks", + "Blocks the scripthash index trails the tip (tip: accept sh_lag=).", + chain.query.sh_lag_heights(), + ); + } +} + +#[derive(Default)] +struct Exposition(String); + +impl Exposition { + fn family(&mut self, name: &str, kind: &str, help: &str) { + let _ = writeln!(self.0, "# HELP {name} {help}\n# TYPE {name} {kind}"); + } + + fn sample(&mut self, name: &str, labels: &str, value: impl Display) { + let _ = writeln!(self.0, "{name}{labels} {value}"); + } + + fn gauge(&mut self, name: &str, help: &str, value: impl Display) { + self.family(name, "gauge", help); + self.sample(name, "", value); + } +} diff --git a/crates/rbitcoin-node/src/run.rs b/crates/rbitcoin-node/src/run.rs index 017b3c594..191e31bd5 100644 --- a/crates/rbitcoin-node/src/run.rs +++ b/crates/rbitcoin-node/src/run.rs @@ -218,10 +218,10 @@ pub fn run_node(config: NodeConfig) -> Result { /// hold the process). #[allow(clippy::cognitive_complexity)] // node bring-up / P2P follow loop pub async fn run_p2p(config: NodeConfig) -> Result<(), NodeError> { - let status = NodeStatus::new(config.shindex); + let status = NodeStatus::new(config.network, config.shindex); let _health = match config.listen.health { Some(addr) => Some( - run_health(addr, Arc::clone(&status)) + run_health(addr, Arc::clone(&status), config.metrics) .await .map_err(|e| NodeError::Config(format!("health listen {addr}: {e}")))?, ), @@ -293,7 +293,7 @@ pub async fn run_p2p(config: NodeConfig) -> Result<(), NodeError> { } } .map_err(|e| NodeError::Config(format!("p2p start: {e}")))?; - status.attach_chain(&node.hub); + status.attach_p2p(&node.hub, &node.peers); for extra in &config.listen.p2p_extra { let bound = node .add_listen(*extra) @@ -386,6 +386,7 @@ pub async fn run_p2p(config: NodeConfig) -> Result<(), NodeError> { .map_err(|e| NodeError::Config(format!("mempool open join: {e}")))? .map_err(NodeError::Config)?; node.peers.attach_mempool(&mempool); + status.attach_mempool(&mempool); if config.listen.proxy.is_some() || config.listen.onion.is_some() { mempool.set_isolated_broadcast(true); } diff --git a/crates/rbitcoin-rpc/src/methods/net.rs b/crates/rbitcoin-rpc/src/methods/net.rs index 2b478e83a..8560a1bfb 100644 --- a/crates/rbitcoin-rpc/src/methods/net.rs +++ b/crates/rbitcoin-rpc/src/methods/net.rs @@ -342,8 +342,7 @@ pub(crate) fn localaddresses_json(ctx: &RpcContext) -> Value { pub(crate) fn getnetworkinfo(ctx: &RpcContext) -> Value { let (cin, cout, timeoffset) = if let Some(hub) = ctx.peers.as_ref() { let rows = hub.snapshot(); - let cin = rows.iter().filter(|p| p.inbound).count() as u64; - let cout = rows.iter().filter(|p| !p.inbound).count() as u64; + let (cin, cout) = rbitcoin_net::connection_counts(&rows); (cin, cout, outbound_median_time_offset(&rows)) } else { (0, ctx.connections.load(Ordering::Relaxed), 0) diff --git a/crates/rbitcoin-test/tests/cross_surface.rs b/crates/rbitcoin-test/tests/cross_surface.rs index 92d84cd7d..7f41de1e1 100644 --- a/crates/rbitcoin-test/tests/cross_surface.rs +++ b/crates/rbitcoin-test/tests/cross_surface.rs @@ -13,6 +13,7 @@ use rbitcoin_primitives::{Height, Network}; use rbitcoin_query::Query; use rbitcoin_test::{build_mature_regtest_with_spend, TestDatadir}; use serde_json::{json, Value}; +use std::collections::HashMap; use std::net::SocketAddr; use std::str::FromStr; use std::sync::atomic::Ordering; @@ -182,6 +183,87 @@ async fn readyz_after_startup(health_addr: SocketAddr) -> (u16, String) { } } +/// `GET /metrics` in the Prometheus text format: one value per series. +async fn scrape_metrics(health_addr: SocketAddr) -> HashMap { + let mut stream = TcpStream::connect(health_addr) + .await + .expect("metrics connect"); + stream + .write_all(b"GET /metrics HTTP/1.1\r\nHost: 127.0.0.1\r\nConnection: close\r\n\r\n") + .await + .unwrap(); + let mut text = String::new(); + stream.read_to_string(&mut text).await.unwrap(); + let (head, body) = text.split_once("\r\n\r\n").expect("http headers"); + assert!(head.starts_with("HTTP/1.1 200"), "{head}"); + assert!( + head.to_ascii_lowercase() + .contains("content-type: text/plain; version=0.0.4; charset=utf-8"), + "{head}" + ); + body.lines() + .filter(|l| !l.is_empty() && !l.starts_with('#')) + .map(|l| { + let (series, v) = l.rsplit_once(' ').unwrap_or_else(|| panic!("sample: {l}")); + let v = v.parse().unwrap_or_else(|_| panic!("sample value: {l}")); + (series.to_string(), v) + }) + .collect() +} + +/// Every `/metrics` gauge that names an RPC field equals that field. +async fn pin_metrics_equal_rpc( + health_addr: SocketAddr, + rpc_addr: SocketAddr, + ready: bool, +) -> HashMap { + let chain = jsonrpc(rpc_addr, "getblockchaininfo", json!([])).await["result"].clone(); + let net = jsonrpc(rpc_addr, "getnetworkinfo", json!([])).await["result"].clone(); + let mempool = jsonrpc(rpc_addr, "getmempoolinfo", json!([])).await["result"].clone(); + let m = scrape_metrics(health_addr).await; + let num = |v: &Value| v.as_f64().unwrap_or_else(|| panic!("number: {v}")); + let flag = |b: bool| if b { 1.0 } else { 0.0 }; + let build = format!( + "rbitcoin_build_info{{version=\"{}\",network=\"regtest\"}}", + env!("CARGO_PKG_VERSION") + ); + for (series, want) in [ + (build.as_str(), 1.0), + ("rbitcoin_phase{phase=\"following\"}", 1.0), + ("rbitcoin_phase{phase=\"opening\"}", 0.0), + ("rbitcoin_ready", flag(ready)), + ("rbitcoin_blocks", num(&chain["blocks"])), + ("rbitcoin_headers", num(&chain["headers"])), + ("rbitcoin_tip_time_seconds", num(&chain["time"])), + ( + "rbitcoin_initial_block_download", + flag(chain["initialblockdownload"] == true), + ), + ( + "rbitcoin_connections{direction=\"in\"}", + num(&net["connections_in"]), + ), + ( + "rbitcoin_connections{direction=\"out\"}", + num(&net["connections_out"]), + ), + ("rbitcoin_mempool_transactions", num(&mempool["size"])), + ("rbitcoin_mempool_bytes", num(&mempool["bytes"])), + ] { + assert_eq!(m.get(series), Some(&want), "{series}: {m:?}"); + } + assert!(m["rbitcoin_scripthash_lag_blocks"] <= 6.0, "{m:?}"); + let now = std::time::SystemTime::now() + .duration_since(std::time::UNIX_EPOCH) + .unwrap() + .as_secs_f64(); + let started = m["process_start_time_seconds"]; + assert!(started <= now && started > now - 600.0, "{m:?}"); + #[cfg(any(target_os = "linux", target_os = "macos"))] + assert!(m["process_resident_memory_bytes"] > 0.0, "{m:?}"); + m +} + async fn pin_address_prefix_404(esplora_addr: SocketAddr) { let (st, body) = http_get(esplora_addr, "/address-prefix/bc1").await; assert_eq!(st, 404, "address-prefix stays 404: {body}"); @@ -1093,6 +1175,8 @@ async fn readyz_names_a_listener_that_did_not_bind() { readyz_after_startup(health_addr).await, (503, "not ready: rpc not listening".into()) ); + let (st, body) = http_get(health_addr, "/metrics").await; + assert_eq!(st, 404, "/metrics without --metrics: {body}"); let stopped = tokio::time::timeout(Duration::from_secs(20), node).await; assert!(matches!(stopped, Ok(Ok(Ok(())))), "{stopped:?}"); } @@ -1157,6 +1241,7 @@ async fn esplora_broadcast_visible_in_rpc_and_electrum() { cfg.listen.esplora = Some(rbitcoin_esplora::EsploraListen::Tcp(esplora_addr)); cfg.rpc.listen = Some(rpc_addr); cfg.listen.health = Some(health_addr); + cfg.metrics = true; // mempool's CORE_RPC.SOCKET_PATH reaches the node from another user. let rpc_sock = td.path().join("run").join("rpc.sock"); cfg.apply_kv("rpc_socket", rpc_sock.to_str().unwrap()) @@ -1759,6 +1844,11 @@ async fn esplora_broadcast_visible_in_rpc_and_electrum() { let (st, body) = http_get(esplora_addr, "/fee-estimates").await; assert_eq!(st, 503, "GET /fee-estimates: {body}"); + let metrics_in_ibd = pin_metrics_equal_rpc(health_addr, rpc_addr, false).await; + assert!( + metrics_in_ibd["rbitcoin_mempool_transactions"] > 0.0, + "{metrics_in_ibd:?}" + ); let tip_before = jsonrpc(rpc_addr, "getbestblockhash", json!([])).await; let tip_hash = tip_before["result"].as_str().expect("tip hash").to_string(); pin_waitforblockheight_timeout_zero_behind(rpc_addr, 106, &tip_hash).await; @@ -1895,6 +1985,7 @@ async fn esplora_broadcast_visible_in_rpc_and_electrum() { "{chain}" ); assert_eq!(http_get(health_addr, "/readyz").await, (200, "ok".into())); + pin_metrics_equal_rpc(health_addr, rpc_addr, true).await; let relay_parent = acs_spend( relay_cb, diff --git a/crates/rbitcoin-test/tests/scenarios.rs b/crates/rbitcoin-test/tests/scenarios.rs index c6e8a2834..7c2faebb8 100644 --- a/crates/rbitcoin-test/tests/scenarios.rs +++ b/crates/rbitcoin-test/tests/scenarios.rs @@ -235,6 +235,7 @@ fn pin_validate_refusals(td: &TestDatadir) { ("chainwork-not-hex", &["--min-chain-work=test"][..]), ("challenge-off-signet", &["--signet-challenge", "51"]), ("tweaks-and-pruning", &["--sp-tweaks", "--prune-seqsigwit"]), + ("metrics-without-health", &["--metrics"]), ] { assert!(exit_is(smoke(td, name, args), 1), "{name} must refuse"); } From c177a2cff429f3d1b5e19dc424d81fc6b940c96a Mon Sep 17 00:00:00 2001 From: Benjamen Keroack Date: Mon, 28 Sep 2026 12:30:09 -0600 Subject: [PATCH 04/11] net: tip: perf meters are running totals; /metrics exports them The Esplora, Electrum, historical block serve, and mempool accept meters were sample-and-reset: the 5 s perf tick swapped each atomic to zero, so no reader could see a running count and Prometheus had nothing to scrape. PerfCounter keeps the running total and a window mark. The perf tick takes the change since its previous sample (the DEBUG line prints the same numbers); /metrics reads the total: - rbitcoin_esplora_requests_total, rbitcoin_esplora_request_seconds_total - rbitcoin_electrum_requests_total, rbitcoin_electrum_request_seconds_total - rbitcoin_block_serve_total, rbitcoin_block_serve_bytes_total - rbitcoin_mempool_accepts_total, rbitcoin_mempool_rejects_total A window maximum cannot come from totals, so PerfMax stays a reset and is not exported. Esplora and Electrum share one RequestMeter in place of two copies of the same statics and compare-exchange loop. No extra work per request; the window mark is touched only by the tick. Co-Authored-By: Claude Opus 5.5 --- crates/rbitcoin-electrum/src/lib.rs | 4 +- crates/rbitcoin-electrum/src/server.rs | 35 ++---- crates/rbitcoin-esplora/src/lib.rs | 3 +- crates/rbitcoin-esplora/src/server.rs | 36 +++--- crates/rbitcoin-net/src/lib.rs | 6 +- crates/rbitcoin-net/src/perf_meter.rs | 128 ++++++++++++++++++++ crates/rbitcoin-net/src/serve_perf.rs | 46 +++---- crates/rbitcoin-net/src/tx_relay.rs | 45 ++++--- crates/rbitcoin-node/src/health/metrics.rs | 53 ++++++++ crates/rbitcoin-test/tests/cross_surface.rs | 24 +++- 10 files changed, 284 insertions(+), 96 deletions(-) create mode 100644 crates/rbitcoin-net/src/perf_meter.rs diff --git a/crates/rbitcoin-electrum/src/lib.rs b/crates/rbitcoin-electrum/src/lib.rs index 83310a6eb..f66edb869 100644 --- a/crates/rbitcoin-electrum/src/lib.rs +++ b/crates/rbitcoin-electrum/src/lib.rs @@ -14,8 +14,8 @@ mod tweaks; mod unspent; pub use server::{ - electrum_scripthash_hex, parse_electrum_request_line, run_electrum, sample_reset_perf, - ElectrumConfig, ElectrumHandle, ServeLimits, TipNotify, + electrum_scripthash_hex, parse_electrum_request_line, perf_totals, run_electrum, + sample_reset_perf, ElectrumConfig, ElectrumHandle, ServeLimits, TipNotify, }; pub use tweaks::DEFAULT_TWEAKS_MIN_DUST; pub use unspent::{ diff --git a/crates/rbitcoin-electrum/src/server.rs b/crates/rbitcoin-electrum/src/server.rs index 9e4daa013..9c5483436 100644 --- a/crates/rbitcoin-electrum/src/server.rs +++ b/crates/rbitcoin-electrum/src/server.rs @@ -6,6 +6,7 @@ use bitcoin::consensus::Encodable; use bitcoin::hashes::Hash; use rbitcoin_consensus::ChainParams; +use rbitcoin_net::RequestMeter; use rbitcoin_net::{BlockingRegion, MempoolHub}; use rbitcoin_primitives::{Fk, Height}; use rbitcoin_query::{ChainView, ChainViewKind, HistoryFilter, Query, ShJoinSlot}; @@ -13,7 +14,7 @@ use rbitcoin_store::{script_hash, StoreError}; use serde_json::{json, Value}; use std::collections::{HashMap, HashSet}; use std::net::SocketAddr; -use std::sync::atomic::{AtomicBool, AtomicU64, Ordering}; +use std::sync::atomic::{AtomicBool, Ordering}; use std::sync::{Arc, Mutex, OnceLock}; use std::time::{Duration, Instant}; use tokio::io::{AsyncBufReadExt, AsyncRead, AsyncWrite, AsyncWriteExt, BufReader}; @@ -35,32 +36,22 @@ pub fn parse_electrum_request_line(line: &str) -> Option { serde_json::from_str(line).ok() } -/// Tip-follow 5s DEBUG `tip: perf`: JSON-RPC request count this window. -static METER_REQ: AtomicU64 = AtomicU64::new(0); -/// Sum of dispatch walls (µs). -static METER_US: AtomicU64 = AtomicU64::new(0); -/// Max single dispatch wall (µs). -static METER_MAX_US: AtomicU64 = AtomicU64::new(0); +/// Electrum JSON-RPC requests: the `tip: perf` window and the `/metrics` totals. +static METER: RequestMeter = RequestMeter::new(); -/// Sample-and-reset Electrum request meters: `(count, sum_us, max_us)`. +/// Electrum request meters since the previous sample: `(count, sum_us, max_us)`. +/// Running totals are untouched ([`perf_totals`]). pub fn sample_reset_perf() -> (u64, u64, u64) { - ( - METER_REQ.swap(0, Ordering::Relaxed), - METER_US.swap(0, Ordering::Relaxed), - METER_MAX_US.swap(0, Ordering::Relaxed), - ) + METER.take_window() +} + +/// Running Electrum totals for `/metrics`: `(requests, sum_us)`. +pub fn perf_totals() -> (u64, u64) { + METER.totals() } fn meter_dispatch_wall(us: u64) { - METER_REQ.fetch_add(1, Ordering::Relaxed); - METER_US.fetch_add(us, Ordering::Relaxed); - let mut cur = METER_MAX_US.load(Ordering::Relaxed); - while us > cur { - match METER_MAX_US.compare_exchange_weak(cur, us, Ordering::Relaxed, Ordering::Relaxed) { - Ok(_) => break, - Err(c) => cur = c, - } - } + METER.note(us); } /// Max simultaneous query-surface clients (Electrum / future Esplora). diff --git a/crates/rbitcoin-esplora/src/lib.rs b/crates/rbitcoin-esplora/src/lib.rs index 351222cce..d660fae3c 100644 --- a/crates/rbitcoin-esplora/src/lib.rs +++ b/crates/rbitcoin-esplora/src/lib.rs @@ -13,5 +13,6 @@ mod server; mod tx_json; pub use server::{ - run_esplora, sample_reset_perf, BlockTemplateFn, EsploraConfig, EsploraHandle, EsploraListen, + perf_totals, run_esplora, sample_reset_perf, BlockTemplateFn, EsploraConfig, EsploraHandle, + EsploraListen, }; diff --git a/crates/rbitcoin-esplora/src/server.rs b/crates/rbitcoin-esplora/src/server.rs index 56a3d8ef3..41102cf74 100644 --- a/crates/rbitcoin-esplora/src/server.rs +++ b/crates/rbitcoin-esplora/src/server.rs @@ -14,6 +14,7 @@ use bitcoin::consensus::Encodable; use bitcoin::Network; use rbitcoin_electrum::ServeLimits; use rbitcoin_net::MempoolHub; +use rbitcoin_net::RequestMeter; use rbitcoin_primitives::Height; use rbitcoin_query::{ChainView, ChainViewKind, Query, ShJoinSlot}; use rbitcoin_store::StoreError; @@ -22,7 +23,7 @@ use serde_json::{json, Value}; use std::collections::HashMap; use std::net::SocketAddr; use std::path::PathBuf; -use std::sync::atomic::{AtomicBool, AtomicU64, Ordering}; +use std::sync::atomic::{AtomicBool, Ordering}; use std::sync::{Arc, Mutex, OnceLock}; use std::time::{Duration, Instant}; use tokio::net::TcpListener; @@ -31,20 +32,18 @@ use tower::limit::ConcurrencyLimitLayer; use tower_http::limit::RequestBodyLimitLayer; use tower_http::timeout::TimeoutLayer; -/// Tip-follow 5s DEBUG `tip: perf`: REST request count this window. -static METER_REQ: AtomicU64 = AtomicU64::new(0); -/// Sum of REST handler walls (µs). -static METER_US: AtomicU64 = AtomicU64::new(0); -/// Max single REST request wall (µs). -static METER_MAX_US: AtomicU64 = AtomicU64::new(0); +/// Esplora REST requests: the `tip: perf` window and the `/metrics` totals. +static METER: RequestMeter = RequestMeter::new(); -/// Sample-and-reset Esplora REST request meters: `(count, sum_us, max_us)`. +/// Esplora REST request meters since the previous sample: `(count, sum_us, max_us)`. +/// Running totals are untouched ([`perf_totals`]). pub fn sample_reset_perf() -> (u64, u64, u64) { - ( - METER_REQ.swap(0, Ordering::Relaxed), - METER_US.swap(0, Ordering::Relaxed), - METER_MAX_US.swap(0, Ordering::Relaxed), - ) + METER.take_window() +} + +/// Running Esplora REST totals for `/metrics`: `(requests, sum_us)`. +pub fn perf_totals() -> (u64, u64) { + METER.totals() } async fn meter_rest(req: Request, next: Next) -> Response { @@ -58,15 +57,7 @@ async fn meter_rest(req: Request, next: Next) -> Response { let resp = next.run(req).await; let elapsed = t0.elapsed(); let us = elapsed.as_micros() as u64; - METER_REQ.fetch_add(1, Ordering::Relaxed); - METER_US.fetch_add(us, Ordering::Relaxed); - let mut cur = METER_MAX_US.load(Ordering::Relaxed); - while us > cur { - match METER_MAX_US.compare_exchange_weak(cur, us, Ordering::Relaxed, Ordering::Relaxed) { - Ok(_) => break, - Err(c) => cur = c, - } - } + METER.note(us); let status = resp.status(); let err = if status.is_success() { None @@ -1141,6 +1132,7 @@ mod tests { use rbitcoin_primitives::{Fk, Height}; use rbitcoin_query::{Query, TxApply}; use rbitcoin_store::{HeaderRecord, InputRecord, OutputRecord, TxRecord}; + use std::sync::atomic::AtomicU64; use tokio::io::{AsyncReadExt, AsyncWriteExt}; use tokio::net::TcpStream; diff --git a/crates/rbitcoin-net/src/lib.rs b/crates/rbitcoin-net/src/lib.rs index 9d40dd884..384ec78b7 100644 --- a/crates/rbitcoin-net/src/lib.rs +++ b/crates/rbitcoin-net/src/lib.rs @@ -20,6 +20,7 @@ mod netgroup; mod peer; mod peer_dos; mod peers; +mod perf_meter; mod reactor; mod seeds; mod serve_perf; @@ -63,6 +64,7 @@ pub use peers::{ pick_stale_follow_evict, DialRequest, DialTarget, LivePeer, PeerConnType, PeerHub, PeerInfo, PeerOut, PingAction, }; +pub use perf_meter::RequestMeter; pub(crate) use rbitcoin_mempool::MempoolGraphStats; pub use rbitcoin_mempool::{AcceptError, Selected}; pub use reactor::BlockingRegion; @@ -71,7 +73,9 @@ pub use seeds::{ resolve_dns_seeds, resolve_fixed_seeds, socks_dns_seed_dests, AddrMan, PeerEntry, PeerFlags, MAX_ADDR_MAN, }; -pub use serve_perf::{format_serve_perf, sample_reset_serve_perf, ServePerfSample}; +pub use serve_perf::{ + format_serve_perf, sample_reset_serve_perf, serve_perf_totals, ServePerfSample, +}; pub use service::P2PNode; pub use socks::{install_i2p_dialer, Dialer}; pub use tx_relay::{ diff --git a/crates/rbitcoin-net/src/perf_meter.rs b/crates/rbitcoin-net/src/perf_meter.rs new file mode 100644 index 000000000..003c69ec3 --- /dev/null +++ b/crates/rbitcoin-net/src/perf_meter.rs @@ -0,0 +1,128 @@ +//! Process meters behind the 5s DEBUG `tip: perf` line and `/metrics`. + +use std::sync::atomic::{AtomicU64, Ordering}; + +/// Monotonic event count. `/metrics` reads the running total; the `tip: perf` +/// line takes the change since its previous sample. +#[derive(Debug, Default)] +pub(crate) struct PerfCounter { + total: AtomicU64, + window_mark: AtomicU64, +} + +impl PerfCounter { + pub const fn new() -> Self { + Self { + total: AtomicU64::new(0), + window_mark: AtomicU64::new(0), + } + } + + pub fn add(&self, n: u64) { + self.total.fetch_add(n, Ordering::Relaxed); + } + + pub fn total(&self) -> u64 { + self.total.load(Ordering::Relaxed) + } + + /// Count since the previous call. The running total is untouched. + pub fn take_window(&self) -> u64 { + let now = self.total(); + now.saturating_sub(self.window_mark.fetch_max(now, Ordering::Relaxed)) + } +} + +/// Largest value noted since the previous [`Self::take`]. A window maximum +/// cannot be derived from running totals, so this one resets. +#[derive(Debug, Default)] +pub(crate) struct PerfMax(AtomicU64); + +impl PerfMax { + pub const fn new() -> Self { + Self(AtomicU64::new(0)) + } + + pub fn note(&self, v: u64) { + self.0.fetch_max(v, Ordering::Relaxed); + } + + pub fn take(&self) -> u64 { + self.0.swap(0, Ordering::Relaxed) + } +} + +/// Request count, summed wall (µs), and window max wall for one serve +/// surface (Esplora REST, Electrum JSON-RPC). +#[derive(Debug, Default)] +pub struct RequestMeter { + requests: PerfCounter, + wall_us: PerfCounter, + max_us: PerfMax, +} + +impl RequestMeter { + pub const fn new() -> Self { + Self { + requests: PerfCounter::new(), + wall_us: PerfCounter::new(), + max_us: PerfMax::new(), + } + } + + pub fn note(&self, wall_us: u64) { + self.requests.add(1); + self.wall_us.add(wall_us); + self.max_us.note(wall_us); + } + + /// `(count, sum_us, max_us)` since the previous call (`tip: perf`). + pub fn take_window(&self) -> (u64, u64, u64) { + ( + self.requests.take_window(), + self.wall_us.take_window(), + self.max_us.take(), + ) + } + + /// Running `(requests, sum_us)` (`/metrics`). + pub fn totals(&self) -> (u64, u64) { + (self.requests.total(), self.wall_us.total()) + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn a_window_sample_leaves_the_running_total() { + let c = PerfCounter::new(); + c.add(1); + c.add(2); + assert_eq!((c.take_window(), c.total()), (3, 3)); + assert_eq!((c.take_window(), c.total()), (0, 3)); + c.add(1); + assert_eq!((c.take_window(), c.total()), (1, 4)); + } + + #[test] + fn a_request_window_keeps_the_totals() { + let m = RequestMeter::new(); + m.note(10); + m.note(30); + assert_eq!(m.take_window(), (2, 40, 30)); + m.note(5); + assert_eq!(m.take_window(), (1, 5, 5)); + assert_eq!(m.totals(), (3, 45)); + } + + #[test] + fn a_window_max_resets_per_sample() { + let m = PerfMax::new(); + m.note(5); + m.note(3); + assert_eq!(m.take(), 5); + assert_eq!(m.take(), 0); + } +} diff --git a/crates/rbitcoin-net/src/serve_perf.rs b/crates/rbitcoin-net/src/serve_perf.rs index c68f44bbf..3860105f1 100644 --- a/crates/rbitcoin-net/src/serve_perf.rs +++ b/crates/rbitcoin-net/src/serve_perf.rs @@ -1,12 +1,12 @@ //! Historical `getdata` witness-block serve meters for the 5s `tip: perf` line. -use std::sync::atomic::{AtomicU64, Ordering}; +use crate::perf_meter::{PerfCounter, PerfMax}; -static SERVE_N: AtomicU64 = AtomicU64::new(0); -static SERVE_BYTES: AtomicU64 = AtomicU64::new(0); -static SERVE_TX: AtomicU64 = AtomicU64::new(0); -static SERVE_WALL_NS: AtomicU64 = AtomicU64::new(0); -static SERVE_MAX_NS: AtomicU64 = AtomicU64::new(0); +static SERVE_N: PerfCounter = PerfCounter::new(); +static SERVE_BYTES: PerfCounter = PerfCounter::new(); +static SERVE_TX: PerfCounter = PerfCounter::new(); +static SERVE_WALL_NS: PerfCounter = PerfCounter::new(); +static SERVE_MAX_NS: PerfMax = PerfMax::new(); /// One 5s window of historical block-serve reconstruct+encode. #[derive(Clone, Copy, Debug, Default, PartialEq, Eq)] @@ -20,30 +20,30 @@ pub struct ServePerfSample { pub(crate) fn note_serve(tx_count: u32, bytes: usize, wall_ns: u128) { let ns = wall_ns.min(u128::from(u64::MAX)) as u64; - SERVE_N.fetch_add(1, Ordering::Relaxed); - SERVE_BYTES.fetch_add(bytes as u64, Ordering::Relaxed); - SERVE_TX.fetch_add(u64::from(tx_count), Ordering::Relaxed); - SERVE_WALL_NS.fetch_add(ns, Ordering::Relaxed); - let mut cur = SERVE_MAX_NS.load(Ordering::Relaxed); - while ns > cur { - match SERVE_MAX_NS.compare_exchange_weak(cur, ns, Ordering::Relaxed, Ordering::Relaxed) { - Ok(_) => break, - Err(c) => cur = c, - } - } + SERVE_N.add(1); + SERVE_BYTES.add(bytes as u64); + SERVE_TX.add(u64::from(tx_count)); + SERVE_WALL_NS.add(ns); + SERVE_MAX_NS.note(ns); } -/// Sample-and-reset serve meters for `DEBUG tip: perf`. +/// Serve meters since the previous sample, for `DEBUG tip: perf`. Running +/// totals are untouched ([`serve_perf_totals`]). pub fn sample_reset_serve_perf() -> ServePerfSample { ServePerfSample { - n: SERVE_N.swap(0, Ordering::Relaxed), - bytes: SERVE_BYTES.swap(0, Ordering::Relaxed), - tx_count: SERVE_TX.swap(0, Ordering::Relaxed), - wall_ns: SERVE_WALL_NS.swap(0, Ordering::Relaxed), - max_ns: SERVE_MAX_NS.swap(0, Ordering::Relaxed), + n: SERVE_N.take_window(), + bytes: SERVE_BYTES.take_window(), + tx_count: SERVE_TX.take_window(), + wall_ns: SERVE_WALL_NS.take_window(), + max_ns: SERVE_MAX_NS.take(), } } +/// Running historical block serves for `/metrics`: `(blocks, bytes)`. +pub fn serve_perf_totals() -> (u64, u64) { + (SERVE_N.total(), SERVE_BYTES.total()) +} + /// `serve n= bytes= tx= avg_us= max_us=` — reconstruct+encode, not BIP324 send. pub fn format_serve_perf(s: &ServePerfSample) -> String { let avg_us = s.wall_ns.checked_div(s.n).unwrap_or(0) / 1_000; diff --git a/crates/rbitcoin-net/src/tx_relay.rs b/crates/rbitcoin-net/src/tx_relay.rs index 600726f8e..f0fcd60db 100644 --- a/crates/rbitcoin-net/src/tx_relay.rs +++ b/crates/rbitcoin-net/src/tx_relay.rs @@ -6,6 +6,7 @@ //! is not in rust-bitcoin 0.32 `NetworkMessage`; the old private `rbtpkg` //! name is gone). +use crate::perf_meter::{PerfCounter, PerfMax}; use arc_swap::ArcSwap; use bitcoin::hashes::Hash; use bitcoin::{Amount, OutPoint, ScriptBuf, Transaction, TxOut, Txid, Wtxid}; @@ -571,10 +572,10 @@ pub struct MempoolHub { tx_snapshot: ArcSwap, tx_snap_dirty: AtomicBool, tx_snap_refreshing: AtomicBool, - meter_accepts: AtomicU64, - meter_rejects: AtomicU64, + meter_accepts: PerfCounter, + meter_rejects: PerfCounter, meter_accept_us: AtomicU64, - meter_accept_max_us: AtomicU64, + meter_accept_max_us: PerfMax, meter_accept_lock_us: AtomicU64, meter_accept_utxo_us: AtomicU64, meter_accept_script_us: AtomicU64, @@ -726,10 +727,10 @@ impl MempoolHub { tx_snapshot: ArcSwap::from_pointee(MempoolTxSnapshot::empty(Instant::now())), tx_snap_dirty: AtomicBool::new(true), tx_snap_refreshing: AtomicBool::new(false), - meter_accepts: AtomicU64::new(0), - meter_rejects: AtomicU64::new(0), + meter_accepts: PerfCounter::new(), + meter_rejects: PerfCounter::new(), meter_accept_us: AtomicU64::new(0), - meter_accept_max_us: AtomicU64::new(0), + meter_accept_max_us: PerfMax::new(), meter_accept_lock_us: AtomicU64::new(0), meter_accept_utxo_us: AtomicU64::new(0), meter_accept_script_us: AtomicU64::new(0), @@ -999,23 +1000,12 @@ impl MempoolHub { fn meter_accept_wall(&self, us: u64, ok: bool) { if ok { - self.meter_accepts.fetch_add(1, Ordering::Relaxed); + self.meter_accepts.add(1); } else { - self.meter_rejects.fetch_add(1, Ordering::Relaxed); + self.meter_rejects.add(1); } self.meter_accept_us.fetch_add(us, Ordering::Relaxed); - let mut cur = self.meter_accept_max_us.load(Ordering::Relaxed); - while us > cur { - match self.meter_accept_max_us.compare_exchange_weak( - cur, - us, - Ordering::Relaxed, - Ordering::Relaxed, - ) { - Ok(_) => break, - Err(c) => cur = c, - } - } + self.meter_accept_max_us.note(us); } fn meter_accept_stages(&self, lock_us: u64, stages: rbitcoin_mempool::AcceptStageUs) { @@ -1029,13 +1019,20 @@ impl MempoolHub { .fetch_add(stages.durable_us, Ordering::Relaxed); } - /// Sample-and-reset mempool/relay counters for the tip-follow 5s DEBUG line. + /// Running mempool accepts and rejects for `/metrics`. A + /// [`Self::sample_reset_perf`] window does not reset them. + pub fn accept_totals(&self) -> (u64, u64) { + (self.meter_accepts.total(), self.meter_rejects.total()) + } + + /// Mempool/relay counters since the previous sample, for the tip-follow + /// 5s DEBUG line. pub fn sample_reset_perf(&self) -> MempoolPerfSample { MempoolPerfSample { - accepts: self.meter_accepts.swap(0, Ordering::Relaxed), - rejects: self.meter_rejects.swap(0, Ordering::Relaxed), + accepts: self.meter_accepts.take_window(), + rejects: self.meter_rejects.take_window(), accept_us: self.meter_accept_us.swap(0, Ordering::Relaxed), - accept_max_us: self.meter_accept_max_us.swap(0, Ordering::Relaxed), + accept_max_us: self.meter_accept_max_us.take(), accept_lock_us: self.meter_accept_lock_us.swap(0, Ordering::Relaxed), accept_utxo_us: self.meter_accept_utxo_us.swap(0, Ordering::Relaxed), accept_script_us: self.meter_accept_script_us.swap(0, Ordering::Relaxed), diff --git a/crates/rbitcoin-node/src/health/metrics.rs b/crates/rbitcoin-node/src/health/metrics.rs index 67f008313..120cf447d 100644 --- a/crates/rbitcoin-node/src/health/metrics.rs +++ b/crates/rbitcoin-node/src/health/metrics.rs @@ -87,10 +87,58 @@ pub(super) fn render(status: &NodeStatus) -> String { "Sum of mempool virtual sizes (getmempoolinfo.bytes).", vbytes, ); + let (accepts, rejects) = mempool.accept_totals(); + out.counter( + "rbitcoin_mempool_accepts_total", + "Transactions the mempool accepted (tip: perf accepts=).", + accepts, + ); + out.counter( + "rbitcoin_mempool_rejects_total", + "Transactions the mempool rejected (tip: perf rejects=).", + rejects, + ); } + let (requests, us) = rbitcoin_esplora::perf_totals(); + out.counter( + "rbitcoin_esplora_requests_total", + "Esplora REST requests (tip: perf esplora req=).", + requests, + ); + out.counter( + "rbitcoin_esplora_request_seconds_total", + "Esplora REST handler wall time in seconds.", + seconds(us), + ); + let (requests, us) = rbitcoin_electrum::perf_totals(); + out.counter( + "rbitcoin_electrum_requests_total", + "Electrum JSON-RPC requests (tip: perf electrum req=).", + requests, + ); + out.counter( + "rbitcoin_electrum_request_seconds_total", + "Electrum dispatch wall time in seconds.", + seconds(us), + ); + let (blocks, bytes) = rbitcoin_net::serve_perf_totals(); + out.counter( + "rbitcoin_block_serve_total", + "Historical blocks served to peers (tip: perf serve n=).", + blocks, + ); + out.counter( + "rbitcoin_block_serve_bytes_total", + "Bytes of historical blocks served to peers (tip: perf serve bytes=).", + bytes, + ); out.0 } +fn seconds(us: u64) -> f64 { + us as f64 / 1e6 +} + fn chain_gauges(out: &mut Exposition, chain: &ChainHub, sh_index: bool) { let tip = chain.query.tip_height(); out.gauge( @@ -141,4 +189,9 @@ impl Exposition { self.family(name, "gauge", help); self.sample(name, "", value); } + + fn counter(&mut self, name: &str, help: &str, value: impl Display) { + self.family(name, "counter", help); + self.sample(name, "", value); + } } diff --git a/crates/rbitcoin-test/tests/cross_surface.rs b/crates/rbitcoin-test/tests/cross_surface.rs index 7f41de1e1..50f3d740b 100644 --- a/crates/rbitcoin-test/tests/cross_surface.rs +++ b/crates/rbitcoin-test/tests/cross_surface.rs @@ -1985,7 +1985,29 @@ async fn esplora_broadcast_visible_in_rpc_and_electrum() { "{chain}" ); assert_eq!(http_get(health_addr, "/readyz").await, (200, "ok".into())); - pin_metrics_equal_rpc(health_addr, rpc_addr, true).await; + let metrics_ready = pin_metrics_equal_rpc(health_addr, rpc_addr, true).await; + for total in [ + "rbitcoin_esplora_requests_total", + "rbitcoin_esplora_request_seconds_total", + "rbitcoin_electrum_requests_total", + "rbitcoin_electrum_request_seconds_total", + "rbitcoin_block_serve_total", + "rbitcoin_block_serve_bytes_total", + "rbitcoin_mempool_accepts_total", + "rbitcoin_mempool_rejects_total", + ] { + assert!( + metrics_ready[total] >= metrics_in_ibd[total], + "{total} counts up: {metrics_in_ibd:?} then {metrics_ready:?}" + ); + } + for total in [ + "rbitcoin_esplora_requests_total", + "rbitcoin_electrum_requests_total", + "rbitcoin_mempool_accepts_total", + ] { + assert!(metrics_in_ibd[total] >= 1.0, "{total}: {metrics_in_ibd:?}"); + } let relay_parent = acs_spend( relay_cb, From 521db7f9089c0b2a283468c338760f69a6d58311 Mon Sep 17 00:00:00 2001 From: Benjamen Keroack Date: Mon, 28 Sep 2026 12:30:25 -0600 Subject: [PATCH 05/11] docs: health probes, readiness reasons, and metrics operations.md gets the two flags and a "Health probes and metrics" section: the routes, the /readyz reasons in check order, a k8s probe snippet that keeps liveness on /healthz, the metric table with the RPC field or log token each one equals, and a scrape config. SECURITY.md scopes the listener. peer-clients.md rank 4 is landed. A changelog.d fragment covers the new flags and the meters becoming running totals. Co-Authored-By: Claude Opus 5.5 --- OPERATOR.md | 2 +- SECURITY.md | 5 ++ changelog.d/health-endpoints.md | 20 ++++++++ docs/operator/operations.md | 87 +++++++++++++++++++++++++++++++++ docs/peer-clients.md | 4 +- 5 files changed, 115 insertions(+), 3 deletions(-) create mode 100644 changelog.d/health-endpoints.md diff --git a/OPERATOR.md b/OPERATOR.md index d978ebeb6..504bd2aaa 100644 --- a/OPERATOR.md +++ b/OPERATOR.md @@ -5,7 +5,7 @@ Start with the task you need to do: | Task | Guide | |------|-------| | Understand node status, build a binary, configure a service | [Setup and installation](docs/operator/setup.md) | -| Try regtest, use the CLI, monitor the node, configure relay and P2P | [Node operations](docs/operator/operations.md) | +| Try regtest, use the CLI, monitor the node (logs, health probes, metrics), configure relay and P2P | [Node operations](docs/operator/operations.md) | | Tune store IO and memory, upgrade the schema, manage optional indexes | [Storage and indexes](docs/operator/storage.md) | | Configure Electrum, Esplora, RPC, or client benchmarking | [Client interfaces](docs/operator/interfaces.md) | | Run signet/mainnet, or tune a constrained host or uplink | [Labs and constrained hosts](docs/operator/field-notes.md) | diff --git a/SECURITY.md b/SECURITY.md index ece41a5b4..478d5b342 100644 --- a/SECURITY.md +++ b/SECURITY.md @@ -64,6 +64,11 @@ that affect consensus, P2P attack surface, or Electrum/query integrity. opt-in (`--esplora-listen`). Internal electrs HTTP (`/internal/*`) is unix listen only, not the public TCP bind. Edge TLS, multi-tenant metering, and API keys are still out of process (see [`OPERATOR.md`](./OPERATOR.md)). +- **Health listener:** `--health-listen` is opt-in, read-only, and + unauthenticated (`/healthz`, `/readyz`, and `/metrics` with `--metrics`). + It binds before the store opens and keeps the same always-on concurrency, + body, and timeout limits. `/metrics` reveals peer counts and mempool size; + keep the port on loopback or a probe-only network. - **Store / archive:** corruption or incorrect spend/scripthash results that mislead a **wallet** backend are in scope. Truncating a sealed mmap after it is mapped is fatal external corruption (the next read can SIGBUS), not diff --git a/changelog.d/health-endpoints.md b/changelog.d/health-endpoints.md new file mode 100644 index 000000000..3570b278d --- /dev/null +++ b/changelog.d/health-endpoints.md @@ -0,0 +1,20 @@ +Added + +- **Health probes.** `--health-listen [ADDR]` (default `127.0.0.1:9332`) + binds before the store opens and serves `GET /healthz` (200 in every + phase) and `GET /readyz` (200 once the node follows the tip with every + configured listener up, the tip within 6 blocks of the best header, and + the scripthash index within 6 blocks of the tip; otherwise 503 with the + reason). RPC, Electrum, and Esplora bind only after catch-up, so a probe + on those would restart a node in the middle of a migration or IBD. +- **Prometheus metrics.** `--metrics` adds `GET /metrics` on the health + listener. Gauges equal their RPC fields (`blocks`, `headers`, + `initialblockdownload`, connections, mempool size) and counters are the + `tip: perf` meters as running totals. + +Changed + +- **`tip: perf` meters are running totals.** Esplora and Electrum + requests, historical block serves, and mempool accepts and rejects count + up for the life of the process. The 5 s DEBUG line still prints the + change since the previous line. diff --git a/docs/operator/operations.md b/docs/operator/operations.md index 12a6f4e4a..f9141568b 100644 --- a/docs/operator/operations.md +++ b/docs/operator/operations.md @@ -100,6 +100,8 @@ Clean smoke: | `--electrum-listen [ADDR]` | `electrum_listen=` | disabled; omit ADDR → `127.0.0.1:50001`. Address/scripthash methods need `--sh-index` | | `--esplora-listen [ADDR\|PATH]` | `esplora_listen=` | disabled (Esplora REST); omit ADDR → `127.0.0.1:3000`; a filesystem path is unix HTTP (mode **0660**, dummy `Host: api` is fine). Address/scripthash methods need `--sh-index` | | `--esplora-onion[=0\|1]` | `esplora_onion=` | **on** — `ADD_ONION` for Esplora when `--tor-control` is set | +| `--health-listen [ADDR]` | `health_listen=` | disabled; omit ADDR → `127.0.0.1:9332`. Unauthenticated `GET /healthz` and `/readyz`, bound before the store opens ([Health probes and metrics](#health-probes-and-metrics)). A bind failure stops the node | +| `--metrics[=0\|1]` | `metrics=` | **off** — Prometheus `GET /metrics` on the health listener. Refused without `--health-listen` | | `--esplora-block-template` | `esplora_block_template=` | **off** — `GET /block-template` is 404; on = GBT JSON (same as RPC template mode) | | `--rpc` | `rpc=` | **off** — unix JSON-RPC `{datadir}/rpc.sock` (mode 0600) | | `--rpc-listen [ADDR]` | `rpc_listen=` | disabled — implies `--rpc`; omit ADDR → `127.0.0.1` and Core-matching RPC port | @@ -344,6 +346,91 @@ uniformly slow pack). Slow or constrained uplinks: [Slow / constrained uplink SH head/body, and spenders are fd pread/pwrite**. Full modality matrix: [`docs/io-modality.md`](docs/io-modality.md). +## Health probes and metrics + +`--health-listen [ADDR]` binds a small HTTP listener at the very start of +`rbitcoin-node`, before the store opens. RPC, Electrum, and Esplora bind +only after the store opens, catch-up finishes, and the scripthash index +materializes; on mainnet those can take minutes (a schema backfill) to +hours (IBD). The health listener answers through all of it. + +| Route | Answer | +|-------|--------| +| `GET /healthz` | `200 ok` in every phase. The process is up and serving HTTP | +| `GET /readyz` | `200 ok` when the node follows the tip and every configured listener is up; otherwise `503 not ready: ` | +| `GET /metrics` | Prometheus text format with `--metrics`; 404 without it | + +Other paths are 404 and other methods 405. Requests carry no body and no +auth. Concurrency, body, and timeout limits are always on. A non-loopback +bind logs a WARN: keep the port on loopback or a probe-only network. + +`/readyz` reports the first failing check, in this order: + +| Reason | Meaning | +|--------|---------| +| `opening` | Store open: schema migration, backfill, spend replay | +| `starting` | P2P, mempool, and proxy bring-up | +| `catch-up` | Initial block download | +| `indexing` | Tip-mode entry (scripthash materialize), follow peers, listeners | +| `stopping` | Shutdown flush | +| `rpc not listening` (also `electrum`, `esplora`) | That listener is configured but did not bind. Its start only warns, so the node keeps following | +| `initial block download` | Same value as RPC `getblockchaininfo.initialblockdownload` (`--max-tip-age`, `--min-chain-work`). Latches off after the first exit, as in Core | +| `tip N blocks behind headers` | `headers - blocks` is over 6 | +| `scripthash index N blocks behind tip` | With `--sh-index`, the index trails the tip by more than 6 (`sh_lag=` on `tip: accept`) | + +Point liveness at `/healthz` and readiness at `/readyz`. Never use `/readyz` +for liveness: a restart during a migration or IBD starts that work over. + +```yaml +livenessProbe: + httpGet: { path: /healthz, port: 9332 } + periodSeconds: 10 +readinessProbe: + httpGet: { path: /readyz, port: 9332 } + periodSeconds: 10 +``` + +A k8s `httpGet` probe connects to the pod IP, so in a pod use +`--health-listen 0.0.0.0:9332` and keep the port off any Service. + +### Metrics + +`--metrics` (needs `--health-listen`) adds `GET /metrics`. Each gauge is a +value the node already publishes on RPC or a log line, under a name that +says which. A scrape reads a few atomics and one peer snapshot, and folds +the mempool once for `rbitcoin_mempool_bytes` (the same fold as +`getmempoolinfo`) on the blocking pool. + +| Metric | Type | Equals | +|--------|------|--------| +| `rbitcoin_build_info{version,network}` | gauge | `1`; version and `network=` from the startup line | +| `rbitcoin_phase{phase}` | gauge | `1` for the current `/readyz` phase, `0` for the others | +| `rbitcoin_ready` | gauge | `1` when `/readyz` is 200 | +| `rbitcoin_blocks` | gauge | `getblockchaininfo.blocks` | +| `rbitcoin_headers` | gauge | `getblockchaininfo.headers` | +| `rbitcoin_tip_time_seconds` | gauge | `getblockchaininfo.time` | +| `rbitcoin_initial_block_download` | gauge | `getblockchaininfo.initialblockdownload` | +| `rbitcoin_connections{direction="in"\|"out"}` | gauge | `getnetworkinfo.connections_in` / `connections_out` | +| `rbitcoin_mempool_transactions` | gauge | `getmempoolinfo.size` | +| `rbitcoin_mempool_bytes` | gauge | `getmempoolinfo.bytes` (virtual size) | +| `rbitcoin_scripthash_lag_blocks` | gauge | `tip: accept sh_lag=` (with `--sh-index`) | +| `rbitcoin_esplora_requests_total` / `_request_seconds_total` | counter | `tip: perf esplora req=` / `avg_us` × `req` | +| `rbitcoin_electrum_requests_total` / `_request_seconds_total` | counter | `tip: perf electrum req=` / `avg_us` × `req` | +| `rbitcoin_block_serve_total` / `_bytes_total` | counter | `tip: perf serve n=` / `bytes=` | +| `rbitcoin_mempool_accepts_total` / `_rejects_total` | counter | `tip: perf accepts=` / `rejects=` | +| `process_resident_memory_bytes` | gauge | `ibd: sizes rss=` (Linux and macOS) | +| `process_start_time_seconds` | gauge | Unix time `rbitcoin-node` started | + +Hub gauges appear once P2P has started. `/metrics` exposes peer and mempool +counts: do not publish it without a firewall. + +```yaml +scrape_configs: + - job_name: rbitcoin + static_configs: + - targets: ["127.0.0.1:9332"] +``` + ## Libre-relay-class policy (mempool + Electrum broadcast) | Rule | Value | diff --git a/docs/peer-clients.md b/docs/peer-clients.md index 3ec355bc8..47d27d040 100644 --- a/docs/peer-clients.md +++ b/docs/peer-clients.md @@ -214,10 +214,10 @@ scheduling a slice. | 1 | Height-1 + spend-pad + 2-block fork vs v31.1 `bitcoind`; BIP324 `v2_contents` ASan + live Core `v2_session` + compact reconstruct vs `getblocktxn` + script-mutating vs Core + compact reorg via `drain_pending` (**landed**; **Q-30** Completed) | satd `block_differential` | **Q-30** / [`TESTING.md`](../TESTING.md) | | 2 | One cross-surface scenario: Esplora `POST /tx` → Electrum history + RPC mempool (**landed**; `esplora_broadcast_visible_in_rpc_and_electrum`) | satd E2E | `rbitcoin-test` `--test cross_surface` | | — | ~~Hornet spec.html vs consensus-tests.md gap hunt~~ **done 2026-09-04** (table in this file; pins in `structure_rule_tests` / `header.rs` / `consensus_rules`) | Hornet | this file + [`consensus-tests.md`](./consensus-tests.md) | -| 4 | `/healthz` (and maybe `/readyz`) on the node listen; Prometheus later as a flag | satd | node / [`OPERATOR.md`](../OPERATOR.md) | +| 4 | `/healthz` + `/readyz` on the node listen; Prometheus as a flag (**landed**; `--health-listen`, `--metrics`) | satd | node / [`operations.md`](./operator/operations.md#health-probes-and-metrics) | | 5 | BIP352 serve: hash-bind tweak batches so a client can audit the stream | satd row idea on our tweaks path | Electrum tweaks / [`OPERATOR.md`](../OPERATOR.md) | -1–2 landed. 4–5 are small product. None require becoming a UTXO node or a +1, 2, and 4 landed. 5 is small product. None require becoming a UTXO node or a Core conf clone. --- From 15bae31f20896a072e9b31c359302129fa4c6daf Mon Sep 17 00:00:00 2001 From: Benjamen Keroack Date: Mon, 28 Sep 2026 13:20:43 -0600 Subject: [PATCH 06/11] node: /readyz gates on a stale tip, past --max-tip-age in_ibd latches off after the first exit and headers - blocks cannot grow when the node has no peers, so a node that lost every peer after IBD answered /readyz 200 with a stale tip while a load balancer kept routing wallets to it. Add the non-latching half of the IBD staleness check to readiness: tip_header age over ChainHub::max_tip_age_secs (--max-tip-age, default 24h) answers 503 "tip stale (last block Ns ago)". rbitcoin_ready inherits the gate through the same readiness(). Generated with [Devin](https://devin.ai) Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com> --- changelog.d/health-endpoints.md | 11 ++++++---- crates/rbitcoin-node/src/health.rs | 33 ++++++++++++++++++++++++++++++ docs/operator/operations.md | 1 + 3 files changed, 41 insertions(+), 4 deletions(-) diff --git a/changelog.d/health-endpoints.md b/changelog.d/health-endpoints.md index 3570b278d..ced1ac62b 100644 --- a/changelog.d/health-endpoints.md +++ b/changelog.d/health-endpoints.md @@ -3,10 +3,13 @@ Added - **Health probes.** `--health-listen [ADDR]` (default `127.0.0.1:9332`) binds before the store opens and serves `GET /healthz` (200 in every phase) and `GET /readyz` (200 once the node follows the tip with every - configured listener up, the tip within 6 blocks of the best header, and - the scripthash index within 6 blocks of the tip; otherwise 503 with the - reason). RPC, Electrum, and Esplora bind only after catch-up, so a probe - on those would restart a node in the middle of a migration or IBD. + configured listener up, the tip within 6 blocks of the best header, the + tip fresher than `--max-tip-age`, and the scripthash index within 6 + blocks of the tip; otherwise 503 with the reason). The stale-tip check + does not latch like `initialblockdownload`, so a node that loses every + peer after IBD goes unready. RPC, Electrum, and Esplora bind only after + catch-up, so a probe on those would restart a node in the middle of a + migration or IBD. - **Prometheus metrics.** `--metrics` adds `GET /metrics` on the health listener. Gauges equal their RPC fields (`blocks`, `headers`, `initialblockdownload`, connections, mempool size) and counters are the diff --git a/crates/rbitcoin-node/src/health.rs b/crates/rbitcoin-node/src/health.rs index 9292224a7..20cf6c405 100644 --- a/crates/rbitcoin-node/src/health.rs +++ b/crates/rbitcoin-node/src/health.rs @@ -142,6 +142,8 @@ impl NodeStatus { in_ibd: false, blocks: 0, headers: 0, + tip_age_secs: None, + max_tip_age_secs: 0, sh_lag: None, }; if phase != Phase::Following { @@ -151,6 +153,10 @@ impl NodeStatus { snap.in_ibd = chain.in_ibd(); snap.blocks = chain.query.tip_height().map_or(0, |h| h.0); snap.headers = chain.best_header_height(); + snap.tip_age_secs = chain + .tip_header() + .map(|h| chain.clock.now_secs().saturating_sub(u64::from(h.time))); + snap.max_tip_age_secs = chain.max_tip_age_secs(); snap.sh_lag = self.sh_index.then(|| chain.query.sh_lag_heights()); } snap @@ -166,6 +172,10 @@ struct ReadySnapshot { in_ibd: bool, blocks: u32, headers: u32, + /// Seconds since the tip block time; `None` with no tip. `in_ibd` latches + /// off after the first exit — this one keeps reporting a stale tip. + tip_age_secs: Option, + max_tip_age_secs: u64, /// `None` without `--sh-index`. sh_lag: Option, } @@ -181,6 +191,11 @@ fn readiness(s: &ReadySnapshot) -> Result<(), String> { if s.in_ibd { return Err("initial block download".into()); } + if let Some(age) = s.tip_age_secs { + if age > s.max_tip_age_secs { + return Err(format!("tip stale (last block {age}s ago)")); + } + } let behind = s.headers.saturating_sub(s.blocks); if behind > READY_LAG_BLOCKS { return Err(format!("tip {behind} blocks behind headers")); @@ -289,6 +304,8 @@ mod tests { in_ibd: false, blocks: 100, headers: 100, + tip_age_secs: Some(0), + max_tip_age_secs: 24 * 60 * 60, sh_lag: None, } } @@ -356,6 +373,22 @@ mod tests { }, Ok(()), )); + let stale = |age: u64| ReadySnapshot { + tip_age_secs: Some(age), + ..following() + }; + cases.push((stale(24 * 60 * 60), Ok(()))); + cases.push(( + stale(24 * 60 * 60 + 1), + Err("tip stale (last block 86401s ago)".into()), + )); + cases.push(( + ReadySnapshot { + tip_age_secs: None, + ..following() + }, + Ok(()), + )); for (snap, want) in cases { assert_eq!(readiness(&snap), want, "{snap:?}"); } diff --git a/docs/operator/operations.md b/docs/operator/operations.md index f9141568b..eae316509 100644 --- a/docs/operator/operations.md +++ b/docs/operator/operations.md @@ -375,6 +375,7 @@ bind logs a WARN: keep the port on loopback or a probe-only network. | `stopping` | Shutdown flush | | `rpc not listening` (also `electrum`, `esplora`) | That listener is configured but did not bind. Its start only warns, so the node keeps following | | `initial block download` | Same value as RPC `getblockchaininfo.initialblockdownload` (`--max-tip-age`, `--min-chain-work`). Latches off after the first exit, as in Core | +| `tip stale (last block Ns ago)` | The non-latching half of that check: the tip block is older than `--max-tip-age` (default 24h). This still fires after IBD has latched off, so a node that later loses every peer stops being ready | | `tip N blocks behind headers` | `headers - blocks` is over 6 | | `scripthash index N blocks behind tip` | With `--sh-index`, the index trails the tip by more than 6 (`sh_lag=` on `tip: accept`) | From a0885d9aad94607e28816ccffbbb1365331dde93 Mon Sep 17 00:00:00 2001 From: Benjamen Keroack Date: Mon, 28 Sep 2026 13:20:52 -0600 Subject: [PATCH 07/11] node: name /metrics in the non-loopback health-bind warn With --metrics the warning listed /healthz and /readyz only, but /metrics is the endpoint that exposes peer counts and mempool size. Generated with [Devin](https://devin.ai) Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com> --- crates/rbitcoin-node/src/health.rs | 5 ++++- 1 file changed, 4 insertions(+), 1 deletion(-) diff --git a/crates/rbitcoin-node/src/health.rs b/crates/rbitcoin-node/src/health.rs index 20cf6c405..4af42c2a9 100644 --- a/crates/rbitcoin-node/src/health.rs +++ b/crates/rbitcoin-node/src/health.rs @@ -229,7 +229,10 @@ pub(crate) async fn run_health( let listener = TcpListener::bind(addr).await?; let local_addr = listener.local_addr()?; if !local_addr.ip().is_loopback() { - warn!("health: {local_addr} is not loopback; /healthz and /readyz are unauthenticated"); + warn!( + "health: {local_addr} is not loopback; /healthz and /readyz{} are unauthenticated", + if metrics { " and /metrics" } else { "" } + ); } let app = router(status, metrics); let task = tokio::spawn(async move { From 8bf62a626fc9c53e85353392e3a22f7f77f1c709 Mon Sep 17 00:00:00 2001 From: Benjamen Keroack Date: Tue, 29 Sep 2026 20:17:13 -0600 Subject: [PATCH 08/11] net: walk header ancestry with no chain guard held best_header_height held the header_tips read guard while header_ancestry_invalid -> prev_of read-locked header_tips again. std's RwLock refuses a new reader while a writer waits, so a note_header_tip arriving between the two reads deadlocked both threads: the writer waits for the held read, the second read waits for the writer. getblockchaininfo already raced this; a Prometheus scrape of rbitcoin_headers during header sync makes it routine, and the stuck blocking task never drops its guard. Both walks that re-read a guard they hold now copy what they need out first: - best_header_height copies the tips, and only walks one that is taller than the best height so far. - chaintips copies header_tips and held_bodies before walking them (prev_of, side_height_and_branchlen and held_path_has_body_gap re-read both, held_bodies through load_side_body). The test races both against writers on the two locks; with either lock's copy reverted in either function it deadlocks within its 20s watchdog. Co-Authored-By: Claude Opus 5.5 --- crates/rbitcoin-net/src/chain.rs | 96 ++++++++++++++++++++++++++++---- 1 file changed, 84 insertions(+), 12 deletions(-) diff --git a/crates/rbitcoin-net/src/chain.rs b/crates/rbitcoin-net/src/chain.rs index 61e0cb6d0..adb668d39 100644 --- a/crates/rbitcoin-net/src/chain.rs +++ b/crates/rbitcoin-net/src/chain.rs @@ -853,12 +853,17 @@ impl ChainHub { for h in self.fork_tips.read().unwrap().iter().copied() { record(&mut out, h, "valid-fork"); } + // The header and held walks below re-read `header_tips` and + // `held_bodies` (`prev_of`, `load_side_body`), so each set is copied + // out and its guard dropped first; see `best_header_height`. { - let headers = self.header_tips.read().unwrap(); + let (covered, hashes): (HashSet, Vec) = { + let headers = self.header_tips.read().unwrap(); + (headers.prevs().collect(), headers.hashes().collect()) + }; // Only header *tips* (a later submitblock of an ancestor must not // re-list that ancestor alongside its descendant). - let covered: HashSet = headers.prevs().collect(); - for hash in headers.hashes() { + for hash in hashes { if covered.contains(&hash) { continue; } @@ -871,10 +876,14 @@ impl ChainHub { } } { - let held = self.held_bodies.read().unwrap(); - let parents: HashSet = - held.blocks().map(|b| b.header.prev_blockhash).collect(); - for hash in held.keys() { + let (parents, hashes): (HashSet, Vec) = { + let held = self.held_bodies.read().unwrap(); + ( + held.blocks().map(|b| b.header.prev_blockhash).collect(), + held.keys().collect(), + ) + }; + for hash in hashes { if parents.contains(&hash) { continue; } @@ -989,14 +998,18 @@ impl ChainHub { } /// Best known header height (may lead `blocks` after `submitheader`). + /// + /// Copies the tips out before walking: `prev_of` read-locks `header_tips` + /// again, and std's `RwLock` refuses a new reader while a writer waits, so + /// a walk under the guard deadlocks against `note_header_tip`. pub fn best_header_height(&self) -> u32 { let mut best = self.tip_height().unwrap_or(0); - let headers = self.header_tips.read().unwrap(); - for (hash, h) in headers.entries() { - if self.header_ancestry_invalid(hash) { - continue; + let tips: Vec<(BlockHash, u32)> = self.header_tips.read().unwrap().entries().collect(); + for (hash, h) in tips { + // Only a taller tip can raise `best`; the walk is for those alone. + if h > best && !self.header_ancestry_invalid(hash) { + best = h; } - best = best.max(h); } best } @@ -4496,6 +4509,65 @@ mod tests { let _ = std::fs::remove_dir_all(dir); } + /// `best_header_height` and `chaintips` walk ancestry through `prev_of`, + /// which read-locks `header_tips` and (via `load_side_body`) + /// `held_bodies`. std's `RwLock` refuses a new reader while a writer + /// waits, so a walk that holds one of those guards deadlocks against a + /// writer: the writer waits for the held read, and the walk's second read + /// waits for the writer. + #[test] + fn ancestry_walks_never_reread_a_held_chain_lock() { + use std::sync::mpsc::{self, RecvTimeoutError}; + use std::thread; + use std::time::{Duration, Instant}; + + let (_dir, hub) = tmp_hub(); + hub.ensure_genesis().unwrap(); + let gen = hub.tip_hash().unwrap(); + let b1 = mine(gen, 1_300_030_000, 1); + hub.accept_block(b1.clone()).unwrap(); + // A header-only tip (`header_tips`) and a parked side body (`held_bodies`). + let child = mine(b1.block_hash(), 1_300_030_100, 2); + hub.ensure_header(&child.header).unwrap(); + let side = mine(gen, 1_300_030_200, 1); + hub.hold_unconnected_body(side.clone()); + assert!(hub.held_body(&side.block_hash()).is_some()); + let tip = child.block_hash(); + + let hub = ChainHub::into_arc(hub); + let stop = Arc::new(AtomicBool::new(false)); + let writer = { + let (hub, stop) = (hub.clone(), stop.clone()); + thread::spawn(move || { + while !stop.load(Ordering::Relaxed) { + drop(hub.header_tips.write().unwrap()); + drop(hub.held_bodies.write().unwrap()); + } + }) + }; + let (done_tx, done_rx) = mpsc::channel(); + { + let hub = hub.clone(); + thread::spawn(move || { + let until = Instant::now() + Duration::from_secs(2); + while Instant::now() < until { + assert_eq!(hub.best_header_height(), 2); + assert!(hub.chaintips().iter().any(|t| t.hash == tip)); + } + let _ = done_tx.send(()); + }); + } + let outcome = done_rx.recv_timeout(Duration::from_secs(20)); + stop.store(true, Ordering::Relaxed); + match outcome { + Ok(()) => writer.join().unwrap(), + Err(RecvTimeoutError::Timeout) => { + panic!("an ancestry walk deadlocked against a chain-lock writer") + } + Err(RecvTimeoutError::Disconnected) => panic!("the walk thread panicked"), + } + } + #[test] fn accept_competing_tip_and_block_at_height_paths() { let (dir, hub) = tmp_hub(); From fd54d9872474e95ce1207272734df9e87bb014c8 Mon Sep 17 00:00:00 2001 From: Benjamen Keroack Date: Tue, 29 Sep 2026 20:25:07 -0600 Subject: [PATCH 09/11] node: health probes answer 503 past their cap; a timed-out check keeps its slot The listener ran tower's ConcurrencyLimitLayer outside the timeout. That limit queues without a bound or a deadline, and axum's Router::layer builds one per route, so /healthz, /readyz and /metrics each had 16. Inside it, a request that hit the 5s timeout was dropped and released its permit, but its spawn_blocking task cannot be cancelled and kept running, so a prober on a non-loopback bind could stack chain and mempool readers past any cap. The timeout is now outermost, and the blocking work is gated in place of the request. /readyz takes one of 4 permits and /metrics the only one, with try_acquire, so past the cap the answer is 503 at once ("not ready: busy", "metrics: a scrape is already running"). The permit moves into the blocking task: it is released when the work ends, not when the request gives up, and /metrics never starts a second mempool fold while one runs. /healthz does no work and is not gated. rbitcoin-node no longer depends on tower. operations.md and SECURITY.md describe the caps and the 503 answers. Co-Authored-By: Claude Opus 5.5 --- Cargo.lock | 1 - SECURITY.md | 7 +- crates/rbitcoin-node/Cargo.toml | 1 - crates/rbitcoin-node/src/health.rs | 190 ++++++++++++++++++++++++----- docs/operator/operations.md | 9 +- 5 files changed, 171 insertions(+), 37 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index 16c449680..64707f78a 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -782,7 +782,6 @@ dependencies = [ "rbitcoin-store", "serde_json", "tokio", - "tower", "tower-http", ] diff --git a/SECURITY.md b/SECURITY.md index 478d5b342..77ba58f49 100644 --- a/SECURITY.md +++ b/SECURITY.md @@ -66,9 +66,10 @@ that affect consensus, P2P attack surface, or Electrum/query integrity. are still out of process (see [`OPERATOR.md`](./OPERATOR.md)). - **Health listener:** `--health-listen` is opt-in, read-only, and unauthenticated (`/healthz`, `/readyz`, and `/metrics` with `--metrics`). - It binds before the store opens and keeps the same always-on concurrency, - body, and timeout limits. `/metrics` reveals peer counts and mempool size; - keep the port on loopback or a probe-only network. + It binds before the store opens. Requests carry no body and time out; + `/readyz` and `/metrics` run a capped number of blocking reads and answer + 503 past the cap rather than queue. `/metrics` reveals peer counts and + mempool size; keep the port on loopback or a probe-only network. - **Store / archive:** corruption or incorrect spend/scripthash results that mislead a **wallet** backend are in scope. Truncating a sealed mmap after it is mapped is fatal external corruption (the next read can SIGBUS), not diff --git a/crates/rbitcoin-node/Cargo.toml b/crates/rbitcoin-node/Cargo.toml index aa2089979..2b714494c 100644 --- a/crates/rbitcoin-node/Cargo.toml +++ b/crates/rbitcoin-node/Cargo.toml @@ -26,7 +26,6 @@ bitcoin = { workspace = true } bitcoin_hashes = { workspace = true } serde_json = { workspace = true } axum = { workspace = true } -tower = { workspace = true } tower-http = { workspace = true } tokio = { workspace = true } mimalloc = { workspace = true } diff --git a/crates/rbitcoin-node/src/health.rs b/crates/rbitcoin-node/src/health.rs index 4af42c2a9..4e615e43d 100644 --- a/crates/rbitcoin-node/src/health.rs +++ b/crates/rbitcoin-node/src/health.rs @@ -20,13 +20,14 @@ use std::sync::atomic::{AtomicU8, Ordering}; use std::sync::{Arc, OnceLock}; use std::time::{Duration, SystemTime}; use tokio::net::TcpListener; -use tokio::task::JoinHandle; -use tower::limit::ConcurrencyLimitLayer; +use tokio::sync::Semaphore; +use tokio::task::{JoinError, JoinHandle}; use tower_http::limit::RequestBodyLimitLayer; use tower_http::timeout::TimeoutLayer; -/// Probes in flight at once. Probers send one request per period. -const MAX_CONCURRENT: usize = 16; +/// `/readyz` checks running at once. Probers send one request per period; +/// past this a probe answers 503 at once instead of queueing. +const READYZ_IN_FLIGHT: usize = 4; /// Per-request wall. Probe timeouts are usually 1s. const REQUEST_TIMEOUT: Duration = Duration::from_secs(5); /// Tip and scripthash index lag that `/readyz` still calls ready. @@ -244,7 +245,31 @@ pub(crate) async fn run_health( Ok(HealthHandle { task }) } +/// Router state: the node's status and the gates on blocking work. +#[derive(Clone)] +struct Health { + status: Arc, + readyz: Arc, + /// One render at a time, so a scrape never starts a second mempool fold + /// while one is running. + scrape: Arc, +} + +impl Health { + fn new(status: Arc) -> Self { + Self { + status, + readyz: Arc::new(Semaphore::new(READYZ_IN_FLIGHT)), + scrape: Arc::new(Semaphore::new(1)), + } + } +} + fn router(status: Arc, metrics: bool) -> Router { + router_for(Health::new(status), metrics) +} + +fn router_for(health: Health, metrics: bool) -> Router { let routes = Router::new() .route("/healthz", get(healthz)) .route("/readyz", get(readyz)); @@ -253,46 +278,72 @@ fn router(status: Arc, metrics: bool) -> Router { } else { routes }; + // Outer → inner: timeout → body (GET only, so none). The timeout covers + // the whole request. The blocking work is capped by `gated`, not by a + // concurrency layer: tower's queues without a bound, and axum builds one + // per route. routes - // Outer → inner: concurrency → body (GET only, so none) → timeout. + .layer(RequestBodyLimitLayer::new(0)) .layer(TimeoutLayer::with_status_code( StatusCode::REQUEST_TIMEOUT, REQUEST_TIMEOUT, )) - .layer(RequestBodyLimitLayer::new(0)) - .layer(ConcurrencyLimitLayer::new(MAX_CONCURRENT)) - .with_state(status) + .with_state(health) +} + +/// Run `f` on the blocking pool under a permit from `gate`, or `None` at +/// once when every permit is taken. +/// +/// The permit moves into the task. A blocking task cannot be cancelled, so +/// a request that times out drops only its wait, and the gate stays closed +/// until `f` returns: timed-out requests cannot stack chain or mempool +/// readers behind the cap. +async fn gated( + gate: &Arc, + f: impl FnOnce() -> T + Send + 'static, +) -> Option> { + let permit = Arc::clone(gate).try_acquire_owned().ok()?; + Some( + tokio::task::spawn_blocking(move || { + let _permit = permit; + let _g = BlockingRegion::enter(); + f() + }) + .await, + ) } async fn healthz() -> &'static str { "ok\n" } -async fn scrape(State(status): State>) -> Response { - let body = tokio::task::spawn_blocking(move || { - let _g = BlockingRegion::enter(); - metrics::render(&status) - }) - .await; - match body { - Ok(body) => ([(header::CONTENT_TYPE, metrics::CONTENT_TYPE)], body).into_response(), - Err(e) => (StatusCode::INTERNAL_SERVER_ERROR, format!("metrics: {e}\n")).into_response(), +async fn scrape(State(health): State) -> Response { + let status = Arc::clone(&health.status); + match gated(&health.scrape, move || metrics::render(&status)).await { + Some(Ok(body)) => ([(header::CONTENT_TYPE, metrics::CONTENT_TYPE)], body).into_response(), + Some(Err(e)) => { + (StatusCode::INTERNAL_SERVER_ERROR, format!("metrics: {e}\n")).into_response() + } + None => ( + StatusCode::SERVICE_UNAVAILABLE, + "metrics: a scrape is already running\n", + ) + .into_response(), } } -async fn readyz(State(status): State>) -> (StatusCode, String) { - let snap = tokio::task::spawn_blocking(move || { - let _g = BlockingRegion::enter(); - status.ready_snapshot() - }) - .await; - match snap.as_ref().map(readiness) { - Ok(Ok(())) => (StatusCode::OK, "ok\n".into()), - Ok(Err(reason)) => ( - StatusCode::SERVICE_UNAVAILABLE, - format!("not ready: {reason}\n"), - ), - Err(e) => (StatusCode::SERVICE_UNAVAILABLE, format!("not ready: {e}\n")), +async fn readyz(State(health): State) -> (StatusCode, String) { + let status = Arc::clone(&health.status); + match gated(&health.readyz, move || status.ready_snapshot()).await { + Some(Ok(snap)) => match readiness(&snap) { + Ok(()) => (StatusCode::OK, "ok\n".into()), + Err(reason) => ( + StatusCode::SERVICE_UNAVAILABLE, + format!("not ready: {reason}\n"), + ), + }, + Some(Err(e)) => (StatusCode::SERVICE_UNAVAILABLE, format!("not ready: {e}\n")), + None => (StatusCode::SERVICE_UNAVAILABLE, "not ready: busy\n".into()), } } @@ -413,6 +464,85 @@ mod tests { ); } + /// A blocking task cannot be cancelled, so a request that times out must + /// leave its permit with the task. Otherwise each timed-out request frees + /// a slot for another reader while the first still runs. + #[tokio::test(flavor = "multi_thread")] + async fn gated_work_keeps_its_permit_past_a_dropped_request() { + let gate = Arc::new(Semaphore::new(1)); + let (release, wait) = std::sync::mpsc::channel::<()>(); + let dropped = tokio::time::timeout( + Duration::from_millis(100), + gated(&gate, move || wait.recv().unwrap()), + ) + .await; + assert!(dropped.is_err(), "the request gave up while the work ran"); + assert_eq!(gate.available_permits(), 0); + let turned_away = tokio::time::timeout(Duration::from_secs(2), gated(&gate, || ())) + .await + .expect("a full gate answers at once"); + assert!(turned_away.is_none()); + + release.send(()).unwrap(); + tokio::time::timeout(Duration::from_secs(5), async { + while gate.available_permits() == 0 { + tokio::time::sleep(Duration::from_millis(10)).await; + } + }) + .await + .expect("the permit returns when the work ends"); + let ran = gated(&gate, || 7).await.expect("a free gate runs"); + assert_eq!(ran.unwrap(), 7); + } + + /// With its gate full, `/readyz` or `/metrics` answers 503 at once rather + /// than queue; with the gate free, it answers from the node. + #[tokio::test(flavor = "multi_thread")] + async fn a_full_gate_answers_503_without_queueing() { + let health = Health::new(NodeStatus::new(Network::Regtest, false)); + let addr = serve(router_for(health.clone(), true)).await; + + let held = Arc::clone(&health.readyz) + .acquire_many_owned(READYZ_IN_FLIGHT as u32) + .await + .unwrap(); + assert_eq!( + get_soon(addr, "/readyz").await, + "HTTP/1.1 503 Service Unavailable|not ready: busy\n" + ); + drop(held); + assert_eq!( + get_soon(addr, "/readyz").await, + "HTTP/1.1 503 Service Unavailable|not ready: opening\n" + ); + + let held = Arc::clone(&health.scrape).acquire_owned().await.unwrap(); + assert_eq!( + get_soon(addr, "/metrics").await, + "HTTP/1.1 503 Service Unavailable|metrics: a scrape is already running\n" + ); + drop(held); + let scraped = get_soon(addr, "/metrics").await; + assert!( + scraped.starts_with("HTTP/1.1 200 OK|# HELP rbitcoin_build_info"), + "{scraped}" + ); + } + + async fn serve(app: Router) -> SocketAddr { + let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); + let addr = listener.local_addr().unwrap(); + tokio::spawn(async move { axum::serve(listener, app).await.unwrap() }); + addr + } + + /// [`get`], failing rather than hanging if the answer waits for a permit. + async fn get_soon(addr: SocketAddr, path: &str) -> String { + tokio::time::timeout(Duration::from_secs(2), get(addr, path)) + .await + .expect("answered without queueing") + } + async fn get(addr: SocketAddr, path: &str) -> String { use tokio::io::{AsyncReadExt, AsyncWriteExt}; let mut s = tokio::net::TcpStream::connect(addr).await.unwrap(); diff --git a/docs/operator/operations.md b/docs/operator/operations.md index eae316509..9b5a4592b 100644 --- a/docs/operator/operations.md +++ b/docs/operator/operations.md @@ -361,13 +361,18 @@ hours (IBD). The health listener answers through all of it. | `GET /metrics` | Prometheus text format with `--metrics`; 404 without it | Other paths are 404 and other methods 405. Requests carry no body and no -auth. Concurrency, body, and timeout limits are always on. A non-loopback -bind logs a WARN: keep the port on loopback or a probe-only network. +auth, and time out after 5 s (408). `/readyz` runs at most 4 checks at +once and `/metrics` one render at a time; past that a request answers 503 +at once (`not ready: busy`, `metrics: a scrape is already running`) +instead of queueing, and a timed-out check keeps its slot until it +returns. A non-loopback bind logs a WARN: keep the port on loopback or a +probe-only network. `/readyz` reports the first failing check, in this order: | Reason | Meaning | |--------|---------| +| `busy` | 4 checks are already running, so this one did not start | | `opening` | Store open: schema migration, backfill, spend replay | | `starting` | P2P, mempool, and proxy bring-up | | `catch-up` | Initial block download | From 6404d8309ec6eb32e2ffeb93c56a73d670c9b0f3 Mon Sep 17 00:00:00 2001 From: Benjamen Keroack Date: Tue, 29 Sep 2026 20:25:43 -0600 Subject: [PATCH 10/11] changelog: one category for the health-endpoints fragment release_changelog_absorb_fragments takes a fragment's first non-empty line as its category and appends the rest under it, so the Changed section was filed under Added with its heading as text. The running totals are what /metrics exports, so they move into the Added metrics bullet. Co-Authored-By: Claude Opus 5.5 --- changelog.d/health-endpoints.md | 14 +++++--------- 1 file changed, 5 insertions(+), 9 deletions(-) diff --git a/changelog.d/health-endpoints.md b/changelog.d/health-endpoints.md index ced1ac62b..3e49396ee 100644 --- a/changelog.d/health-endpoints.md +++ b/changelog.d/health-endpoints.md @@ -12,12 +12,8 @@ Added migration or IBD. - **Prometheus metrics.** `--metrics` adds `GET /metrics` on the health listener. Gauges equal their RPC fields (`blocks`, `headers`, - `initialblockdownload`, connections, mempool size) and counters are the - `tip: perf` meters as running totals. - -Changed - -- **`tip: perf` meters are running totals.** Esplora and Electrum - requests, historical block serves, and mempool accepts and rejects count - up for the life of the process. The 5 s DEBUG line still prints the - change since the previous line. + `initialblockdownload`, connections, mempool size). Counters are the + `tip: perf` meters (Esplora and Electrum requests, historical block + serves, mempool accepts and rejects), which now count up for the life + of the process; the 5 s DEBUG line still prints the change since the + previous line. From faf64e9211ef03937db775b65660d02eacda8181 Mon Sep 17 00:00:00 2001 From: "rearden-grok[bot]" <317016512+rearden-grok[bot]@users.noreply.github.com> Date: Thu, 1 Oct 2026 10:57:25 -0700 Subject: [PATCH 11/11] docs: match /metrics to lifetime counters and byte RSS The operator table treated the Prometheus counters as the 5s tip: perf window and mixed units. Those counters are process totals (seconds, not microseconds). process_resident_memory_bytes is the same RSS reading in bytes; the log lines print integer MiB. A scrape also reads the chain. --- crates/rbitcoin-node/src/health/metrics.rs | 7 ++++--- docs/operator/operations.md | 23 ++++++++++++---------- 2 files changed, 17 insertions(+), 13 deletions(-) diff --git a/crates/rbitcoin-node/src/health/metrics.rs b/crates/rbitcoin-node/src/health/metrics.rs index 120cf447d..f6879beaa 100644 --- a/crates/rbitcoin-node/src/health/metrics.rs +++ b/crates/rbitcoin-node/src/health/metrics.rs @@ -10,8 +10,9 @@ use std::time::UNIX_EPOCH; pub(super) const CONTENT_TYPE: &str = "text/plain; version=0.0.4; charset=utf-8"; /// One scrape. Chain reads may touch the store and the mempool totals take -/// its lock, so call from the blocking pool. Cost: one peer snapshot and one -/// fold over the mempool (the same fold as `getmempoolinfo`). +/// its lock, so call from the blocking pool. Cost: those chain reads +/// (`best_header_height`, tip header, `in_ibd`, scripthash lag), one peer +/// snapshot, and one mempool fold (the same fold as `getmempoolinfo`). pub(super) fn render(status: &NodeStatus) -> String { let mut out = Exposition::default(); out.family( @@ -58,7 +59,7 @@ pub(super) fn render(status: &NodeStatus) -> String { if rss_kb > 0 { out.gauge( "process_resident_memory_bytes", - "Resident memory size in bytes (ibd: sizes rss=).", + "Resident memory size in bytes. ibd: sizes rss= is the same reading in MiB.", rss_kb * 1024, ); } diff --git a/docs/operator/operations.md b/docs/operator/operations.md index 9b5a4592b..6c765685a 100644 --- a/docs/operator/operations.md +++ b/docs/operator/operations.md @@ -401,11 +401,14 @@ A k8s `httpGet` probe connects to the pod IP, so in a pod use ### Metrics -`--metrics` (needs `--health-listen`) adds `GET /metrics`. Each gauge is a -value the node already publishes on RPC or a log line, under a name that -says which. A scrape reads a few atomics and one peer snapshot, and folds -the mempool once for `rbitcoin_mempool_bytes` (the same fold as -`getmempoolinfo`) on the blocking pool. +`--metrics` (needs `--health-listen`) adds `GET /metrics`. Gauges are values +the node already publishes on RPC or a log line, under a name that says +which. Counters are the process-lifetime totals behind the 5s DEBUG +`tip: perf` line; that line still prints only the change since its previous +sample. A scrape runs on the blocking pool: chain reads +(`best_header_height`, the tip header, `in_ibd`, and scripthash lag with +`--sh-index`), one peer snapshot, and one mempool fold for +`rbitcoin_mempool_bytes` (the same fold as `getmempoolinfo`). | Metric | Type | Equals | |--------|------|--------| @@ -420,11 +423,11 @@ the mempool once for `rbitcoin_mempool_bytes` (the same fold as | `rbitcoin_mempool_transactions` | gauge | `getmempoolinfo.size` | | `rbitcoin_mempool_bytes` | gauge | `getmempoolinfo.bytes` (virtual size) | | `rbitcoin_scripthash_lag_blocks` | gauge | `tip: accept sh_lag=` (with `--sh-index`) | -| `rbitcoin_esplora_requests_total` / `_request_seconds_total` | counter | `tip: perf esplora req=` / `avg_us` × `req` | -| `rbitcoin_electrum_requests_total` / `_request_seconds_total` | counter | `tip: perf electrum req=` / `avg_us` × `req` | -| `rbitcoin_block_serve_total` / `_bytes_total` | counter | `tip: perf serve n=` / `bytes=` | -| `rbitcoin_mempool_accepts_total` / `_rejects_total` | counter | `tip: perf accepts=` / `rejects=` | -| `process_resident_memory_bytes` | gauge | `ibd: sizes rss=` (Linux and macOS) | +| `rbitcoin_esplora_requests_total` / `_request_seconds_total` | counter | Lifetime sum of `tip: perf esplora req=`, and of that handler's wall time in seconds. The DEBUG line is the last ~5s window (`req=`, `avg_us` in microseconds) | +| `rbitcoin_electrum_requests_total` / `_request_seconds_total` | counter | Lifetime sum of `tip: perf electrum req=`, and of that handler's wall time in seconds. The DEBUG line is the last ~5s window (`req=`, `avg_us` in microseconds) | +| `rbitcoin_block_serve_total` / `_bytes_total` | counter | Lifetime sum of `tip: perf serve n=` / `bytes=`. That line is the last ~5s window | +| `rbitcoin_mempool_accepts_total` / `_rejects_total` | counter | Lifetime sum of `tip: perf accepts=` / `rejects=`. That line is the last ~5s window | +| `process_resident_memory_bytes` | gauge | Same RSS reading as `ibd: sizes rss=` and `tip: perf rss=`, in bytes (`rss_kb * 1024`). Those lines print integer MiB (`rss_kb / 1024`). Linux and macOS | | `process_start_time_seconds` | gauge | Unix time `rbitcoin-node` started | Hub gauges appear once P2P has started. `/metrics` exposes peer and mempool