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
32 changes: 27 additions & 5 deletions src/bin/electrs.rs
Original file line number Diff line number Diff line change
Expand Up @@ -60,7 +60,13 @@ fn run_server(config: Arc<Config>, salt_rwlock: Arc<RwLock<String>>) -> Result<(
info!("starting electrs");

if let Some(zmq_addr) = config.zmq_addr.as_ref() {
zmq::start(&format!("tcp://{zmq_addr}"), block_hash_notify);
if let Err(e) = zmq::start(&format!("tcp://{zmq_addr}"), &block_hash_notify, &metrics) {
error!(
"ZMQ notifications disabled, falling back to polling: zmq_addr='{}' err='{}'",
zmq_addr,
e.display_chain()
);
}
}

info!("connecting to daemon at {}", config.daemon_rpc_addr);
Expand All @@ -86,7 +92,15 @@ fn run_server(config: Arc<Config>, salt_rwlock: Arc<RwLock<String>>) -> Result<(
);
info!("starting initial sync");
let mut tip = indexer.update(&daemon)?;
info!("initial sync complete, tip at {}", tip);
let daemon_tip = daemon.getbestblockhash()?;
if tip == daemon_tip {
info!("initial sync complete, tip at {}", tip);
} else {
info!(
"initial sync incomplete, serving a partial index and retrying outstanding blocks in the background tip='{}' daemon_tip='{}'",
tip, daemon_tip
);
}

let chain = Arc::new(ChainQuery::new(
Arc::clone(&store),
Expand All @@ -108,10 +122,18 @@ fn run_server(config: Arc<Config>, salt_rwlock: Arc<RwLock<String>>) -> Result<(
Arc::clone(&config),
)));

// A mempool update is retried whenever the tip moves, so index and try again. An update
// that cannot advance the tip never will here, and no listener is bound yet.
while !Mempool::update(&mempool, &daemon, &tip)? {
// Mempool syncing was aborted because the chain tip moved;
// Index the new block(s) and try again.
tip = indexer.update(&daemon)?;
let new_tip = indexer.update(&daemon)?;
if new_tip == tip {
warn!(
"index could not advance, starting up with a partial index tip='{}'",
tip
);
break;
}
tip = new_tip;
}

#[cfg(feature = "liquid")]
Expand Down
97 changes: 75 additions & 22 deletions src/daemon.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1270,7 +1270,14 @@ impl Daemon {
pub fn getblock(&self, blockhash: &BlockHash) -> Result<Block> {
let block =
block_from_value(self.request("getblock", json!([blockhash, /*verbose=*/ false]))?)?;
assert_eq!(block.block_hash(), *blockhash);
let returned = block.block_hash();
if returned != *blockhash {
bail!(
"bitcoind returned wrong block for getblock requested_hash='{}' returned_hash='{}'",
blockhash,
returned
);
}
Ok(block)
}

Expand All @@ -1295,23 +1302,46 @@ impl Daemon {
Err(e) => {
let err_msg = format!("{e:?}");
if err_msg.contains("Block not found on disk")
|| err_msg.contains("Block not available")
|| err_msg.contains("Block not available")
{
// There is a small chance the node returns the header but didn't finish to index the block
log::warn!("getblocks failing with: {e:?} trying {attempts} more time")
if attempts == 0 {
bail!(
"failed to get blocks from bitcoind attempts='0' err='{}'",
err_msg
);
}
log::warn!(
"getblocks failed, retrying err='{}' attempts_left='{}'",
err_msg,
attempts
);
} else {
panic!("failed to get blocks from bitcoind: {}", err_msg);
bail!("failed to get blocks from bitcoind err='{}'", err_msg);
}
}
}
if attempts == 0 {
panic!("failed to get blocks from bitcoind")
}
std::thread::sleep(RETRY_WAIT_DURATION);
};
let mut blocks = vec![];
for value in values {
blocks.push(block_from_value(value)?);
if values.len() != blockhashes.len() {
bail!(
"bitcoind returned wrong number of blocks requested='{}' returned='{}'",
blockhashes.len(),
values.len()
);
}
let mut blocks = Vec::with_capacity(values.len());
for (value, requested) in values.into_iter().zip(blockhashes.iter()) {
let block = block_from_value(value)?;
let returned = block.block_hash();
if returned != *requested {
bail!(
"bitcoind returned wrong block for getblocks requested_hash='{}' returned_hash='{}'",
requested,
returned
);
}
blocks.push(block);
}
Ok(blocks)
}
Expand Down Expand Up @@ -1466,20 +1496,31 @@ impl Daemon {

#[trace]
fn get_all_headers(&self, tip: &BlockHash) -> Result<Vec<BlockHeader>> {
const CHUNK_SIZE: usize = 100_000;

let info: Value = self.request("getblockheader", json!([tip]))?;
let tip_height = info
.get("height")
.expect("missing height")
.as_u64()
.expect("non-numeric height") as usize;
let all_heights: Vec<usize> = (0..=tip_height).collect();
let chunk_size = 100_000;
let mut result = vec![];
for heights in all_heights.chunks(chunk_size) {
let mut headers = self.getblockheaders(&heights)?;
assert!(headers.len() == heights.len());
.and_then(|v| v.as_u64())
.ok_or_else(|| format!("bitcoind returned malformed getblockheader info='{:?}'", info))?;
let tip_height = usize::try_from(tip_height).map_err(|_| {
format!("bitcoind returned out-of-range tip_height='{tip_height}'")
})?;

result.append(&mut headers);
// Chunked so the height allocation stays bounded whatever height bitcoind reports.
let mut result = vec![];
for start in (0..=tip_height).step_by(CHUNK_SIZE) {
let end = start.saturating_add(CHUNK_SIZE - 1).min(tip_height);
let heights: Vec<usize> = (start..=end).collect();
let headers = self.getblockheaders(&heights)?;
if headers.len() != heights.len() {
bail!(
"bitcoind returned wrong number of headers requested='{}' returned='{}'",
heights.len(),
headers.len()
);
}
result.extend(headers);

debug!(
"downloaded {}/{} block headers ({:.0}%)",
Expand All @@ -1491,10 +1532,22 @@ impl Daemon {

let mut blockhash = *DEFAULT_BLOCKHASH;
for header in &result {
assert_eq!(header.prev_blockhash, blockhash);
if header.prev_blockhash != blockhash {
bail!(
"bitcoind returned non-contiguous headers expected_prev='{}' got_prev='{}'",
blockhash,
header.prev_blockhash
);
}
blockhash = header.block_hash();
}
assert_eq!(blockhash, *tip);
if blockhash != *tip {
bail!(
"bitcoind headers do not end at requested tip expected_tip='{}' got_tip='{}'",
tip,
blockhash
);
}
Ok(result)
}

Expand Down
19 changes: 12 additions & 7 deletions src/electrum/server.rs
Original file line number Diff line number Diff line change
Expand Up @@ -585,11 +585,7 @@ impl Connection {

#[trace(method = %method)]
fn handle_command(&mut self, method: &str, params: &[Value], id: &Value) -> Result<Value> {
let timer = self
.stats
.latency
.with_label_values(&[method])
.start_timer();
let started = Instant::now();

let result = match method {
"blockchain.block.header" => self.blockchain_block_header(&params),
Expand Down Expand Up @@ -626,15 +622,24 @@ impl Connection {
"server.add_peer" => self.server_add_peer(&params),

&_ => {
warn!("rpc #{} unknown method {} {:?}", id, method, params);
debug!("rpc #{} unknown method {} {:?}", id, method, params);
self.stats
.latency
.with_label_values(&["unknown"])
.observe(started.elapsed().as_secs_f64());
return Ok(json_rpc_error(
format!("unknown method {}", method),
Some(id),
JsonRpcV2Error::MethodNotFound,
));
}
};
timer.observe_duration();

self.stats
.latency
.with_label_values(&[method])
.observe(started.elapsed().as_secs_f64());

Ok(match result {
Ok(result) => json!({"jsonrpc": "2.0", "id": id, "result": result}),
Err(e) => {
Expand Down
2 changes: 1 addition & 1 deletion src/new_index/db.rs
Original file line number Diff line number Diff line change
Expand Up @@ -388,7 +388,7 @@ impl DB {
}

#[cfg(test)]
fn open_test(path: &Path) -> DB {
pub(crate) fn open_test(path: &Path) -> DB {
let mut db_opts = rocksdb::Options::default();
db_opts.create_if_missing(true);
db_opts.set_prefix_extractor(rocksdb::SliceTransform::create_fixed_prefix(33));
Expand Down
Loading
Loading