diff --git a/src/bin/bal-pusher.rs b/src/bin/bal-pusher.rs index 3fae55a..3bf8c34 100644 --- a/src/bin/bal-pusher.rs +++ b/src/bin/bal-pusher.rs @@ -5,7 +5,6 @@ use bitcoin::Network; use bitcoincore_rpc::{Auth, Client, Error, RpcApi, bitcoin}; use bitcoincore_rpc_json::GetBlockchainInfoResult; -use byteorder::{LittleEndian, ReadBytesExt}; use ed25519_dalek::{Signer as _, SigningKey, pkcs8::DecodePrivateKey}; use log::{debug, error, info, trace, warn}; use serde::Deserialize; @@ -15,21 +14,19 @@ use sqlite::{Connection, Value}; use std::collections::HashMap; use std::env; use std::error::Error as StdError; -use std::io::Cursor; use std::str; use std::{thread, time::Duration}; -use zmq::{Context, DEALER, DONTWAIT, Socket}; +use zmq::{Context, Socket}; use bal_server::db::open_db; use bal_server::validation::is_valid_welist_url; use base64::{Engine as _, engine::general_purpose}; use reqwest::Client as rClient; use std::net::SocketAddr; -use std::time::Instant; use url::Url; const LOCKTIME_THRESHOLD: i64 = 5000000; -const VERSION: &str = "0.0.2"; +const VERSION: &str = env!("CARGO_PKG_VERSION"); #[derive(Debug, Clone, Serialize, Deserialize)] struct MyConfig { db_file: String, @@ -87,7 +84,7 @@ fn get_network_params(cfg: &MyConfig, network: Network) -> &NetworkParams { fn get_network_params_default(network: Network) -> NetworkParams { match network { Network::Testnet => NetworkParams { - host: "http://i27.0.0.1".to_string(), + host: "http://127.0.0.1".to_string(), port: 18332, dir_path: "testnet3/".to_string(), db_field: "testnet".to_string(), @@ -97,7 +94,7 @@ fn get_network_params_default(network: Network) -> NetworkParams { zmq_listener: "tcp://127.0.0.1:23332".to_string(), }, Network::Testnet4 => NetworkParams { - host: "http://i27.0.0.1".to_string(), + host: "http://127.0.0.1".to_string(), port: 48332, dir_path: "testnet4/".to_string(), db_field: "testnet4".to_string(), @@ -205,29 +202,9 @@ fn get_client( } } async fn main_result(cfg: &MyConfig, network_params: &NetworkParams) -> Result<(), Error> { - /*let url = args.next().expect("Usage: "); - let user = args.next().expect("no user given"); - let pass = args.next().expect("no pass given"); - */ - //let network = Network::Regtest match get_client(network_params) { Ok((rpc, bcinfo)) => { info!("connected"); - //let best_block_hash = rpc.get_best_block_hash()?; - //info!("best block hash: {}", best_block_hash); - //let bestblockcount = rpc.get_block_count()?; - //info!("best block height: {}", bestblockcount); - //let best_block_hash_by_height = rpc.get_block_hash(bestblockcount)?; - //info!("best block hash by height: {}", best_block_hash_by_height); - //assert_eq!(best_block_hash_by_height, best_block_hash); - //let from_block= std::cmp::max(0, bestblockcount - 11); - //let mut time_sum:u64=0; - //for i in from_block..bestblockcount{ - // let hash = rpc.get_block_hash(i).unwrap(); - // let block: bitcoin::Block = rpc.get_by_id(&hash).unwrap(); - // time_sum += >::into(block.header.time); - //} - //let average_time = time_sum/11; info!("median time: {}", bcinfo.median_time); //info!("height time: {}",bcinfo.median_time); info!("blocks: {}", bcinfo.blocks); @@ -288,26 +265,10 @@ async fn main_result(cfg: &MyConfig, network_params: &NetworkParams) -> Result<( info!("to be pushed: {}: {}", txid, locktime); match rpc.send_raw_transaction(tx) { Ok(o) => { - /*let mut file = OpenOptions::new() - .append(true) // Set the append option - .create(true) // Create the file if it doesn't exist - .open("valid_txs")?; - let data = format!("{}\t:\t{}\t:\t{}\n",txid,average_time,locktime); - file.write_all(data.as_bytes())?; - drop(file); - */ info!("tx: {} pusshata PUSHED\n{}", txid, o); pushed_txs.push(txid.to_string()); } Err(err) => { - /*let mut file = OpenOptions::new() - .append(true) // Set the append option - .create(true) // Create the file if it doesn't exist - .open("/home/bal/invalid_txs")?; - let data = format!("{}:\t{}\t:\t{}\t:\t{}\n",txid,err,average_time,locktime); - file.write_all(data.as_bytes())?; - drop(file); - */ warn!("Error: {}\n{}", err, txid); //store err in invalid_txs invalid_txs.insert(txid.to_string(), err.to_string()); @@ -317,16 +278,37 @@ async fn main_result(cfg: &MyConfig, network_params: &NetworkParams) -> Result<( for txid in &pushed_txs { let sql = "UPDATE tbl_tx SET status = 1 WHERE txid = ?"; - let mut stmt = db.prepare(sql).unwrap(); - stmt.bind((1, Value::String(txid.clone()))).unwrap(); - let _ = stmt.next(); + match db.prepare(sql) { + Ok(mut stmt) => { + if let Err(e) = stmt.bind((1, Value::String(txid.clone()))) { + error!("Failed to bind txid for status update: {}", e); + continue; + } + let _ = stmt.next(); + } + Err(e) => { + error!("Failed to prepare status update: {}", e); + } + } } for (txid, txerr) in &invalid_txs { let sql = "UPDATE tbl_tx SET status = 2, push_err = ? WHERE txid = ?"; - let mut stmt = db.prepare(sql).unwrap(); - stmt.bind((1, Value::String(txerr.clone()))).unwrap(); - stmt.bind((2, Value::String(txid.clone()))).unwrap(); - let _ = stmt.next(); + match db.prepare(sql) { + Ok(mut stmt) => { + if let Err(e) = stmt.bind((1, Value::String(txerr.clone()))) { + error!("Failed to bind txerr for error update: {}", e); + continue; + } + if let Err(e) = stmt.bind((2, Value::String(txid.clone()))) { + error!("Failed to bind txid for error update: {}", e); + continue; + } + let _ = stmt.next(); + } + Err(e) => { + error!("Failed to prepare error update: {}", e); + } + } } if let Err(e) = send_stats_report(cfg, bcinfo).await { error!("send_stats_report failed: {}", e); @@ -391,31 +373,6 @@ ON CONFLICT(chain) DO UPDATE SET " ); - /* - let sql = format!("CREATE TABLE tbl_stats AS - SELECT - CURRENT_TIMESTAMP AS report_date, - '{chain}' as chain, - (SELECT COUNT(*) FROM tbl_tx WHERE network ='{chain}') AS totals, - (SELECT COUNT(*) FROM tbl_tx WHERE status = 0 AND network ='{chain}') AS waiting, - (SELECT COUNT(*) FROM tbl_tx WHERE status = 1 AND network ='{chain}') AS sent, - (SELECT COUNT(*) FROM tbl_tx WHERE status = 2 AND network ='{chain}') AS failed, - (SELECT SUM(our_fees) FROM tbl_tx WHERE status = 0 AND network ='{chain}') AS waiting_profit, - (SELECT SUM(our_fees) OR 0 FROM tbl_tx WHERE status = 1 AND network ='{chain}') AS sent_profit, - (SELECT SUM(our_fees) FROM tbl_tx WHERE status = 2 AND network ='{chain}') AS missed_profit, - (SELECT COUNT(*) FROM tbl_inp JOIN tbl_tx ON(tbl_inp.txid = tbl_tx.txid) WHERE tbl_tx.status=0 AND tbl_tx.network ='{chain}') AS unique_inputs; - "); - let sql = "UPDATE tbl_stats set - totals = (SELECT COUNT(*) FROM tbl_tx WHERE network ='{chain}'), - waiting = (SELECT COUNT(*) FROM tbl_tx WHERE status = 0 AND network ='{chain}'), - sent = (SELECT COUNT(*) FROM tbl_tx WHERE status = 1 AND network ='{chain}'), - failed = (SELECT COUNT(*) FROM tbl_tx WHERE status = 1 AND network ='{chain}'), - waiting_profit = (SELECT SUM(our_fees) FROM tbl_tx WHERE status = 0 AND network ='{chain}'), - sent_profit = (SELECT SUM(our_fees) FROM tbl_tx WHERE status = 0 AND network ='{chain}'), - missed_profit = (SELECT SUM(our_fees) FROM tbl_tx WHERE status = 0 AND network ='{chain}') - unique_inputs = (SELECT COUNT(*) FROM tbl_inp JOIN tbl_tx ON(tbl_inp.txid = tbl_tx.txid) WHERE tbl_tx.status=0 AND tbl_tx.network ='{chain}') - WHERE chain = '{chain}' - */ if let Err(err) = db.execute(&sql) { error!("error inserting creating stats table {err}"); } else { @@ -609,85 +566,12 @@ fn parse_env_netconfig(cfg_lock: &mut MyConfig, chain: &str) -> NetworkParams { 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(); - let socket = match context.socket(DEALER) { - Ok(sock) => sock, - Err(_) => return false, - }; - - if socket.connect(endpoint).is_err() { - return false; - } - - // Try to send an empty message non-blocking - socket.send("", DONTWAIT).is_ok() -} - -// Add this struct to monitor connection health -#[allow(dead_code)] -struct ConnectionMonitor { - last_message_time: Instant, - timeout: Duration, - consecutive_timeouts: u32, - max_consecutive_timeouts: u32, -} - -#[allow(dead_code)] -impl ConnectionMonitor { - fn new(timeout_secs: u64, max_timeouts: u32) -> Self { - Self { - last_message_time: Instant::now(), - timeout: Duration::from_secs(timeout_secs), - consecutive_timeouts: 0, - max_consecutive_timeouts: max_timeouts, - } - } - - fn update(&mut self) { - self.last_message_time = Instant::now(); - self.consecutive_timeouts = 0; - } - - fn check_connection(&mut self) -> ConnectionStatus { - let elapsed = self.last_message_time.elapsed(); - - if elapsed > self.timeout { - self.consecutive_timeouts += 1; - - if self.consecutive_timeouts >= self.max_consecutive_timeouts { - ConnectionStatus::Lost(elapsed) - } else { - ConnectionStatus::Warning(elapsed) - } - } else { - ConnectionStatus::Healthy - } - } - - fn reset(&mut self) { - self.consecutive_timeouts = 0; - self.last_message_time = Instant::now(); - } -} - -#[allow(dead_code)] -enum ConnectionStatus { - Healthy, - Warning(Duration), - Lost(Duration), -} - #[tokio::main] async fn main() -> std::io::Result<()> { env_logger::init(); @@ -747,15 +631,15 @@ async fn main() -> std::io::Result<()> { Ok(m) => m, Err(e) => { consecutive_timeouts += 1; - if consecutive_timeouts == 1 { - warn!("ZMQ recv timeout or error: {}, retrying...", e); - } else if consecutive_timeouts.is_multiple_of(12) { - warn!( + if consecutive_timeouts.is_multiple_of(720) { + error!( "No ZMQ messages for {}s ({} consecutive timeouts), is bitcoind ZMQ active on {}?", consecutive_timeouts * 5, consecutive_timeouts, zmq_address ); + } else { + trace!("ZMQ recv timeout or error: {}, retrying...", e); } continue; } @@ -783,17 +667,6 @@ async fn main() -> std::io::Result<()> { thread::sleep(Duration::from_millis(100)); // Sleep for 100ms } } -#[allow(dead_code)] -fn seq_to_str(seq: &[u8]) -> String { - if seq.len() == 4 { - let mut rdr = Cursor::new(seq); - let sequence = rdr - .read_u32::() - .expect("Failed to read integer"); - return sequence.to_string(); - } - "Unknown".to_string() -} #[cfg(test)] mod tests { diff --git a/src/bin/bal-server.rs b/src/bin/bal-server.rs index 7d2b7a7..8c2304d 100644 --- a/src/bin/bal-server.rs +++ b/src/bin/bal-server.rs @@ -118,7 +118,7 @@ pub struct StatsResponse { } #[derive(Debug, Clone)] -#[allow(dead_code)] +#[expect(dead_code)] struct ActixConfig { max_body_size: usize, timeout_secs: u64, @@ -413,13 +413,8 @@ async fn echo_stats(path: web::Path, data: web::Data) -> impl unique_inputs, }); } - match serde_json::to_string(&stats) { - Ok(json_data) => { - debug!("echo info reply: {}", json_data); - HttpResponse::Ok().json(stats) - } - Err(_err) => HttpResponse::InternalServerError().body("error"), - } + debug!("echo stats reply for chain: {}", netconfig.name); + HttpResponse::Ok().json(stats) } async fn echo_search(body: Bytes, data: web::Data) -> impl Responder { @@ -531,7 +526,6 @@ fn parse_request_transactions( known_addresses: &HashSet, ) -> Vec<(ParsedTx, String, u64)> { let mut result: Vec<(ParsedTx, String, u64)> = Vec::new(); - let mut union_tx = true; for line in strbody.split('\n') { if line.is_empty() { @@ -615,11 +609,6 @@ fn parse_request_transactions( trace!("willexecutor output not found for tx {}, skipping", txid); continue; } - if !union_tx { - // This is only used for SQL building later; we track it in the caller - } else { - union_tx = false; - } result.push(( ParsedTx { txid, @@ -911,15 +900,6 @@ async fn main() -> std::io::Result<()> { cfg: cfg.clone(), }); - // Initialize networks - { - let db = data.db.lock().unwrap(); - for network in NETWORKS { - let netconfig = data.cfg.get_net_config(network); - insert_xpub(&db, &netconfig.name.to_string(), &netconfig.address); - } - } - let bind_address = data.cfg.bind_address.clone(); let bind_port = data.cfg.bind_port; diff --git a/src/db.rs b/src/db.rs index b096352..9e3fe89 100644 --- a/src/db.rs +++ b/src/db.rs @@ -164,7 +164,7 @@ pub fn create_database(db: &Connection) { let _ = db.execute("DROP INDEX IF EXISTS idx_stats_chain;"); let _ = db.execute("CREATE UNIQUE INDEX IF NOT EXISTS idx_stats_chain ON tbl_stats(chain);"); - let _ = db.execute("UPDATE tbl_tx set network='bitcoin' where network='mainnet');"); + let _ = db.execute("UPDATE tbl_tx set network='bitcoin' where network='mainnet';"); } /* pub fn get_xpub_id(db: &Connection, network: &String, xpub: &String) -> Option{ @@ -181,13 +181,14 @@ pub fn create_database(db: &Connection) { 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, - Err(e) => { - error!("Failed to prepare xpub insert statement: {}", e); - return; - } - }; + let mut stmt = + match db.prepare("INSERT OR IGNORE INTO tbl_xpub(network,xpub) VALUES(?, ?);") { + Ok(s) => s, + Err(e) => { + error!("Failed to prepare xpub insert statement: {}", e); + return; + } + }; if let Err(e) = stmt.bind((1, Value::String(network.to_string()))) { error!("Failed to bind network parameter for xpub insert: {}", e); return;