diff --git a/Cargo.lock b/Cargo.lock index b58d267a..944211b4 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -49,15 +49,6 @@ version = "0.13.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "9e1b586273c5702936fe7b7d6896644d8be71e6314cfe09d3167c95f712589e8" -[[package]] -name = "base64-compat" -version = "1.0.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "5a8d4d2746f89841e49230dd26917df1876050f95abafafbe34f47cb534b88d7" -dependencies = [ - "byteorder", -] - [[package]] name = "bech32" version = "0.9.1" @@ -118,12 +109,6 @@ version = "1.3.2" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "bef38d45163c2f1dde094a7dfd33ccf595c92905c8f8f4fdc18d06fb1037718a" -[[package]] -name = "byteorder" -version = "1.4.3" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "14c189c53d098945499cdfa7ecc63567cf3886b3332b312a5b4585d8d3a6a610" - [[package]] name = "cc" version = "1.0.73" @@ -226,13 +211,13 @@ checksum = "4217ad341ebadf8d8e724e264f13e593e0648f5b3e94b3896a5df283be015ecc" [[package]] name = "jsonrpc" -version = "0.12.1" +version = "0.16.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "7f8423b78fc94d12ef1a4a9d13c348c9a78766dda0cc18817adf0faf77e670c8" +checksum = "34efde8d2422fb79ed56db1d3aea8fa5b583351d15a26770cdee2f88813dd702" dependencies = [ - "base64-compat", + "base64", + "minreq", "serde", - "serde_derive", "serde_json", ] @@ -308,6 +293,17 @@ dependencies = [ "adler", ] +[[package]] +name = "minreq" +version = "2.8.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "3de406eeb24aba36ed3829532fa01649129677186b44a49debec0ec574ca7da7" +dependencies = [ + "log", + "serde", + "serde_json", +] + [[package]] name = "object" version = "0.30.4" diff --git a/Cargo.toml b/Cargo.toml index fac64f57..d98e19e6 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -49,7 +49,7 @@ backtrace = "0.3" rusqlite = { version = "0.27", features = ["bundled", "unlock_notify"] } # To talk to bitcoind -jsonrpc = "0.12" +jsonrpc = { version = "0.16", features = ["minreq_http"], default-features = false } # Used for daemonization libc = { version = "0.2", optional = true } diff --git a/src/bitcoin/d/mod.rs b/src/bitcoin/d/mod.rs index 798694ff..9d3d670c 100644 --- a/src/bitcoin/d/mod.rs +++ b/src/bitcoin/d/mod.rs @@ -23,7 +23,8 @@ use std::{ use jsonrpc::{ arg, client::Client, - simple_http::{self, SimpleHttpTransport}, + minreq, + minreq_http::{self, MinreqHttpTransport}, }; use miniscript::{ @@ -74,16 +75,41 @@ impl BitcoindError { /// Is it a timeout of any kind? pub fn is_timeout(&self) -> bool { - match self { - BitcoindError::Server(jsonrpc::Error::Transport(ref e)) => { - match e.downcast_ref::() { - Some(simple_http::Error::Timeout) => true, - Some(simple_http::Error::SocketError(e)) => e.kind() == io::ErrorKind::TimedOut, - _ => false, - } + if let BitcoindError::Server(jsonrpc::Error::Transport(ref e)) = self { + if let Some(minreq_http::Error::Minreq(minreq::Error::IoError(e))) = + e.downcast_ref::() + { + return e.kind() == io::ErrorKind::TimedOut; } - _ => false, } + false + } + + /// Is it an error that can be recovered from? + pub fn is_transient(&self) -> bool { + if let BitcoindError::Server(jsonrpc::Error::Transport(ref e)) = self { + if let Some(ref e) = e.downcast_ref::() { + // Bitcoind is overloaded + if let minreq_http::Error::Http(minreq_http::HttpError { status_code, .. }) = e { + return status_code == &503; + } + // Bitcoind may have been restarted + return matches!(e, minreq_http::Error::Minreq(minreq::Error::IoError(_))); + } + } + false + } + + /// Is it an error that has to do with our credentials? + pub fn is_unauthorized(&self) -> bool { + if let BitcoindError::Server(jsonrpc::Error::Transport(ref e)) = self { + if let Some(minreq_http::Error::Http(minreq_http::HttpError { status_code, .. })) = + e.downcast_ref::() + { + return status_code == &402; + } + } + false } } @@ -131,8 +157,8 @@ impl From for BitcoindError { } } -impl From for BitcoindError { - fn from(e: simple_http::Error) -> Self { +impl From for BitcoindError { + fn from(e: minreq_http::Error) -> Self { jsonrpc::error::Error::Transport(Box::new(e)).into() } } @@ -204,19 +230,20 @@ impl BitcoinD { ) -> Result { let cookie_string = fs::read_to_string(&config.cookie_path).map_err(BitcoindError::CookieFile)?; + let node_url = format!("http://{}", config.addr); let watchonly_url = format!("http://{}/wallet/{}", config.addr, watchonly_wallet_path); // Create a dummy bitcoind with clients using a low timeout to sanity check the connection. let dummy_node_client = Client::with_transport( - SimpleHttpTransport::builder() - .url(&config.addr.to_string()) + MinreqHttpTransport::builder() + .url(&node_url) .map_err(BitcoindError::from)? .timeout(Duration::from_secs(3)) .cookie_auth(cookie_string.clone()) .build(), ); let sendonly_client = Client::with_transport( - SimpleHttpTransport::builder() + MinreqHttpTransport::builder() .url(&watchonly_url) .map_err(BitcoindError::from)? .timeout(Duration::from_secs(1)) @@ -224,7 +251,7 @@ impl BitcoinD { .build(), ); let dummy_wo_client = Client::with_transport( - SimpleHttpTransport::builder() + MinreqHttpTransport::builder() .url(&watchonly_url) .map_err(BitcoindError::from)? .timeout(Duration::from_secs(3)) @@ -242,15 +269,15 @@ impl BitcoinD { // Now the connection is checked, create the clients with an appropriate timeout. let node_client = Client::with_transport( - SimpleHttpTransport::builder() - .url(&config.addr.to_string()) + MinreqHttpTransport::builder() + .url(&node_url) .map_err(BitcoindError::from)? .timeout(Duration::from_secs(RPC_SOCKET_TIMEOUT)) .cookie_auth(cookie_string.clone()) .build(), ); let sendonly_client = Client::with_transport( - SimpleHttpTransport::builder() + MinreqHttpTransport::builder() .url(&watchonly_url) .map_err(BitcoindError::from)? .timeout(Duration::from_secs(1)) @@ -258,7 +285,7 @@ impl BitcoinD { .build(), ); let watchonly_client = Client::with_transport( - SimpleHttpTransport::builder() + MinreqHttpTransport::builder() .url(&watchonly_url) .map_err(BitcoindError::from)? .timeout(Duration::from_secs(RPC_SOCKET_TIMEOUT)) @@ -307,20 +334,23 @@ impl BitcoinD { Ok(res) => return Ok(res), Err(e) => { if e.is_warming_up() { + // Always retry when bitcoind is warming up, it'll be available eventually. + std::thread::sleep(Duration::from_secs(1)); error = Some(e) - } else if let BitcoindError::Server(jsonrpc::Error::Transport(ref err)) = e { - match err.downcast_ref::() { - Some(simple_http::Error::Timeout) - | Some(simple_http::Error::SocketError(_)) - | Some(simple_http::Error::HttpErrorCode(503)) => { - if i <= self.retries { - std::thread::sleep(Duration::from_secs(1)); - log::debug!("Retrying RPC request to bitcoind: attempt #{}", i); - } - error = Some(e); - } - _ => return Err(e), + } else if e.is_unauthorized() { + // FIXME: it should be trivial for us to cache the cookie path and simply + // refresh the credentials when this happens. Unfortunately this means + // making the BitcoinD struct mutable... + log::error!("Denied access to bitcoind. Most likely bitcoind was restarted from under us and the cookie changed."); + return Err(e); + } else if e.is_transient() { + // If we start hitting transient errors retry requests for a limited time. + log::warn!("Transient error when sending request to bitcoind: {}", e); + if i <= self.retries { + std::thread::sleep(Duration::from_secs(1)); + log::debug!("Retrying RPC request to bitcoind: attempt #{}", i); } + error = Some(e); } else { return Err(e); } diff --git a/src/lib.rs b/src/lib.rs index abb52019..d0e61392 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -513,10 +513,12 @@ mod tests { "HTTP/1.1 200\n\r\n{\"jsonrpc\":\"2.0\",\"id\":1,\"result\":[]}\n".as_bytes(); // Read the first echo, respond to it - let (mut stream, _) = server.accept().unwrap(); - read_til_json_end(&mut stream); - stream.write_all(echo_resp).unwrap(); - stream.flush().unwrap(); + { + let (mut stream, _) = server.accept().unwrap(); + read_til_json_end(&mut stream); + stream.write_all(echo_resp).unwrap(); + stream.flush().unwrap(); + } // Read the second echo, respond to it let (mut stream, _) = server.accept().unwrap(); @@ -549,23 +551,27 @@ mod tests { // Send them responses for the calls involved when creating a fresh wallet fn complete_wallet_creation(server: &net::TcpListener) { - let net_resp = - ["HTTP/1.1 200\n\r\n{\"jsonrpc\":\"2.0\",\"id\":1,\"result\":[]}\n".as_bytes()] - .concat(); - let (mut stream, _) = server.accept().unwrap(); - read_til_json_end(&mut stream); - stream.write_all(&net_resp).unwrap(); - stream.flush().unwrap(); + { + let net_resp = + ["HTTP/1.1 200\n\r\n{\"jsonrpc\":\"2.0\",\"id\":1,\"result\":[]}\n".as_bytes()] + .concat(); + let (mut stream, _) = server.accept().unwrap(); + read_til_json_end(&mut stream); + stream.write_all(&net_resp).unwrap(); + stream.flush().unwrap(); + } - let net_resp = [ + { + let net_resp = [ "HTTP/1.1 200\n\r\n{\"jsonrpc\":\"2.0\",\"id\":1,\"result\":{\"name\":\"dummy\"}}\n" .as_bytes(), - ] - .concat(); - let (mut stream, _) = server.accept().unwrap(); - read_til_json_end(&mut stream); - stream.write_all(&net_resp).unwrap(); - stream.flush().unwrap(); + ] + .concat(); + let (mut stream, _) = server.accept().unwrap(); + read_til_json_end(&mut stream); + stream.write_all(&net_resp).unwrap(); + stream.flush().unwrap(); + } let net_resp = [ "HTTP/1.1 200\n\r\n{\"jsonrpc\":\"2.0\",\"id\":1,\"result\":[{\"success\":true}]}\n" @@ -580,12 +586,14 @@ mod tests { // Send them a dummy result to loadwallet. fn complete_wallet_loading(server: &net::TcpListener) { - let listwallets_resp = - "HTTP/1.1 200\n\r\n{\"jsonrpc\":\"2.0\",\"id\":1,\"result\":[]}\n".as_bytes(); - let (mut stream, _) = server.accept().unwrap(); - read_til_json_end(&mut stream); - stream.write_all(listwallets_resp).unwrap(); - stream.flush().unwrap(); + { + let listwallets_resp = + "HTTP/1.1 200\n\r\n{\"jsonrpc\":\"2.0\",\"id\":1,\"result\":[]}\n".as_bytes(); + let (mut stream, _) = server.accept().unwrap(); + read_til_json_end(&mut stream); + stream.write_all(listwallets_resp).unwrap(); + stream.flush().unwrap(); + } let loadwallet_resp = "HTTP/1.1 200\n\r\n{\"jsonrpc\":\"2.0\",\"id\":1,\"result\":{\"name\":\"dummy\"}}\n" diff --git a/tests/test_misc.py b/tests/test_misc.py index 1d420da1..92394c1a 100644 --- a/tests/test_misc.py +++ b/tests/test_misc.py @@ -4,6 +4,8 @@ from fixtures import * from test_framework.serializations import PSBT from test_framework.utils import wait_for, RpcError, OLD_LIANAD_PATH, LIANAD_PATH +from threading import Thread + def receive_and_send(lianad, bitcoind): n_coins = len(lianad.rpc.listcoins()["coins"]) @@ -16,11 +18,13 @@ def receive_and_send(lianad, bitcoind): wait_for(lambda: len(lianad.rpc.listcoins()["coins"]) == n_coins + 3) # Create a spend that will create a change output, sign and broadcast it. - outpoints = [next( - c["outpoint"] - for c in lianad.rpc.listcoins()["coins"] - if c["spend_info"] is None - )] + outpoints = [ + next( + c["outpoint"] + for c in lianad.rpc.listcoins()["coins"] + if c["spend_info"] is None + ) + ] destinations = { bitcoind.rpc.getnewaddress(): 200_000, } @@ -218,3 +222,30 @@ def test_migration(lianad_multisig, bitcoind): receive_and_send(lianad, bitcoind) spend_txs = lianad.rpc.listspendtxs()["spend_txs"] assert len(spend_txs) == 2 and all(s["updated_at"] is not None for s in spend_txs) + + +def test_retry_on_workqueue_exceeded(lianad, bitcoind): + """Make sure we retry requests to bitcoind if it is temporarily overloaded.""" + # Start by reducing the work queue to a single slot. Note we need to stop lianad + # as we don't support yet restarting a bitcoind due to the cookie file getting + # overwritten. + lianad.stop() + bitcoind.cmd_line += ["-rpcworkqueue=1", "-rpcthreads=1"] + bitcoind.stop() + bitcoind.start() + lianad.start() + + # Stuck the bitcoind RPC server working queue with a command that takes 5 seconds + # to be replied to, and make lianad send it a request. Make sure we detect this is + # a transient HTTP 503 error and we retry the request. Once the 5 seconds are past + # our request succeeds and we get the reply to the lianad RPC command. + t = Thread(target=bitcoind.rpc.waitfornewblock, args=(5_000,)) + t.start() + lianad.rpc.getinfo() + lianad.wait_for_logs( + [ + "Transient error when sending request to bitcoind.*(status: 503, body: Work queue depth exceeded)", + "Retrying RPC request to bitcoind", + ] + ) + t.join()