fix(pusher): improve welist response logging, error handling, and add request timeout
- Always log HTTP status code and response body from welist (info level) - Log send_stats_report errors at call site instead of silently discarding - Add 10s timeout to reqwest client to prevent indefinite hangs - Apply clippy fixes (is_empty, if-let chains, dead_code, etc.)
This commit is contained in:
@@ -6,7 +6,6 @@ use bitcoincore_rpc::{Auth, Client, Error, RpcApi, bitcoin};
|
||||
use bitcoincore_rpc_json::GetBlockchainInfoResult;
|
||||
|
||||
use byteorder::{LittleEndian, ReadBytesExt};
|
||||
use hex;
|
||||
use log::{debug, error, info, trace, warn};
|
||||
use serde::Deserialize;
|
||||
use serde::Serialize;
|
||||
@@ -23,10 +22,8 @@ use zmq::{Context, DEALER, DONTWAIT, Socket};
|
||||
use bal_server::db::open_db;
|
||||
use bal_server::validation::is_valid_welist_url;
|
||||
use base64::{Engine as _, engine::general_purpose};
|
||||
use openssl::hash::MessageDigest;
|
||||
use openssl::pkey::PKey;
|
||||
use openssl::sign::Signer;
|
||||
use openssl::sign::Verifier;
|
||||
use reqwest::Client as rClient;
|
||||
use std::fs;
|
||||
use std::time::Instant;
|
||||
@@ -143,7 +140,7 @@ fn get_network_params_default(network: Network) -> NetworkParams {
|
||||
}
|
||||
|
||||
fn get_cookie_filename(network: &NetworkParams) -> Result<String, Box<dyn StdError>> {
|
||||
if network.cookie_file != "" {
|
||||
if !network.cookie_file.is_empty() {
|
||||
Ok(network.cookie_file.clone())
|
||||
} else {
|
||||
match env::var_os("HOME") {
|
||||
@@ -161,12 +158,12 @@ fn get_cookie_filename(network: &NetworkParams) -> Result<String, Box<dyn StdErr
|
||||
}
|
||||
}
|
||||
fn get_client_from_username(
|
||||
url: &String,
|
||||
url: &str,
|
||||
network: &NetworkParams,
|
||||
) -> Result<(Client, GetBlockchainInfoResult), Box<dyn StdError>> {
|
||||
if network.rpc_user != "" {
|
||||
if !network.rpc_user.is_empty() {
|
||||
match Client::new(
|
||||
&url[..],
|
||||
url,
|
||||
Auth::UserPass(network.rpc_user.to_string(), network.rpc_pass.to_string()),
|
||||
) {
|
||||
Ok(client) => match client.get_blockchain_info() {
|
||||
@@ -180,30 +177,30 @@ fn get_client_from_username(
|
||||
}
|
||||
}
|
||||
fn get_client_from_cookie(
|
||||
url: &String,
|
||||
url: &str,
|
||||
network: &NetworkParams,
|
||||
) -> Result<(Client, GetBlockchainInfoResult), Box<dyn StdError>> {
|
||||
match get_cookie_filename(network) {
|
||||
Ok(cookie) => match Client::new(&url[..], Auth::CookieFile(cookie.into())) {
|
||||
Ok(cookie) => match Client::new(url, Auth::CookieFile(cookie.into())) {
|
||||
Ok(client) => match client.get_blockchain_info() {
|
||||
Ok(bcinfo) => Ok((client, bcinfo)),
|
||||
Err(err) => Err(err.into()),
|
||||
},
|
||||
Err(err) => Err(err.into()),
|
||||
},
|
||||
Err(err) => Err(err.into()),
|
||||
Err(err) => Err(err),
|
||||
}
|
||||
}
|
||||
fn get_client(
|
||||
network: &NetworkParams,
|
||||
) -> Result<(Client, GetBlockchainInfoResult), Box<dyn StdError>> {
|
||||
let url = format!("{}:{}/", network.host, &network.port);
|
||||
let url = format!("{}:{}/", network.host, network.port);
|
||||
debug!("trying to connect to bitcoin daemon:{url}");
|
||||
match get_client_from_username(&url, network) {
|
||||
Ok(client) => Ok(client),
|
||||
Err(_) => match get_client_from_cookie(&url, &network) {
|
||||
Err(_) => match get_client_from_cookie(&url, network) {
|
||||
Ok(client) => Ok(client),
|
||||
Err(err) => Err(err.into()),
|
||||
Err(err) => Err(err),
|
||||
},
|
||||
}
|
||||
}
|
||||
@@ -265,7 +262,7 @@ async fn main_result(cfg: &MyConfig, network_params: &NetworkParams) -> Result<(
|
||||
let mut invalid_txs: std::collections::HashMap<String, String> = HashMap::new();
|
||||
for row_result in match query_tx.bind::<&[(_, Value)]>(
|
||||
&[
|
||||
(":locktime_threshold", (LOCKTIME_THRESHOLD as i64).into()),
|
||||
(":locktime_threshold", LOCKTIME_THRESHOLD.into()),
|
||||
(":bestblock_time", (average_time as i64).into()),
|
||||
(":bestblock_height", (bcinfo.blocks as i64).into()),
|
||||
(":network", network_params.db_field.clone().into()),
|
||||
@@ -331,7 +328,9 @@ async fn main_result(cfg: &MyConfig, network_params: &NetworkParams) -> Result<(
|
||||
stmt.bind((2, Value::String(txid.clone()))).unwrap();
|
||||
let _ = stmt.next();
|
||||
}
|
||||
let _ = send_stats_report(cfg, bcinfo).await;
|
||||
if let Err(e) = send_stats_report(cfg, bcinfo).await {
|
||||
error!("send_stats_report failed: {}", e);
|
||||
}
|
||||
let _ = calculate_stats(&db, network_params.db_field.clone()).await;
|
||||
}
|
||||
Err(erx) => {
|
||||
@@ -437,7 +436,10 @@ async fn send_stats_report(
|
||||
);
|
||||
return Ok(());
|
||||
}
|
||||
let client = rClient::new();
|
||||
let client = rClient::builder()
|
||||
.timeout(Duration::from_secs(10))
|
||||
.build()
|
||||
.unwrap_or_else(|_| rClient::new());
|
||||
let url = format!("{}/ping", welist_url);
|
||||
debug!("welist url: {}", url);
|
||||
let chain = bcinfo.chain.to_string().to_lowercase();
|
||||
@@ -446,7 +448,7 @@ async fn send_stats_report(
|
||||
cfg.url, chain, bcinfo.blocks, bcinfo.median_time, bcinfo.best_block_hash
|
||||
);
|
||||
trace!("message to be sent: {}", message);
|
||||
let sign = sign_message(cfg.ssl_key_path.as_str(), &message.as_str());
|
||||
let sign = sign_message(cfg.ssl_key_path.as_str(), message.as_str());
|
||||
let response = client
|
||||
.post(url)
|
||||
.header("User-Agent", format!("bal-pusher/{}", VERSION))
|
||||
@@ -461,16 +463,12 @@ async fn send_stats_report(
|
||||
}))
|
||||
.send()
|
||||
.await?;
|
||||
if !response.status().is_success() {
|
||||
warn!(
|
||||
"Non-success response: {} {}",
|
||||
response.status(),
|
||||
response.status().canonical_reason().unwrap_or("")
|
||||
);
|
||||
}
|
||||
|
||||
let body = &(response.text().await?);
|
||||
info!("Report to welist({})\tSent: {}", welist_url, body);
|
||||
let status = response.status();
|
||||
let body = response.text().await?;
|
||||
info!(
|
||||
"Report to welist({}) status={} body={}",
|
||||
welist_url, status, body
|
||||
);
|
||||
} else {
|
||||
debug!("Not sending stats");
|
||||
}
|
||||
@@ -484,9 +482,7 @@ fn sign_message(private_key_path: &str, message: &str) -> String {
|
||||
|
||||
let signature = signer.sign_oneshot_to_vec(message.as_bytes()).unwrap();
|
||||
|
||||
let signature_b64 = general_purpose::STANDARD.encode(&signature);
|
||||
|
||||
signature_b64
|
||||
general_purpose::STANDARD.encode(&signature)
|
||||
}
|
||||
|
||||
fn parse_env(cfg: &mut MyConfig) {
|
||||
@@ -505,71 +501,45 @@ fn parse_env_netconfig(cfg_lock: &mut MyConfig, chain: &str) -> NetworkParams {
|
||||
"testnet4" => &mut cfg_lock.testnet4,
|
||||
&_ => &mut cfg_lock.mainnet,
|
||||
};
|
||||
match env::var(format!("BAL_PUSHER_{}_HOST", chain.to_uppercase())) {
|
||||
Ok(value) => {
|
||||
cfg.host = value;
|
||||
}
|
||||
Err(_) => {}
|
||||
if let Ok(value) = env::var(format!("BAL_PUSHER_{}_HOST", chain.to_uppercase())) {
|
||||
cfg.host = value;
|
||||
}
|
||||
match env::var(format!("BAL_PUSHER_{}_PORT", chain.to_uppercase())) {
|
||||
Ok(value) => match value.parse::<u64>() {
|
||||
Ok(value) => match u16::try_from(value) {
|
||||
Ok(port) => cfg.port = port,
|
||||
Err(e) => {
|
||||
error!(
|
||||
"Port value {} exceeds u16 range for chain {}: {}",
|
||||
value, chain, e
|
||||
);
|
||||
}
|
||||
},
|
||||
Err(_) => {}
|
||||
},
|
||||
Err(_) => {}
|
||||
}
|
||||
match env::var(format!("BAL_PUSHER_{}_DIR_PATH", chain.to_uppercase())) {
|
||||
Ok(value) => {
|
||||
cfg.dir_path = value;
|
||||
if let Ok(value) = env::var(format!("BAL_PUSHER_{}_PORT", chain.to_uppercase()))
|
||||
&& let Ok(port_num) = value.parse::<u64>()
|
||||
{
|
||||
if let Ok(port) = u16::try_from(port_num) {
|
||||
cfg.port = port;
|
||||
} else {
|
||||
error!(
|
||||
"Port value {} exceeds u16 range for chain {}",
|
||||
port_num, chain
|
||||
);
|
||||
}
|
||||
Err(_) => {}
|
||||
}
|
||||
match env::var(format!("BAL_PUSHER_{}_DB_FIELD", chain.to_uppercase())) {
|
||||
Ok(value) => {
|
||||
cfg.db_field = value;
|
||||
}
|
||||
Err(_) => {}
|
||||
if let Ok(value) = env::var(format!("BAL_PUSHER_{}_DIR_PATH", chain.to_uppercase())) {
|
||||
cfg.dir_path = value;
|
||||
}
|
||||
match env::var(format!("BAL_PUSHER_{}_COOKIE_FILE", chain.to_uppercase())) {
|
||||
Ok(value) => {
|
||||
cfg.cookie_file = value;
|
||||
}
|
||||
Err(_) => {}
|
||||
if let Ok(value) = env::var(format!("BAL_PUSHER_{}_DB_FIELD", chain.to_uppercase())) {
|
||||
cfg.db_field = value;
|
||||
}
|
||||
match env::var(format!("BAL_PUSHER_{}_RPC_USER", chain.to_uppercase())) {
|
||||
Ok(value) => {
|
||||
cfg.rpc_user = value;
|
||||
}
|
||||
Err(_) => {}
|
||||
if let Ok(value) = env::var(format!("BAL_PUSHER_{}_COOKIE_FILE", chain.to_uppercase())) {
|
||||
cfg.cookie_file = value;
|
||||
}
|
||||
match env::var(format!("BAL_PUSHER_{}_RPC_PASSWORD", chain.to_uppercase())) {
|
||||
Ok(value) => {
|
||||
cfg.rpc_pass = value;
|
||||
}
|
||||
Err(_) => {}
|
||||
if let Ok(value) = env::var(format!("BAL_PUSHER_{}_RPC_USER", chain.to_uppercase())) {
|
||||
cfg.rpc_user = value;
|
||||
}
|
||||
println!(
|
||||
"{}",
|
||||
format!("BAL_PUSHER_{}_ZMQ_HASHBLOCK", chain.to_uppercase())
|
||||
);
|
||||
match env::var(format!("BAL_PUSHER_{}_ZMQ_HASHBLOCK", chain.to_uppercase())) {
|
||||
Ok(value) => {
|
||||
println!("value:{}", value);
|
||||
cfg.zmq_listener = value;
|
||||
}
|
||||
Err(_) => {}
|
||||
if let Ok(value) = env::var(format!("BAL_PUSHER_{}_RPC_PASSWORD", chain.to_uppercase())) {
|
||||
cfg.rpc_pass = value;
|
||||
}
|
||||
println!("BAL_PUSHER_{}_ZMQ_HASHBLOCK", chain.to_uppercase());
|
||||
if let Ok(value) = env::var(format!("BAL_PUSHER_{}_ZMQ_HASHBLOCK", chain.to_uppercase())) {
|
||||
println!("value:{}", value);
|
||||
cfg.zmq_listener = value;
|
||||
}
|
||||
cfg.clone()
|
||||
}
|
||||
|
||||
#[allow(dead_code)]
|
||||
fn check_zmq_connection(endpoint: &str) -> bool {
|
||||
trace!("check zmq connection");
|
||||
let context = Context::new();
|
||||
@@ -587,6 +557,7 @@ fn check_zmq_connection(endpoint: &str) -> bool {
|
||||
}
|
||||
|
||||
// Add this struct to monitor connection health
|
||||
#[allow(dead_code)]
|
||||
struct ConnectionMonitor {
|
||||
last_message_time: Instant,
|
||||
timeout: Duration,
|
||||
@@ -594,6 +565,7 @@ struct ConnectionMonitor {
|
||||
max_consecutive_timeouts: u32,
|
||||
}
|
||||
|
||||
#[allow(dead_code)]
|
||||
impl ConnectionMonitor {
|
||||
fn new(timeout_secs: u64, max_timeouts: u32) -> Self {
|
||||
Self {
|
||||
@@ -631,6 +603,7 @@ impl ConnectionMonitor {
|
||||
}
|
||||
}
|
||||
|
||||
#[allow(dead_code)]
|
||||
enum ConnectionStatus {
|
||||
Healthy,
|
||||
Warning(Duration),
|
||||
@@ -642,7 +615,6 @@ async fn main() -> std::io::Result<()> {
|
||||
env_logger::init();
|
||||
let mut cfg = MyConfig::default();
|
||||
|
||||
let dbfile = env::var("BAL_PUSHER_DB_FILE").unwrap();
|
||||
parse_env(&mut cfg);
|
||||
let mut args = std::env::args();
|
||||
let _exe_name = args.next().unwrap();
|
||||
@@ -686,9 +658,6 @@ async fn main() -> std::io::Result<()> {
|
||||
|
||||
let _ = main_result(&cfg, network_params).await;
|
||||
info!("waiting new blocks..");
|
||||
let mut last_seq: Vec<u8> = [0; 4].to_vec();
|
||||
let mut counter = 0;
|
||||
let max = 100;
|
||||
socket.set_rcvtimeo(5000).unwrap(); // 5 seconds timeout
|
||||
loop {
|
||||
let message = match socket.recv_multipart(0) {
|
||||
@@ -700,8 +669,6 @@ async fn main() -> std::io::Result<()> {
|
||||
};
|
||||
let topic = message[0].clone();
|
||||
let body = message[1].clone();
|
||||
let seq = message[2].clone();
|
||||
last_seq = seq;
|
||||
debug!(
|
||||
"ZMQ:GET TOPIC: {}",
|
||||
String::from_utf8(topic.clone()).expect("invalid topic")
|
||||
@@ -714,7 +681,8 @@ async fn main() -> std::io::Result<()> {
|
||||
thread::sleep(Duration::from_millis(100)); // Sleep for 100ms
|
||||
}
|
||||
}
|
||||
fn seq_to_str(seq: &Vec<u8>) -> String {
|
||||
#[allow(dead_code)]
|
||||
fn seq_to_str(seq: &[u8]) -> String {
|
||||
if seq.len() == 4 {
|
||||
let mut rdr = Cursor::new(seq);
|
||||
let sequence = rdr
|
||||
|
||||
Reference in New Issue
Block a user