Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 2 additions & 0 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

2 changes: 1 addition & 1 deletion OPERATOR.md
Original file line number Diff line number Diff line change
Expand Up @@ -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) |
Expand Down
6 changes: 6 additions & 0 deletions SECURITY.md
Original file line number Diff line number Diff line change
Expand Up @@ -64,6 +64,12 @@ 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. 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
Expand Down
19 changes: 19 additions & 0 deletions changelog.d/health-endpoints.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,19 @@
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, 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). 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.
4 changes: 2 additions & 2 deletions crates/rbitcoin-electrum/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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::{
Expand Down
35 changes: 13 additions & 22 deletions crates/rbitcoin-electrum/src/server.rs
Original file line number Diff line number Diff line change
Expand Up @@ -6,14 +6,15 @@
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};
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};
Expand All @@ -35,32 +36,22 @@ pub fn parse_electrum_request_line(line: &str) -> Option<Value> {
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).
Expand Down
3 changes: 2 additions & 1 deletion crates/rbitcoin-esplora/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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,
};
36 changes: 14 additions & 22 deletions crates/rbitcoin-esplora/src/server.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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;
Expand All @@ -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 {
Expand All @@ -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
Expand Down Expand Up @@ -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;

Expand Down
96 changes: 84 additions & 12 deletions crates/rbitcoin-net/src/chain.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<BlockHash>, Vec<BlockHash>) = {
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<BlockHash> = headers.prevs().collect();
for hash in headers.hashes() {
for hash in hashes {
if covered.contains(&hash) {
continue;
}
Expand All @@ -871,10 +876,14 @@ impl ChainHub {
}
}
{
let held = self.held_bodies.read().unwrap();
let parents: HashSet<BlockHash> =
held.blocks().map(|b| b.header.prev_blockhash).collect();
for hash in held.keys() {
let (parents, hashes): (HashSet<BlockHash>, Vec<BlockHash>) = {
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;
}
Expand Down Expand Up @@ -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
}
Expand Down Expand Up @@ -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();
Expand Down
11 changes: 8 additions & 3 deletions crates/rbitcoin-net/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,7 @@ mod netgroup;
mod peer;
mod peer_dos;
mod peers;
mod perf_meter;
mod reactor;
mod seeds;
mod serve_perf;
Expand Down Expand Up @@ -59,9 +60,11 @@ 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 use perf_meter::RequestMeter;
pub(crate) use rbitcoin_mempool::MempoolGraphStats;
pub use rbitcoin_mempool::{AcceptError, Selected};
pub use reactor::BlockingRegion;
Expand All @@ -70,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::{
Expand Down
7 changes: 7 additions & 0 deletions crates/rbitcoin-net/src/peers.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down
Loading
Loading