From cd24eda1114aa74d0148b2bf6e50d51b35387316 Mon Sep 17 00:00:00 2001 From: svatantrya Date: Sat, 18 Jul 2026 22:58:42 -0400 Subject: [PATCH] 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.) --- .gitignore | 3 - Cargo.lock | 2 +- Cargo.toml | 2 +- src/bin/bal-pusher.rs | 150 +++++++++++++------------------- src/db.rs | 4 +- src/xpub.rs | 32 ++----- tests/db_path_validation.rs | 3 +- tests/panic_regression_tests.rs | 7 +- tests/secret_leakage_tests.rs | 94 ++++++++------------ 9 files changed, 111 insertions(+), 186 deletions(-) diff --git a/.gitignore b/.gitignore index 905c922..50e779f 100644 --- a/.gitignore +++ b/.gitignore @@ -8,9 +8,6 @@ .env.production .env.secret -# Shell scripts that load env vars (contain secrets, local only) -bal-pusher.sh -bal-server.sh # Private keys - NEVER commit to git # Only public_key.pem should be tracked (if needed) diff --git a/Cargo.lock b/Cargo.lock index f0f2700..317029b 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -312,7 +312,7 @@ checksum = "ace50bade8e6234aa140d9a2f552bbee1db4d353f69b8217bc503490fc1a9f26" [[package]] name = "bal_server" -version = "0.3.0" +version = "0.3.1" dependencies = [ "actix-governor", "actix-rt", diff --git a/Cargo.toml b/Cargo.toml index 250f4ea..eb91e64 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "bal_server" -version = "0.3.0" +version = "0.3.1" edition = "2024" # See more keys and their definitions at https://doc.rust-lang.org/cargo/reference/manifest.html diff --git a/src/bin/bal-pusher.rs b/src/bin/bal-pusher.rs index 98e5154..3973e77 100644 --- a/src/bin/bal-pusher.rs +++ b/src/bin/bal-pusher.rs @@ -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> { - 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 Result<(Client, GetBlockchainInfoResult), Box> { - 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> { 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> { - 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 = 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::() { - 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::() + { + 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 = [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) -> String { +#[allow(dead_code)] +fn seq_to_str(seq: &[u8]) -> String { if seq.len() == 4 { let mut rdr = Cursor::new(seq); let sequence = rdr diff --git a/src/db.rs b/src/db.rs index a18c623..b096352 100644 --- a/src/db.rs +++ b/src/db.rs @@ -178,8 +178,8 @@ pub fn create_database(db: &Connection) { } } */ -pub fn insert_xpub(db: &Connection, network: &String, xpub: &String) { - if xpub != "" { +pub fn insert_xpub(db: &Connection, network: &str, xpub: &str) { + if !xpub.is_empty() { trace!("going to insert: {} xpub:{}", network, xpub); let mut stmt = match db.prepare("INSERT INTO tbl_xpub(network,xpub) VALUES(?, ?);") { Ok(s) => s, diff --git a/src/xpub.rs b/src/xpub.rs index 38777c5..0874926 100644 --- a/src/xpub.rs +++ b/src/xpub.rs @@ -11,6 +11,7 @@ use sha2::{Digest, Sha256}; use std::str::FromStr; // Mainnet (BIP44/BIP49/BIP84) +#[allow(dead_code)] enum BS58Prefix { Xpub, Ypub, @@ -102,44 +103,29 @@ fn calc_checksum(desc: &str) -> Result { Ok(checksum) } -pub fn get_bitcoincore_descriptor(xpub: &String) -> String { +pub fn get_bitcoincore_descriptor(xpub: &str) -> String { let fingerprint = match calculate_fingerprint(xpub) { Ok(f) => f, Err(_) => return String::new(), // Invalid xpub, return empty descriptor }; - let mut bip = 84; - let cpub = xpub.to_string(); - match &xpub[0..4] { - "vpub" => { - bip = 84; - } - "zpub" => { - bip = 84; - } - &_ => { - bip = 84; - } - }; let xpub_converted = match convert_xpub(xpub) { Ok(c) => c, Err(_) => return String::new(), // Invalid xpub, return empty descriptor }; let descriptor = format!("wpkh([{}/84h/0h/0h]{}/0/*)", fingerprint, xpub_converted); - let descriptor = match calc_checksum(&descriptor) { + match calc_checksum(&descriptor) { Ok(checksum) => { let clean_descriptor = descriptor.split('#').next().unwrap_or(&descriptor); format!("{}#{}", clean_descriptor, checksum) } Err(err) => { eprintln!("Error: {}", err); - "".to_string() + String::new() } - }; - descriptor - //format!("{}#{}",descriptor,checksum) + } } -fn convert_xpub(xpub: &String) -> Result { +fn convert_xpub(xpub: &str) -> Result { if xpub.len() >= 4 && (&xpub[0..4] == "xpub" || &xpub[0..4] == "ypub" || &xpub[0..4] == "zpub") { convert_to(xpub, BS58Prefix::Xpub) @@ -165,7 +151,7 @@ fn base58check_decode(s: &str) -> Result, String> { return Err("Data troppo corta".to_string()); } let (payload, checksum) = data.split_at(data.len() - 4); - let hash = Sha256::digest(&Sha256::digest(payload)); + let hash = Sha256::digest(Sha256::digest(payload)); if hash[0..4] != checksum[..] { return Err("Checksum invalido".to_string()); } @@ -173,7 +159,7 @@ fn base58check_decode(s: &str) -> Result, String> { } fn base58check_encode(data: &[u8]) -> String { - let checksum = &Sha256::digest(&Sha256::digest(data))[0..4]; + let checksum = &Sha256::digest(Sha256::digest(data))[0..4]; let full = [data, checksum].concat(); bs58::encode(full).into_string() } @@ -205,7 +191,7 @@ pub fn new_address_from_xpub( ) -> Result<(String, String), Box> { let xpub = Xpub::from_str(&convert_to(zpub, BS58Prefix::Xpub)?)?; let path = format!("m/0/{}", index); - let derivation_path = DerivationPath::from_str(&path.as_str())?; + let derivation_path = DerivationPath::from_str(path.as_str())?; let secp = Secp256k1::new(); let derived_xpub = xpub.derive_pub(&secp, &derivation_path)?; let public_key = derived_xpub.public_key; diff --git a/tests/db_path_validation.rs b/tests/db_path_validation.rs index f156437..324b98d 100644 --- a/tests/db_path_validation.rs +++ b/tests/db_path_validation.rs @@ -1,7 +1,6 @@ use bal_server::db::open_db; use sqlite::State; use std::fs; -use std::path::Path; #[test] fn test_open_db_blocks_traversal() { @@ -74,7 +73,7 @@ fn test_open_db_rejects_symlink() { let _ = fs::remove_file(real); let _ = fs::remove_file(link); fs::File::create(real).unwrap(); - fs::soft_link(real, link).unwrap(); + std::os::unix::fs::symlink(real, link).unwrap(); let res = open_db(link); assert!(res.is_err(), "Symlink DB path should be rejected"); diff --git a/tests/panic_regression_tests.rs b/tests/panic_regression_tests.rs index 96d902f..d43473b 100644 --- a/tests/panic_regression_tests.rs +++ b/tests/panic_regression_tests.rs @@ -36,11 +36,8 @@ fn test_db_null_unwrap_or() { let mut found_value = None; let _ = db.iterate("SELECT * FROM test_stats;", |pairs| { - let row: HashMap<_, _> = pairs - .into_iter() - .map(|(k, v)| (k.to_string(), v.map(|s| s))) - .collect(); - let totals = row["totals"].clone().unwrap_or("0").to_string(); + let row: HashMap<_, _> = pairs.iter().map(|(k, v)| (k.to_string(), *v)).collect(); + let totals = row["totals"].unwrap_or("0").to_string(); found_value = Some(totals); true }); diff --git a/tests/secret_leakage_tests.rs b/tests/secret_leakage_tests.rs index 7016f73..6011ca6 100644 --- a/tests/secret_leakage_tests.rs +++ b/tests/secret_leakage_tests.rs @@ -1,5 +1,4 @@ use std::fs; -use std::path::Path; #[test] fn test_gitignore_protection_env() { @@ -23,30 +22,20 @@ fn test_gitignore_protection_env() { ]; for pattern in required_patterns { - let has_exact = gitignore.contains(&pattern); - let has_wildcard = gitignore.contains(&format!("*.env.local")) - || gitignore.contains(&format!(".env.local")); + let has_wildcard = gitignore.contains("*.env.local") || gitignore.contains(".env.local"); let has_env = gitignore.contains("*.env") || gitignore.contains(".env"); - // For .env.local, either .env.local or *.env.local is acceptable let is_env_local = pattern == "*.env.local" || pattern == ".env.local"; if is_env_local { assert!( has_wildcard, ".gitignore must contain pattern '*.env.local' or '.env.local' to protect secrets", ); - } else if pattern == ".env.production" || pattern == ".env.secret" { - assert!( - gitignore.contains(pattern), - ".gitignore must contain pattern '{}' to protect secrets", - pattern - ); - } else if pattern == "*.env" { - assert!( - has_env, - ".gitignore must contain pattern '*.env' or '.env' to protect secrets", - ); - } else if pattern == ".env" { + } else if pattern == ".env.production" + || pattern == ".env.secret" + || pattern == "*.env" + || pattern == ".env" + { assert!( has_env, ".gitignore must contain pattern '*.env' or '.env' to protect secrets", @@ -65,7 +54,6 @@ fn test_gitignore_protection_env() { #[test] fn test_no_private_key_in_git() { - // Check that .gitignore includes private_key.pem let gitignore = match fs::read_to_string(".gitignore") { Ok(c) => c, Err(e) => { @@ -88,20 +76,17 @@ fn test_no_private_key_in_git() { ".gitignore must block chiave_privata.key" ); - // Check that no private key files are tracked by git let output = std::process::Command::new("git") - .args(&["ls-files", "*.pem", "*.key"]) + .args(["ls-files", "*.pem", "*.key"]) .output() .expect("Failed to run git ls-files"); let tracked_keys = String::from_utf8(output.stdout).unwrap(); let tracked_keys: Vec<&str> = tracked_keys.lines().collect(); - // Only non-empty entries and only public_key.pem should be tracked for tracked in tracked_keys.iter().filter(|s| !s.is_empty()) { if !tracked.contains("public_key.pem") { - assert!( - false, + panic!( "Private key file is tracked by git: {}. Remove it with git rm --cached", tracked ); @@ -113,47 +98,41 @@ fn test_no_private_key_in_git() { #[test] fn test_no_token_in_source_files() { - // Scan source files for hardcoded tokens let mut found_issues = Vec::new(); - // Scan .sh files for hardcoded 40-char hex strings for entry in fs::read_dir(".").unwrap().filter_map(|e| e.ok()) { let path = entry.path(); if !path.is_file() { continue; } - if let Some(ext) = path.extension() { - if ext == "sh" { - let content = fs::read_to_string(&path).unwrap(); - for (line_num, line) in content.lines().enumerate() { - // Skip comments and example/template files - if line.trim().starts_with("#") - || line.to_lowercase().contains("example") - || line.to_lowercase().contains("template") + if let Some(ext) = path.extension() + && ext == "sh" + { + let content = fs::read_to_string(&path).unwrap(); + for (line_num, line) in content.lines().enumerate() { + if line.trim().starts_with('#') + || line.to_lowercase().contains("example") + || line.to_lowercase().contains("template") + { + continue; + } + if line.trim().len() >= 40 { + let hex_chars: Vec<_> = line + .trim() + .chars() + .filter(|c| c.is_ascii_hexdigit()) + .collect(); + if (40..=64).contains(&hex_chars.len()) + && (line.to_lowercase().contains("token") + || line.to_lowercase().contains("api") + || line.to_lowercase().contains("secret")) { - continue; - } - // Check for 40-64 hex chars that could be API tokens (not in .env.example comments) - if line.trim().len() >= 40 { - let hex_chars = line - .trim() - .chars() - .filter(|c| c.is_ascii_hexdigit()) - .collect::>(); - if hex_chars.len() >= 40 && hex_chars.len() <= 64 { - // Check if it looks like it's part of a TOKEN assignment - if line.to_lowercase().contains("token") - || line.to_lowercase().contains("api") - || line.to_lowercase().contains("secret") - { - found_issues.push(format!( - "Potential hardcoded token in {}: line {}: {}", - path.display(), - line_num + 1, - line.trim() - )); - } - } + found_issues.push(format!( + "Potential hardcoded token in {}: line {}: {}", + path.display(), + line_num + 1, + line.trim() + )); } } } @@ -165,8 +144,7 @@ fn test_no_token_in_source_files() { for issue in &found_issues { println!(" {}", issue); } - assert!( - false, + panic!( "Found potential hardcoded tokens in shell scripts: {:?}", found_issues );