fix(pushtxs): skip invalid txs instead of aborting batch, anonymize error responses

- parse_request_transactions now skips txs without willexecutor output
  instead of returning an error that aborts the entire batch
- returns 'error' (400) only when NO txs are valid
- all HTTP error bodies replaced with generic 'error' to avoid leaking
  internal details
This commit is contained in:
2026-07-18 21:35:39 -04:00
parent 190cac929e
commit 8ce3f6a445
2 changed files with 59 additions and 49 deletions

14
bal-pusher.sh Normal file
View File

@@ -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

View File

@@ -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<AppState>) -> 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(&param.as_str()) {
return HttpResponse::NotFound().body("Unknown network");
return HttpResponse::NotFound().body("error");
}
info!("echo info!!!{}", param);
let netconfig = data.cfg.get_net_config(&param);
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<String>, data: web::Data<AppState>) -> impl Responder {
let param = path.into_inner();
if !NETWORKS.contains(&param.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(&param);
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<StatsResponse> = 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<String>, data: web::Data<AppState>) -> 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<String>, data: web::Data<AppState>) -> 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<AppState>) -> 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<AppState>) -> 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<String>,
) -> Result<Vec<(ParsedTx, String, u64)>, 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(&param.as_str()) {
return HttpResponse::NotFound().body("Unknown network");
return HttpResponse::NotFound().body("error");
}
let netconfig = data.cfg.get_net_config(&param);
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<String> = 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));
}
};