diff --git a/bal-pusher.sh b/bal-pusher.sh new file mode 100644 index 0000000..42ee3ae --- /dev/null +++ b/bal-pusher.sh @@ -0,0 +1,14 @@ +RUST_LOG=trace + +export BAL_PUSHER_DB_FILE="$(pwd)/bal.db" +#export BAL_PUSHER_BITCOIN_COOKIE_FILE=/~/.bitcoin/.cookie +#export BAL_PUSHER_REGTEST_COOKIE_FILE=/~/.bitcoin/regtest/.cookie +#export BAL_PUSHER_TESTNET_COOKIE_FILE=/~/.bitcoin/testnet3/.cookie +#export BAL_PUSHER_SIGNET_COOKIE_FILE=/~/.bitcoin/signet/.cookie + +export BAL_PUSHER_REGTEST_ZMQ_HASHBLOCK=tcp://127.0.0.1:18443 +export BAL_PUSHER_SEND_STATS=true +export WELIST_SERVER_URL=http://localhost:8085 +export BAL_SERVER_URL="http://127.0.0.1:9133" +export SSL_KEY_PATH="$(pwd)/private_key.pem" +cargo run --bin=bal-pusher regtest diff --git a/src/bin/bal-server.rs b/src/bin/bal-server.rs index bfbabcc..7d2b7a7 100644 --- a/src/bin/bal-server.rs +++ b/src/bin/bal-server.rs @@ -7,7 +7,6 @@ use chrono::Utc; use hex_conservative::FromHex; use log::{debug, error, info, trace}; use serde::{Deserialize, Serialize}; -use serde_json; use sqlite::State; use sqlite::{Connection, Value}; use std::collections::{HashMap, HashSet}; @@ -119,6 +118,7 @@ pub struct StatsResponse { } #[derive(Debug, Clone)] +#[allow(dead_code)] struct ActixConfig { max_body_size: usize, timeout_secs: u64, @@ -208,7 +208,7 @@ async fn echo_pub_key(data: web::Data) -> impl Responder { "Failed to read public key file {}: {}", data.cfg.pub_key_path, e ); - HttpResponse::InternalServerError().body("Failed to read public key file") + HttpResponse::InternalServerError().body("error") } } } @@ -224,13 +224,13 @@ async fn echo_info( ) -> impl Responder { let param = path.into_inner(); if !NETWORKS.contains(¶m.as_str()) { - return HttpResponse::NotFound().body("Unknown network"); + return HttpResponse::NotFound().body("error"); } info!("echo info!!!{}", param); let netconfig = data.cfg.get_net_config(¶m); if !netconfig.enabled { debug!("network disabled {}", param); - return HttpResponse::BadRequest().body("network disabled"); + return HttpResponse::BadRequest().body("error"); } let remote_addr = req .headers() @@ -257,7 +257,7 @@ async fn echo_info( Ok(g) => g, Err(_p) => { error!("DB mutex poisoned in echo_info (lookup phase)"); - return HttpResponse::InternalServerError().body("DB mutex poisoned"); + return HttpResponse::InternalServerError().body("error"); } }; match get_last_used_address_by_ip( @@ -275,10 +275,7 @@ async fn echo_info( version: VERSION.to_string(), }); } - None => { - let next = get_next_address_index(&db, &netconfig.name, &netconfig.address); - next - } + None => get_next_address_index(&db, &netconfig.name, &netconfig.address), } }; // lock released @@ -288,8 +285,7 @@ async fn echo_info( Ok(address) => address, Err(e) => { error!("Failed to derive address from xpub: {}", e); - return HttpResponse::BadRequest() - .body(format!("Failed to derive address: {}", e)); + return HttpResponse::BadRequest().body("error"); } }; @@ -299,7 +295,7 @@ async fn echo_info( Ok(g) => g, Err(_p) => { error!("DB mutex poisoned in echo_info (save phase)"); - return HttpResponse::InternalServerError().body("DB mutex poisoned"); + return HttpResponse::InternalServerError().body("error"); } }; save_new_address(&db, next_idx.0, &derived.0, &derived.1, &remote_addr); @@ -322,30 +318,30 @@ async fn echo_info( debug!("echo info reply: {}", json_data); HttpResponse::Ok().json(info) } - Err(err) => HttpResponse::InternalServerError().body(format!("error:{}", err)), + Err(_err) => HttpResponse::InternalServerError().body("error"), } } async fn echo_stats(path: web::Path, data: web::Data) -> impl Responder { let param = path.into_inner(); if !NETWORKS.contains(¶m.as_str()) { - return HttpResponse::NotFound().body("Unknown network"); + return HttpResponse::NotFound().body("error"); } info!("echo stats!!! {}", data.cfg.expose_stats); let netconfig = data.cfg.get_net_config(¶m); if !netconfig.enabled { debug!("network disabled {}", param); - return HttpResponse::BadRequest().body("network disabled"); + return HttpResponse::BadRequest().body("error"); } if !data.cfg.expose_stats { - return HttpResponse::Forbidden().body("Stats not exposed"); + return HttpResponse::Forbidden().body("error"); } let mut stats: Vec = vec![]; let db = match data.db.lock() { Ok(g) => g, Err(_p) => { error!("DB mutex poisoned in echo_stats"); - return HttpResponse::InternalServerError().body("DB mutex poisoned"); + return HttpResponse::InternalServerError().body("error"); } }; let mut stmt = match db.prepare( @@ -354,12 +350,12 @@ async fn echo_stats(path: web::Path, data: web::Data) -> impl Ok(s) => s, Err(e) => { error!("Failed to prepare stats query: {}", e); - return HttpResponse::InternalServerError().body("Database error"); + return HttpResponse::InternalServerError().body("error"); } }; if let Err(e) = stmt.bind((1, Value::String(netconfig.name.clone()))) { error!("Failed to bind chain in stats query: {}", e); - return HttpResponse::InternalServerError().body("Database error"); + return HttpResponse::InternalServerError().body("error"); } while let Ok(State::Row) = stmt.next() { let report_date = stmt.read("report_date").unwrap_or("0".to_string()); @@ -422,7 +418,7 @@ async fn echo_stats(path: web::Path, data: web::Data) -> impl debug!("echo info reply: {}", json_data); HttpResponse::Ok().json(stats) } - Err(err) => HttpResponse::InternalServerError().body(format!("error:{}", err)), + Err(_err) => HttpResponse::InternalServerError().body("error"), } } @@ -431,33 +427,33 @@ async fn echo_search(body: Bytes, data: web::Data) -> impl Responder { let strbody = match std::str::from_utf8(&body) { Ok(s) => s, Err(_) => { - return HttpResponse::BadRequest().body("Invalid UTF-8 body"); + return HttpResponse::BadRequest().body("error"); } }; info!("{}", strbody); if strbody.is_empty() || strbody.len() != 64 || !strbody.chars().all(|c| c.is_ascii_hexdigit()) { - return HttpResponse::BadRequest().body("Invalid txid"); + return HttpResponse::BadRequest().body("error"); } let db = match data.db.lock() { Ok(g) => g, Err(_p) => { error!("DB mutex poisoned in echo_search"); - return HttpResponse::InternalServerError().body("DB mutex poisoned"); + return HttpResponse::InternalServerError().body("error"); } }; let mut statement = match db.prepare("SELECT * FROM tbl_tx WHERE txid = ? LIMIT 1") { Ok(s) => s, Err(e) => { error!("Failed to prepare statement: {}", e); - return HttpResponse::InternalServerError().body("Database error"); + return HttpResponse::InternalServerError().body("error"); } }; if let Err(e) = statement.bind((1, strbody)) { error!("Failed to bind parameter: {}", e); - return HttpResponse::InternalServerError().body("Database error"); + return HttpResponse::InternalServerError().body("error"); } if let Ok(State::Row) = statement.next() { @@ -503,11 +499,14 @@ async fn echo_search(body: Bytes, data: web::Data) -> impl Responder { } } match serde_json::to_string(&response_data) { - Ok(json_data) => HttpResponse::Ok().json(json_data), - Err(_) => HttpResponse::BadRequest().body("Bad data received"), + Ok(json_data) => { + debug!("echo search reply: {}", json_data); + HttpResponse::Ok().json(&response_data) + } + Err(_) => HttpResponse::BadRequest().body("error"), } } else { - HttpResponse::BadRequest().body("Bad data received") + HttpResponse::BadRequest().body("error") } } @@ -524,13 +523,13 @@ struct ParsedTx { } /// Parse all transactions from the request body **without** needing the DB lock. -/// Returns `Ok(parsed_txs)` if at least one tx was valid, or `Err(HttpResponse)` for early failure. +/// Skips transactions that don't have a valid willexecutor output. fn parse_request_transactions( strbody: &str, _req_time: i64, netconfig: &NetConfig, known_addresses: &HashSet, -) -> Result, HttpResponse> { +) -> Vec<(ParsedTx, String, u64)> { let mut result: Vec<(ParsedTx, String, u64)> = Vec::new(); let mut union_tx = true; @@ -613,8 +612,8 @@ fn parse_request_transactions( } if !found { - error!("willexecutor output not found for tx {}", txid); - return Err(HttpResponse::BadRequest().body("Bad data received")); + 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 @@ -636,7 +635,7 @@ fn parse_request_transactions( )); } - Ok(result) + result } async fn echo_push( @@ -648,24 +647,24 @@ async fn echo_push( let strbody = match std::str::from_utf8(&body) { Ok(s) => s, Err(_) => { - return HttpResponse::BadRequest().body("Invalid UTF-8 body"); + return HttpResponse::BadRequest().body("error"); } }; let param = path.into_inner(); if !NETWORKS.contains(¶m.as_str()) { - return HttpResponse::NotFound().body("Unknown network"); + return HttpResponse::NotFound().body("error"); } let netconfig = data.cfg.get_net_config(¶m); if !netconfig.enabled { trace!("network not enabled {}", &netconfig.name); - return HttpResponse::BadRequest().body("Network not enabled"); + return HttpResponse::BadRequest().body("error"); } let req_time = match Utc::now().timestamp_nanos_opt() { Some(t) => t, None => { error!("Invalid timestamp"); - return HttpResponse::BadRequest().body("Invalid timestamp"); + return HttpResponse::BadRequest().body("error"); } }; @@ -675,7 +674,7 @@ async fn echo_push( Ok(g) => g, Err(_p) => { error!("DB mutex poisoned acquiring addresses in echo_push"); - return HttpResponse::InternalServerError().body("DB mutex poisoned"); + return HttpResponse::InternalServerError().body("error"); } }; if netconfig.xpub { @@ -683,7 +682,7 @@ async fn echo_push( Ok(addrs) => addrs, Err(e) => { error!("Failed to load addresses from xpub: {}", e); - return HttpResponse::InternalServerError().body("Database error"); + return HttpResponse::InternalServerError().body("error"); } } } else { @@ -692,12 +691,9 @@ async fn echo_push( }; // lock released here // Parse all transactions (CPU-bound, no DB needed) - let parsed = match parse_request_transactions(strbody, req_time, netconfig, &known_addresses) { - Ok(v) => v, - Err(resp) => return resp, - }; + let parsed = parse_request_transactions(strbody, req_time, netconfig, &known_addresses); if parsed.is_empty() { - return HttpResponse::Ok().body("thx"); + return HttpResponse::BadRequest().body("error"); } let all_txids: Vec = parsed.iter().map(|(p, _, _)| p.txid.clone()).collect(); @@ -708,14 +704,14 @@ async fn echo_push( Ok(g) => g, Err(_p) => { error!("DB mutex poisoned in echo_push duplicate check"); - return HttpResponse::InternalServerError().body("DB mutex poisoned"); + return HttpResponse::InternalServerError().body("error"); } }; match check_duplicate_txids(&db, &all_txids) { Ok(dups) => dups, Err(e) => { error!("Duplicate check failed: {}", e); - return HttpResponse::InternalServerError().body("Database error"); + return HttpResponse::InternalServerError().body("error"); } } }; // lock released here @@ -731,7 +727,7 @@ async fn echo_push( Ok(g) => g, Err(_p) => { error!("DB mutex poisoned in echo_push insert phase"); - return HttpResponse::InternalServerError().body("DB mutex poisoned"); + return HttpResponse::InternalServerError().body("error"); } }; @@ -818,7 +814,7 @@ async fn echo_push( if let Err(err) = execute_insert(&db, sqltxs, ptx, sqlinps, pinps, sqlouts, pouts) { error!("execute_insert failed: {}", err); - return HttpResponse::BadRequest().body("Bad data received"); + return HttpResponse::BadRequest().body("error"); } } // lock released @@ -900,7 +896,7 @@ async fn main() -> std::io::Result<()> { let db = match open_db(&cfg.db_file) { Ok(c) => c, Err(e) => { - return Err(std::io::Error::new(std::io::ErrorKind::Other, e)); + return Err(std::io::Error::other(e)); } };