fix: resolve PostgreSQL type mismatches and add willexecutor debug logging

- Fix get_next_address_index: use try_get::<i32> for PG SERIAL/INTEGER columns
- Fix search_tx: use try_get::<i32> for PG status column
- Fix execute_insert: parse locktime (String→i64) and in_vout (String→i32)
  before binding to PG INTEGER columns
- Fix execute_insert: bind tbl_out vout as i32, amount as String for PG
- Fix save_new_address: cast xpub i64 to i32 for PG INTEGER column
- Fix get_pending_txs: cast i64 bind params to i32, use try_get::<i32> for reads
- Fix get_stats: use try_get::<i32> for all numeric PG INTEGER columns
- Add trace logging in parse_request_transactions for xpub address matching
- Add trace logging in get_all_addresses_by_xpub for query debugging
This commit is contained in:
2026-08-18 03:14:03 -04:00
parent 8dc344cbd1
commit 8f764f06b2
65 changed files with 4879 additions and 1188 deletions

View File

@@ -10,7 +10,6 @@ use log::{debug, error, info, trace, warn};
use serde::Deserialize;
use serde::Serialize;
use serde_json::json;
use sqlite::{Connection, Value};
use std::collections::HashMap;
use std::env;
use std::error::Error as StdError;
@@ -18,7 +17,7 @@ use std::str;
use std::{thread, time::Duration};
use zmq::{Context, Socket};
use bal_server::db::open_db;
use bal_server::db::{calculate_and_upsert_stats, get_pending_txs, open_database};
use bal_server::validation::is_valid_welist_url;
use base64::{Engine as _, engine::general_purpose};
use reqwest::Client as rClient;
@@ -29,7 +28,9 @@ const LOCKTIME_THRESHOLD: i64 = 5000000;
const VERSION: &str = env!("CARGO_PKG_VERSION");
#[derive(Debug, Clone, Serialize, Deserialize)]
struct MyConfig {
db_backend: String,
db_file: String,
pg_dsn: String,
bitcoin_dir: String,
regtest: NetworkParams,
testnet: NetworkParams,
@@ -44,7 +45,9 @@ struct MyConfig {
impl Default for MyConfig {
fn default() -> Self {
MyConfig {
db_backend: env::var("BAL_PUSHER_DB_BACKEND").unwrap_or("sqlite".to_string()),
db_file: env::var("BAL_PUSHER_DB_FILE").unwrap_or("bal.db".to_string()),
pg_dsn: env::var("BAL_PUSHER_PG_DSN").unwrap_or_default(),
bitcoin_dir: env::var("BAL_PUSHER_BITCOIN_DIR").unwrap_or("".to_string()),
regtest: get_network_params_default(Network::Regtest),
testnet: get_network_params_default(Network::Testnet),
@@ -206,114 +209,79 @@ async fn main_result(cfg: &MyConfig, network_params: &NetworkParams) -> Result<(
Ok((rpc, bcinfo)) => {
info!("connected");
info!("median time: {}", bcinfo.median_time);
//info!("height time: {}",bcinfo.median_time);
info!("blocks: {}", bcinfo.blocks);
debug!("best block hash: {}", bcinfo.best_block_hash);
let average_time = bcinfo.median_time;
let db = match open_db(&cfg.db_file) {
Ok(c) => c,
let connection_string = match cfg.db_backend.as_str() {
"sqlite" => cfg.db_file.clone(),
"postgresql" => cfg.pg_dsn.clone(),
other => {
error!("Unknown DB backend: {}", other);
return Ok(());
}
};
let db = match open_database(&cfg.db_backend, &connection_string).await {
Ok(pool) => pool,
Err(e) => {
error!("Fatal: {}", e);
std::process::exit(1);
}
};
info!("db open {}", &cfg.db_file);
info!("db open {}", &connection_string);
let sqlquery = "SELECT * FROM tbl_tx WHERE network = :network AND status = :status AND ( locktime < :bestblock_height OR locktime > :locktime_threshold AND locktime < :bestblock_time);";
let query_tx = match db.prepare(sqlquery) {
Ok(q) => q.into_iter(),
let pending_txs = match get_pending_txs(
&db,
&network_params.db_field,
LOCKTIME_THRESHOLD,
bcinfo.blocks as i64,
average_time as i64,
)
.await
{
Ok(txs) => txs,
Err(e) => {
warn!("tbl_tx not ready yet (tables may not exist): {}", e);
warn!("Failed to query pending transactions: {}", e);
return Ok(());
}
};
trace!("query_tx: {}", sqlquery);
trace!(":locktime_threshold: {}", LOCKTIME_THRESHOLD);
trace!(":bestblock_time: {}", average_time);
trace!(":bestblock_height: {}", bcinfo.blocks);
trace!(":network: {}", network_params.db_field.clone());
trace!(":status: {}", 0);
//let query_tx = db.prepare("SELECT * FROM tbl_tx where status = :status").unwrap().into_iter();
let mut pushed_txs: Vec<String> = Vec::new();
let mut invalid_txs: std::collections::HashMap<String, String> = HashMap::new();
for row_result in match query_tx.bind::<&[(_, Value)]>(
&[
(":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()),
(":status", 0.into()),
][..],
) {
Ok(bound) => bound,
Err(e) => {
error!("Failed to bind query parameters: {}", e);
return Ok(());
}
} {
let row = match row_result {
Ok(r) => r,
Err(e) => {
warn!("Failed to read row: {}", e);
continue;
}
};
let tx = row.read::<&str, _>("tx");
let txid = row.read::<&str, _>("txid");
let locktime = row.read::<i64, _>("locktime");
let mut invalid_txs: HashMap<String, String> = HashMap::new();
for row in &pending_txs {
let txid = &row.txid;
let tx = &row.tx;
let locktime = row.locktime;
info!("to be pushed: {}: {}", txid, locktime);
match rpc.send_raw_transaction(tx) {
match rpc.send_raw_transaction(tx.as_str()) {
Ok(o) => {
info!("tx: {} pusshata PUSHED\n{}", txid, o);
pushed_txs.push(txid.to_string());
pushed_txs.push(txid.clone());
}
Err(err) => {
warn!("Error: {}\n{}", err, txid);
//store err in invalid_txs
invalid_txs.insert(txid.to_string(), err.to_string());
invalid_txs.insert(txid.clone(), err.to_string());
}
};
}
for txid in &pushed_txs {
let sql = "UPDATE tbl_tx SET status = 1 WHERE txid = ?";
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);
}
if let Err(e) = bal_server::db::update_tx_status(&db, txid, 1, None).await {
error!("Failed to update tx status: {}", e);
}
}
for (txid, txerr) in &invalid_txs {
let sql = "UPDATE tbl_tx SET status = 2, push_err = ? WHERE txid = ?";
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) = bal_server::db::update_tx_status(&db, txid, 2, Some(txerr)).await {
error!("Failed to update tx status: {}", e);
}
}
if let Err(e) = send_stats_report(cfg, bcinfo).await {
error!("send_stats_report failed: {}", e);
}
if let Err(e) = calculate_stats(&db, network_params.db_field.clone()).await {
if let Err(e) = calculate_and_upsert_stats(&db, &network_params.db_field).await {
warn!("calculate_stats failed: {e}");
}
}
@@ -325,65 +293,8 @@ async fn main_result(cfg: &MyConfig, network_params: &NetworkParams) -> Result<(
}
Ok(())
}
async fn calculate_stats(db: &Connection, chain: String) -> Result<(), reqwest::Error> {
// Validate chain to prevent SQL injection via environment variable tampering
if !chain
.chars()
.all(|c| c.is_alphanumeric() || c == '-' || c == '_')
|| chain.is_empty()
{
error!("Invalid chain name: {chain}");
return Ok(());
}
//let sql = "drop table if exists tbl_stats;";
let sql = format!("DELETE FROM tbl_stats WHERE chain = '{chain}';");
if let Err(err) = db.execute(&sql) {
error!("error deleting from tbl_stats where chain:{chain} error: {err}");
}
let sql = format!(
"INSERT INTO tbl_stats (
report_date, chain, totals, waiting, sent, failed,
waiting_profit, sent_profit, missed_profit, unique_inputs
)
VALUES (
CURRENT_TIMESTAMP,
'{chain}',
(SELECT COUNT(*) FROM tbl_tx WHERE network = '{chain}'),
(SELECT COUNT(*) FROM tbl_tx WHERE status = 0 AND network = '{chain}'),
(SELECT COUNT(*) FROM tbl_tx WHERE status = 1 AND network = '{chain}'),
(SELECT COUNT(*) FROM tbl_tx WHERE status = 2 AND network = '{chain}'),
(SELECT IFNULL(SUM(our_fees),0) FROM tbl_tx WHERE status = 0 AND network = '{chain}'),
(SELECT IFNULL(SUM(our_fees),0) FROM tbl_tx WHERE status = 1 AND network = '{chain}'),
(SELECT IFNULL(SUM(our_fees),0) FROM tbl_tx WHERE status = 2 AND network = '{chain}'),
(SELECT COUNT(DISTINCT tbl_inp.in_txid)
FROM tbl_inp
JOIN tbl_tx ON tbl_inp.txid = tbl_tx.txid
WHERE tbl_tx.status = 0 AND tbl_tx.network = '{chain}')
)
ON CONFLICT(chain) DO UPDATE SET
report_date = excluded.report_date,
totals = excluded.totals,
waiting = excluded.waiting,
sent = excluded.sent,
failed = excluded.failed,
waiting_profit = excluded.waiting_profit,
sent_profit = excluded.sent_profit,
missed_profit = excluded.missed_profit,
unique_inputs = excluded.unique_inputs;
"
);
if let Err(err) = db.execute(&sql) {
error!("error inserting creating stats table {err}");
} else {
info!("tbl_stats creation success");
}
Ok(())
}
/// Parse the `(host, port)` pair from a base URL like `https://host[:port]`.
///
/// Falls back to the scheme's well-known default port (443 for `https`,
/// 80 for plain `http`), or to 443 when the scheme is unknown.
fn parse_host_port(base_url: &str) -> Option<(String, u16)> {
let url = Url::parse(base_url).ok()?;
let host = url
@@ -396,8 +307,6 @@ fn parse_host_port(base_url: &str) -> Option<(String, u16)> {
}
/// Resolve `host:port` and return the first IPv6 (AAAA) address, if any.
///
/// Returns `None` when the host has no IPv6 address.
async fn resolve_first_ipv6(host: &str, port: u16) -> Option<SocketAddr> {
use std::net::ToSocketAddrs;
let host = host.to_string();
@@ -412,14 +321,6 @@ async fn resolve_first_ipv6(host: &str, port: u16) -> Option<SocketAddr> {
.flatten()
}
/// Build the HTTP client used for welist reports.
///
/// When `BAL_PUSHER_PREFER_IPV6` is truthy, the welist host is resolved and
/// the client is pinned to its first IPv6 address (the original hostname is
/// still used for the `Host` header and TLS SNI). This works around networks
/// where the IPv4 route to the welist host is broken while IPv6 works: the
/// default connector may otherwise pick the broken family and the request
/// stalls. When the variable is unset (the default), behavior is unchanged.
async fn welist_http_client(welist_url: &str) -> rClient {
let prefer_ipv6 = env::var("BAL_PUSHER_PREFER_IPV6")
.unwrap_or("false".to_string())
@@ -521,6 +422,15 @@ fn sign_message(private_key_path: &str, message: &str) -> String {
}
fn parse_env(cfg: &mut MyConfig) {
if let Ok(value) = env::var("BAL_PUSHER_DB_BACKEND") {
cfg.db_backend = value;
}
if let Ok(value) = env::var("BAL_PUSHER_DB_FILE") {
cfg.db_file = value;
}
if let Ok(value) = env::var("BAL_PUSHER_PG_DSN") {
cfg.pg_dsn = value;
}
cfg.regtest = parse_env_netconfig(cfg, "regtest");
cfg.signet = parse_env_netconfig(cfg, "signet");
cfg.testnet = parse_env_netconfig(cfg, "testnet");
@@ -528,7 +438,6 @@ fn parse_env(cfg: &mut MyConfig) {
drop(parse_env_netconfig(cfg, "bitcoin"));
}
fn parse_env_netconfig(cfg_lock: &mut MyConfig, chain: &str) -> NetworkParams {
//fn parse_env_netconfig(cfg_lock: &MutexGuard<MyConfig>, chain: &str) -> &NetworkParams{
let cfg = match chain {
"regtest" => &mut cfg_lock.regtest,
"signet" => &mut cfg_lock.signet,

View File

@@ -7,16 +7,14 @@ use chrono::Utc;
use hex_conservative::FromHex;
use log::{debug, error, info, trace};
use serde::{Deserialize, Serialize};
use sqlite::State;
use sqlite::{Connection, Value};
use std::collections::{HashMap, HashSet};
use std::env;
use std::fs;
use std::sync::Mutex;
use bal_server::db::{
check_duplicate_txids, create_database, execute_insert, get_all_addresses_by_xpub,
get_last_used_address_by_ip, get_next_address_index, insert_xpub, open_db, save_new_address,
DatabasePool, InsertInpData, InsertOutData, InsertTxData, check_duplicate_txids,
create_database, get_all_addresses_by_xpub, get_last_used_address_by_ip,
get_next_address_index, get_stats, insert_xpub, open_database, save_new_address, search_tx,
};
use bal_server::xpub::new_address_from_xpub;
@@ -56,7 +54,9 @@ struct MyConfig {
info: String,
bind_address: String,
bind_port: u16,
db_backend: String,
db_file: String,
pg_dsn: String,
pub_key_path: String,
expose_stats: bool,
}
@@ -71,7 +71,9 @@ impl Default for MyConfig {
mainnet: NetConfig::default_network("bitcoin".to_string(), Network::Bitcoin),
bind_address: "127.0.0.1".to_string(),
bind_port: 9137,
db_backend: "sqlite".to_string(),
db_file: "bal.db".to_string(),
pg_dsn: String::new(),
info: "Will Executor Server".to_string(),
pub_key_path: "public_key.pem".to_string(),
expose_stats: env::var("BAL_SERVER_EXPOSE_STATS")
@@ -192,7 +194,7 @@ fn parse_actix_config() -> ActixConfig {
}
struct AppState {
db: Mutex<Connection>,
db: DatabasePool,
cfg: MyConfig,
}
@@ -272,35 +274,26 @@ async fn echo_info(
address
}
true => {
// Lock #1: fetch existing address OR atomically claim next index
let next_idx = {
let db = match data.db.lock() {
Ok(g) => g,
Err(_p) => {
error!("DB mutex poisoned in echo_info (lookup phase)");
return HttpResponse::InternalServerError().body("error");
}
};
match get_last_used_address_by_ip(
&db,
&netconfig.name,
&netconfig.address,
&remote_addr,
) {
Some(address) => {
return HttpResponse::Ok().json(InfoResponse {
address,
base_fee: netconfig.fixed_fee,
chain: netconfig.network.to_string(),
info: data.cfg.info.to_string(),
version: VERSION.to_string(),
});
}
None => get_next_address_index(&db, &netconfig.name, &netconfig.address),
}
}; // lock released
if let Some(address) = get_last_used_address_by_ip(
&data.db,
&netconfig.name,
&netconfig.address,
&remote_addr,
)
.await
{
return HttpResponse::Ok().json(InfoResponse {
address,
base_fee: netconfig.fixed_fee,
chain: netconfig.network.to_string(),
info: data.cfg.info.to_string(),
version: VERSION.to_string(),
});
}
let next_idx =
get_next_address_index(&data.db, &netconfig.name, &netconfig.address).await;
// Derive address (CPU-bound, no lock held)
let derived =
match new_address_from_xpub(&netconfig.address, next_idx.1, netconfig.network) {
Ok(address) => address,
@@ -310,20 +303,10 @@ async fn echo_info(
}
};
// Lock #2: save the newly derived address
{
let db = match data.db.lock() {
Ok(g) => g,
Err(_p) => {
error!("DB mutex poisoned in echo_info (save phase)");
return HttpResponse::InternalServerError().body("error");
}
};
save_new_address(&db, next_idx.0, &derived.0, &derived.1, &remote_addr);
debug!("save new address {} {}", derived.0, derived.1);
trace!("next {} {}", next_idx.0, next_idx.1);
derived.0
} // lock released
save_new_address(&data.db, next_idx.0, &derived.0, &derived.1, &remote_addr).await;
debug!("save new address {} {}", derived.0, derived.1);
trace!("next {} {}", next_idx.0, next_idx.1);
derived.0
}
};
let info = InfoResponse {
@@ -357,83 +340,31 @@ async fn echo_stats(path: web::Path<String>, data: web::Data<AppState>) -> impl
if !data.cfg.expose_stats {
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("error");
}
};
let mut stmt = match db.prepare(
"SELECT report_date, chain, totals, waiting, sent, failed, waiting_profit, sent_profit, missed_profit, unique_inputs FROM tbl_stats WHERE chain = ?"
) {
Ok(s) => s,
let stats_rows = match get_stats(&data.db, &netconfig.name).await {
Ok(rows) => rows,
Err(e) => {
error!("Failed to prepare stats query: {}", e);
error!("Failed to query stats: {}", e);
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("error");
}
while let Ok(State::Row) = stmt.next() {
let report_date = stmt.read("report_date").unwrap_or("0".to_string());
let chain = stmt.read("chain").unwrap_or("?".to_string());
let totals = stmt
.read("totals")
.unwrap_or("0".to_string())
.parse::<i64>()
.unwrap_or(0);
let waiting = stmt
.read("waiting")
.unwrap_or("0".to_string())
.parse::<i64>()
.unwrap_or(0);
let sent = stmt
.read("sent")
.unwrap_or("0".to_string())
.parse::<i64>()
.unwrap_or(0);
let failed = stmt
.read("failed")
.unwrap_or("0".to_string())
.parse::<i64>()
.unwrap_or(0);
let waiting_profit = stmt
.read("waiting_profit")
.unwrap_or("0".to_string())
.parse::<i64>()
.unwrap_or(0);
let sent_profit = stmt
.read("sent_profit")
.unwrap_or("0".to_string())
.parse::<i64>()
.unwrap_or(0);
let missed_profit = stmt
.read("missed_profit")
.unwrap_or("0".to_string())
.parse::<i64>()
.unwrap_or(0);
let unique_inputs = stmt
.read("unique_inputs")
.unwrap_or("0".to_string())
.parse::<i64>()
.unwrap_or(0);
stats.push(StatsResponse {
report_date,
chain,
totals,
waiting,
sent,
failed,
waiting_profit,
sent_profit,
missed_profit,
unique_inputs,
});
}
let stats: Vec<StatsResponse> = stats_rows
.into_iter()
.map(|row| StatsResponse {
report_date: row.report_date,
chain: row.chain,
totals: row.totals,
waiting: row.waiting,
sent: row.sent,
failed: row.failed,
waiting_profit: row.waiting_profit,
sent_profit: row.sent_profit,
missed_profit: row.missed_profit,
unique_inputs: row.unique_inputs,
})
.collect();
debug!("echo stats reply for chain: {}", netconfig.name);
HttpResponse::Ok().json(stats)
}
@@ -453,93 +384,42 @@ async fn echo_search(body: Bytes, data: web::Data<AppState>) -> impl Responder {
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("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("error");
}
};
if let Err(e) = statement.bind((1, strbody)) {
error!("Failed to bind parameter: {}", e);
return HttpResponse::InternalServerError().body("error");
}
match search_tx(&data.db, strbody).await {
Ok(Some(row)) => {
let mut response_data = HashMap::new();
response_data.insert("status", row.status);
response_data.insert("tx", row.tx);
response_data.insert("our_address", row.our_address);
response_data.insert("our_fees", row.our_fees);
response_data.insert("time", row.reqid);
if let Ok(State::Row) = statement.next() {
let mut response_data = HashMap::new();
match statement.read::<String, _>("status") {
Ok(value) => {
response_data.insert("status", value);
}
Err(e) => {
error!("Error reading status: {}", e);
match serde_json::to_string(&response_data) {
Ok(json_data) => {
debug!("echo search reply: {}", json_data);
HttpResponse::Ok().json(&response_data)
}
Err(_) => HttpResponse::BadRequest().body("error"),
}
}
match statement.read::<String, _>("tx") {
Ok(value) => {
response_data.insert("tx", value);
}
Err(e) => {
error!("Error reading tx: {}", e);
}
Ok(None) => HttpResponse::BadRequest().body("error"),
Err(e) => {
error!("Failed to search tx: {}", e);
HttpResponse::InternalServerError().body("error")
}
match statement.read::<String, _>("our_address") {
Ok(value) => {
response_data.insert("our_address", value);
}
Err(e) => {
error!("Error reading address: {}", e);
}
}
match statement.read::<String, _>("our_fees") {
Ok(value) => {
response_data.insert("our_fees", value);
}
Err(e) => {
error!("Error reading fees: {}", e);
}
}
match statement.read::<String, _>("reqid") {
Ok(value) => {
response_data.insert("time", value);
}
Err(e) => {
error!("Error reading reqid: {}", e);
}
}
match serde_json::to_string(&response_data) {
Ok(json_data) => {
debug!("echo search reply: {}", json_data);
HttpResponse::Ok().json(&response_data)
}
Err(_) => HttpResponse::BadRequest().body("error"),
}
} else {
HttpResponse::BadRequest().body("error")
}
}
/// Holds a transaction that has already been parsed and validated outside the DB lock.
#[derive(Clone)]
struct ParsedTx {
txid: String,
wtxid: String,
ntxid: String,
raw_hex: String, // the original line
raw_hex: String,
locktime: String,
inputs: Vec<(String, String)>, // (in_txid, in_vout)
outputs: Vec<(usize, String, u64)>, // (idx, script_pubkey, amount_sat)
inputs: Vec<(String, String)>,
outputs: Vec<(usize, String, u64)>,
}
/// Parse all transactions from the request body **without** needing the DB lock.
/// Skips transactions that don't have a valid willexecutor output.
fn parse_request_transactions(
strbody: &str,
_req_time: i64,
@@ -576,7 +456,6 @@ fn parse_request_transactions(
let wtxid = tx.compute_wtxid();
let locktime = tx.lock_time.to_string();
// Collect inputs
let mut inputs: Vec<(String, String)> = Vec::with_capacity(tx.input.len());
for input in tx.input {
inputs.push((
@@ -585,7 +464,6 @@ fn parse_request_transactions(
));
}
// Collect outputs and find which one is ours + its amount
let mut outputs: Vec<(usize, String, u64)> = Vec::with_capacity(tx.output.len());
let mut found = false;
let mut our_address = String::new();
@@ -601,13 +479,25 @@ fn parse_request_transactions(
netconfig.network,
) {
Ok(addr) => addr.to_string(),
Err(_) => continue, // skip un-decodable outputs
Err(_) => continue,
};
let expected_ours = if netconfig.xpub {
if known_addresses.contains(&address) {
trace!(
"output {} address {} found in known_addresses (total: {})",
idx,
&address,
known_addresses.len()
);
address.clone()
} else {
trace!(
"output {} address {} NOT in known_addresses (total: {}), skipping",
idx,
&address,
known_addresses.len()
);
continue;
}
} else {
@@ -619,6 +509,11 @@ fn parse_request_transactions(
our_fees = amount;
found = true;
trace!("address and fees are correct {}: {}", our_address, our_fees);
} else if address == expected_ours {
trace!(
"output {} address matches but amount {} < fixed_fee {}, skipping",
idx, amount, netconfig.fixed_fee
);
}
}
@@ -679,26 +574,17 @@ async fn echo_push(
};
// ===== PHASE 1: parse all transactions WITHOUT the DB lock =====
let known_addresses: HashSet<String> = {
let db = match data.db.lock() {
Ok(g) => g,
Err(_p) => {
error!("DB mutex poisoned acquiring addresses in echo_push");
let known_addresses: HashSet<String> = if netconfig.xpub {
match get_all_addresses_by_xpub(&data.db, &netconfig.address).await {
Ok(addrs) => addrs,
Err(e) => {
error!("Failed to load addresses from xpub: {}", e);
return HttpResponse::InternalServerError().body("error");
}
};
if netconfig.xpub {
match get_all_addresses_by_xpub(&db, &netconfig.address) {
Ok(addrs) => addrs,
Err(e) => {
error!("Failed to load addresses from xpub: {}", e);
return HttpResponse::InternalServerError().body("error");
}
}
} else {
HashSet::new()
}
}; // lock released here
} else {
HashSet::new()
};
// Parse all transactions (CPU-bound, no DB needed)
let parsed = parse_request_transactions(strbody, req_time, netconfig, &known_addresses);
@@ -709,134 +595,86 @@ async fn echo_push(
let all_txids: Vec<String> = parsed.iter().map(|(p, _, _)| p.txid.clone()).collect();
// ===== PHASE 2: check duplicates in a single batch query =====
let duplicates = {
let db = match data.db.lock() {
Ok(g) => g,
Err(_p) => {
error!("DB mutex poisoned in echo_push duplicate check");
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("error");
}
let duplicates = match check_duplicate_txids(&data.db, &all_txids).await {
Ok(dups) => dups,
Err(e) => {
error!("Duplicate check failed: {}", e);
return HttpResponse::InternalServerError().body("error");
}
}; // lock released here
};
let all_present = all_txids.iter().all(|t| duplicates.contains(t));
if all_present {
return HttpResponse::Ok().body("already present");
}
// ===== PHASE 3: build insert statements and execute (single DB lock, minimal time) =====
// ===== PHASE 3: build insert data and execute (single transaction, minimal time) =====
let mut tx_data = Vec::new();
let mut inp_data = Vec::new();
let mut out_data = Vec::new();
for (parsed, our_address, our_fees) in &parsed {
if duplicates.contains(&parsed.txid) {
continue;
}
tx_data.push(InsertTxData {
txid: parsed.txid.clone(),
wtxid: parsed.wtxid.clone(),
ntxid: parsed.ntxid.clone(),
raw_hex: parsed.raw_hex.clone(),
locktime: parsed.locktime.clone(),
reqid: req_time.to_string(),
network: netconfig.name.clone(),
our_address: our_address.clone(),
our_fees: our_fees.to_string(),
});
for (in_txid, in_vout) in &parsed.inputs {
inp_data.push(InsertInpData {
txid: parsed.txid.clone(),
in_txid: in_txid.clone(),
in_vout: in_vout.clone(),
});
}
for (idx, script, amount) in &parsed.outputs {
out_data.push(InsertOutData {
txid: parsed.txid.clone(),
vout: i64::try_from(*idx).unwrap_or(-1),
script_pubkey: script.clone(),
amount: i64::try_from(*amount).unwrap_or(0),
});
}
}
if tx_data.is_empty() {
return HttpResponse::Ok().body("already present");
}
if let Err(err) = bal_server::db::execute_insert(&data.db, &tx_data, &inp_data, &out_data).await
{
let db = match data.db.lock() {
Ok(g) => g,
Err(_p) => {
error!("DB mutex poisoned in echo_push insert phase");
return HttpResponse::InternalServerError().body("error");
}
};
let sqltxshead = "INSERT INTO tbl_tx (txid, wtxid, ntxid, tx, locktime, reqid, network, our_address, our_fees)".to_string();
let mut sqltxs = String::new();
let sqlinpshead = "INSERT INTO tbl_inp (txid, in_txid, in_vout )".to_string();
let mut sqlinps = String::new();
let sqloutshead = "INSERT INTO tbl_out (txid, vout, script_pubkey, amount )".to_string();
let mut sqlouts = String::new();
let mut union_tx = true;
let mut union_inps = true;
let mut union_outs = true;
let mut ptx: Vec<(usize, Value)> = vec![];
let mut pinps: Vec<(usize, Value)> = vec![];
let mut pouts: Vec<(usize, Value)> = vec![];
let mut linenum = 1usize;
let mut lineinp = 1usize;
let mut lineout = 1usize;
for (parsed, our_address, our_fees) in &parsed {
if duplicates.contains(&parsed.txid) {
continue;
}
if !union_tx {
sqltxs.push_str(" UNION ALL");
} else {
union_tx = false;
}
sqltxs.push_str(" SELECT ?, ?, ?, ?, ?, ?, ?, ?, ?");
ptx.push((linenum, Value::String(parsed.txid.clone())));
ptx.push((linenum + 1, Value::String(parsed.wtxid.clone())));
ptx.push((linenum + 2, Value::String(parsed.ntxid.clone())));
ptx.push((linenum + 3, Value::String(parsed.raw_hex.clone())));
ptx.push((linenum + 4, Value::String(parsed.locktime.clone())));
ptx.push((linenum + 5, Value::String(req_time.to_string())));
ptx.push((linenum + 6, Value::String(netconfig.name.clone())));
ptx.push((linenum + 7, Value::String(our_address.clone())));
ptx.push((linenum + 8, Value::String(our_fees.to_string())));
linenum += 9;
for (in_txid, in_vout) in &parsed.inputs {
if !union_inps {
sqlinps.push_str(" UNION ALL");
} else {
union_inps = false;
}
sqlinps.push_str(" SELECT ?, ?, ?");
pinps.push((lineinp, Value::String(parsed.txid.clone())));
pinps.push((lineinp + 1, Value::String(in_txid.clone())));
pinps.push((lineinp + 2, Value::String(in_vout.clone())));
lineinp += 3;
}
for (idx, script, amount) in &parsed.outputs {
if !union_outs {
sqlouts.push_str(" UNION ALL");
} else {
union_outs = false;
}
sqlouts.push_str(" SELECT ?, ?, ?, ?");
pouts.push((lineout, Value::String(parsed.txid.clone())));
pouts.push((
lineout + 1,
Value::Integer(i64::try_from(*idx).unwrap_or(-1)),
));
pouts.push((lineout + 2, Value::String(script.clone())));
pouts.push((
lineout + 3,
Value::Integer(i64::try_from(*amount).unwrap_or(0)),
));
lineout += 4;
}
}
if sqltxs.is_empty() {
return HttpResponse::Ok().body("already present");
}
let sqltxs = format!("{}{};", sqltxshead, sqltxs);
let sqlinps = format!("{}{};", sqlinpshead, sqlinps);
let sqlouts = format!("{}{};", sqloutshead, sqlouts);
if let Err(err) = execute_insert(&db, sqltxs, ptx, sqlinps, pinps, sqlouts, pouts) {
error!("execute_insert failed: {}", err);
return HttpResponse::BadRequest().body("error");
}
} // lock released
error!("execute_insert failed: {}", err);
return HttpResponse::BadRequest().body("error");
}
HttpResponse::Ok().body("thx")
}
fn parse_env(data: &MyConfig) -> MyConfig {
let mut cfg = data.clone();
if let Ok(value) = env::var("BAL_SERVER_DB_BACKEND") {
debug!("BAL_SERVER_DB_BACKEND: {}", value);
cfg.db_backend = value;
}
if let Ok(value) = env::var("BAL_SERVER_DB_FILE") {
debug!("BAL_SERVER_DB_FILE: {}", value);
cfg.db_file = value;
}
if let Ok(value) = env::var("BAL_SERVER_PG_DSN") {
debug!("BAL_SERVER_PG_DSN: {}", value);
cfg.pg_dsn = value;
}
if let Ok(value) = env::var("BAL_SERVER_BIND_ADDRESS") {
debug!("BAL_SERVER_BIND_ADDRESS: {}", value);
cfg.bind_address = value;
@@ -889,10 +727,10 @@ fn parse_env_netconfig(cfg: &mut MyConfig, chain: &str) {
}
}
fn init_network(db: &Connection, cfg: &MyConfig) {
async fn init_network(pool: &DatabasePool, cfg: &MyConfig) {
for network in NETWORKS {
let netconfig = cfg.get_net_config(network);
insert_xpub(db, &netconfig.name.to_string(), &netconfig.address);
insert_xpub(pool, &netconfig.name.to_string(), &netconfig.address).await;
}
}
@@ -903,32 +741,44 @@ async fn main() -> std::io::Result<()> {
let actix_cfg = parse_actix_config();
let cfg = parse_env(&cfg);
let db = match open_db(&cfg.db_file) {
Ok(c) => c,
let connection_string = match cfg.db_backend.as_str() {
"sqlite" => cfg.db_file.clone(),
"postgresql" => cfg.pg_dsn.clone(),
other => {
return Err(std::io::Error::other(format!(
"Unknown DB backend: {}",
other
)));
}
};
let db = match open_database(&cfg.db_backend, &connection_string).await {
Ok(pool) => pool,
Err(e) => {
return Err(std::io::Error::other(e));
}
};
// Create database tables
create_database(&db);
create_database(&db)
.await
.map_err(|e| std::io::Error::other(format!("Failed to create database: {}", e)))?;
// Initialize networks
init_network(&db, &cfg);
init_network(&db, &cfg).await;
let data = web::Data::new(AppState {
db: Mutex::new(db),
db,
cfg: cfg.clone(),
});
let bind_address = data.cfg.bind_address.clone();
let bind_port = data.cfg.bind_port;
// Use a single global rate limiter with the most conservative settings (1 req/sec)
// Per-endpoint rate limiting requires advanced configuration with explicit types
let governor_conf = GovernorConfigBuilder::const_default()
.seconds_per_request(actix_cfg.rate_limit_pushtxs.0) // Most restrictive: 1 req/sec
.burst_size(actix_cfg.rate_limit_pushtxs.1) // Burst: 3
.seconds_per_request(actix_cfg.rate_limit_pushtxs.0)
.burst_size(actix_cfg.rate_limit_pushtxs.1)
.finish()
.unwrap();

412
src/db.rs
View File

@@ -1,412 +0,0 @@
use log::{error, info, trace, warn};
use sqlite::{Connection, Error, State, Value};
use std::collections::HashSet;
use std::path::Path;
use std::thread;
use std::time::Duration;
/// Check which txids are already present in the database in a single batch query.
/// Returns a HashSet of txids that already exist (duplicates).
/// This is O(1) per query regardless of the number of txids, replacing the N+1 pattern.
pub fn check_duplicate_txids(db: &Connection, txids: &[String]) -> Result<HashSet<String>, Error> {
if txids.is_empty() {
return Ok(HashSet::new());
}
// Build a single query with all txids using IN clause placeholders
// SQLite supports up to 1000 parameters per statement, so we chunk for safety
let mut duplicates = HashSet::new();
let chunk_size = 500; // Safe chunk size for SQLite parameters
for chunk in txids.chunks(chunk_size) {
let placeholders = chunk.iter().map(|_| "?").collect::<Vec<_>>().join(",");
let sql = format!("SELECT txid FROM tbl_tx WHERE txid IN ({})", placeholders);
let mut stmt = db.prepare(sql)?;
for (i, txid) in chunk.iter().enumerate() {
stmt.bind((i + 1, Value::String(txid.clone())))?;
}
while let Ok(State::Row) = stmt.next() {
if let Ok(txid) = stmt.read::<String, _>("txid") {
duplicates.insert(txid);
}
}
}
Ok(duplicates)
}
/// Validates and opens the SQLite database, enforcing security best practices:
/// - Path must not contain `..` (directory traversal).
/// - Absolute paths must not target known system directories.
/// - If the file exists, it must be a regular file (not a symlink or device).
/// - WAL journal mode is enabled for safe concurrent access.
/// - Synchronous is set to NORMAL for performance with safety.
///
/// Returns `Err` on validation failure or open error to prevent panics.
pub fn open_db(path: &str) -> Result<Connection, String> {
let p = Path::new(path);
// Prevent directory traversal
for component in p.components() {
if component == std::path::Component::ParentDir {
return Err("Database path may not contain '..'".to_string());
}
}
// If absolute, block known sensitive system directories
if p.is_absolute() {
let path_str = p.to_str().unwrap_or("");
let forbidden = [
"/etc", "/proc", "/sys", "/dev", "/usr", "/bin", "/sbin", "/lib", "/opt",
];
for prefix in &forbidden {
if path_str.starts_with(prefix) {
return Err(format!(
"Absolute database path under {} is forbidden",
prefix
));
}
}
}
// If file exists, must be a regular file (not a symlink, device, etc.)
if p.exists() {
if p.is_symlink() {
return Err("Database path must not be a symlink".to_string());
}
let metadata = std::fs::metadata(p)
.map_err(|e| format!("Cannot access database file metadata: {}", e))?;
if !metadata.is_file() {
return Err(
"Database path must point to a regular file, not a directory or device".to_string(),
);
}
}
let conn = sqlite::open(path).map_err(|e| format!("Failed to open SQLite database: {}", e))?;
// Set busy timeout BEFORE WAL mode so SQLite waits instead of failing immediately.
// This handles the race where two processes (server + pusher) open the same DB
// and both try to enable WAL mode concurrently.
conn.execute("PRAGMA busy_timeout = 5000;")
.map_err(|e| format!("Failed to set busy_timeout: {}", e))?;
// Retry WAL mode up to 5 times (handles concurrent open from bal-pusher).
let mut wal_ok = false;
for attempt in 0..5 {
match conn.execute("PRAGMA journal_mode = WAL;") {
Ok(_) => {
wal_ok = true;
break;
}
Err(e) => {
warn!(
"WAL mode attempt {}/5 failed: {}, retrying in 100ms...",
attempt + 1,
e
);
thread::sleep(Duration::from_millis(100));
}
}
}
if !wal_ok {
// WAL might already be enabled by another process; this is not fatal.
warn!("Could not set WAL mode after retries — may already be enabled by another process");
}
conn.execute("PRAGMA synchronous = NORMAL;")
.map_err(|e| format!("Failed to set synchronous NORMAL: {}", e))?;
Ok(conn)
}
/// Loads all known addresses for a given xpub into a HashSet for fast
/// in-memory lookup during transaction validation (replaces N+1 query).
pub fn get_all_addresses_by_xpub(db: &Connection, xpub: &str) -> Result<HashSet<String>, Error> {
let mut stmt = db.prepare(
"SELECT a.address FROM tbl_address a JOIN tbl_xpub x ON a.xpub = x.id WHERE x.xpub = ?",
)?;
stmt.bind((1, Value::String(xpub.to_string())))?;
let mut addresses = HashSet::new();
while let Ok(State::Row) = stmt.next() {
match stmt.read::<String, _>("address") {
Ok(addr) => {
addresses.insert(addr);
}
Err(e) => {
error!("Failed to read address column: {}", e);
}
}
}
Ok(addresses)
}
pub fn create_database(db: &Connection) {
info!("database sanity check");
let _ = db.execute("CREATE TABLE IF NOT EXISTS tbl_tx (txid PRIMARY KEY, date_creation TIMESTAMP DEFAULT CURRENT_TIMESTAMP, date_update TIMESTAMP DEFAULT CURRENT_TIMESTAMP, wtxid, ntxid, tx, locktime integer, network, network_fees, reqid, our_fees, our_address, status integer DEFAULT 0);");
let _ = db.execute("ALTER TABLE tbl_tx ADD COLUMN push_err TEXT");
let _ = db.execute("CREATE TABLE IF NOT EXISTS tbl_inp(id, txid, in_txid, in_vout);");
let _ = db.execute("CREATE UNIQUE INDEX ON tbl_inp(txid,in_txid,in_vout);");
let _ =
db.execute("CREATE TABLE IF NOT EXISTS tbl_out(id, txid, script_pubkey, amount, vout);");
let _ = db.execute("CREATE UNIQUE INDEX ON tbl_out(txid, script_pubkey, amount, vout);");
let _ = db.execute("CREATE TABLE IF NOT EXISTS tbl_xpub (id INTEGER PRIMARY KEY , network TEXT, xpub TEXT, date_create TIMESTAMP DEFAULT CURRENT_TIMESTAMP,path_idx INTEGER DEFAULT -1);");
let _ = db.execute("CREATE UNIQUE INDEX idx_xpub ON tbl_xpub (network, xpub)");
let _ = db.execute("CREATE TABLE IF NOT EXISTS tbl_address (address TEXT PRIMARY_KEY, path TEXT NOT NULL, date_create TIMESTAMP DEFAULT CURRENT_TIMESTAMP, xpub INTEGER,remote_address TEXT);");
let _ = db.execute("CREATE TABLE IF NOT EXISTS tbl_stats (report_date TEXT, chain TEXT, totals INTEGER, waiting INTEGER, sent INTEGER, failed INTEGER, waiting_profit INTEGER, sent_profit INTEGER, missed_profit INTEGER, unique_inputs INTEGER);");
// UNIQUE index required for ON CONFLICT(chain) DO UPDATE in calculate_stats
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';");
}
/*
pub fn get_xpub_id(db: &Connection, network: &String, xpub: &String) -> Option<i64>{
let mut stmt = db.prepare("SELECT * FROM tbl_xpub where network = ? and xpub = ?;").unwrap();
let _ = stmt.bind((1,Value::String(network.to_string()))).unwrap();
let _ = stmt.bind((2,Value::String(xpub.to_string()))).unwrap();
if let Ok(State::Row) = stmt.next(){
return Some(stmt.read::<i64, _>("id").unwrap());
} else {
return None;
}
}
*/
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 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;
}
if let Err(e) = stmt.bind((2, Value::String(xpub.to_string()))) {
error!("Failed to bind xpub parameter: {}", e);
return;
}
if let Err(e) = stmt.next() {
error!("Failed to insert xpub: {}", e);
}
}
}
pub fn get_last_used_address_by_ip(
db: &Connection,
network: &String,
xpub: &String,
address: &String,
) -> Option<String> {
let mut stmt = match db.prepare("SELECT tbl_address.address FROM tbl_xpub join tbl_address on(tbl_xpub.id = tbl_address.xpub) where tbl_xpub.network = ? and tbl_address.remote_address = ? and tbl_xpub.xpub = ? ORDER BY tbl_address.date_create DESC LIMIT 1;") {
Ok(s) => s,
Err(e) => {
error!("Failed to prepare address query: {}", e);
return None;
}
};
if let Err(e) = stmt.bind((1, Value::String(network.to_string()))) {
error!("Failed to bind network parameter: {}", e);
return None;
}
if let Err(e) = stmt.bind((2, Value::String(address.to_string()))) {
error!("Failed to bind address parameter: {}", e);
return None;
}
if let Err(e) = stmt.bind((3, Value::String(xpub.to_string()))) {
error!("Failed to bind xpub parameter: {}", e);
return None;
}
if let Ok(State::Row) = stmt.next() {
match stmt.read::<String, _>("address") {
Ok(addr) => Some(addr),
Err(e) => {
error!("Failed to read address column: {}", e);
None
}
}
} else {
None
}
}
pub fn get_next_address_index(db: &Connection, network: &String, xpub: &String) -> (i64, i64) {
let mut stmt = match db.prepare("UPDATE tbl_xpub SET path_idx = path_idx + 1 WHERE network = ? and xpub= ? RETURNING path_idx,id;") {
Ok(s) => s,
Err(e) => {
error!("Failed to prepare xpub index update: {}", e);
return (0, 0);
}
};
if let Err(e) = stmt.bind((1, Value::String(network.to_string()))) {
error!("Failed to bind network parameter: {}", e);
return (0, 0);
}
if let Err(e) = stmt.bind((2, Value::String(xpub.to_string()))) {
error!("Failed to bind xpub parameter: {}", e);
return (0, 0);
}
match stmt.next() {
Ok(State::Row) => match stmt.read::<i64, _>("path_idx") {
Ok(next) => match stmt.read::<i64, _>("id") {
Ok(id) => (id, next),
Err(e) => {
error!("Failed to read id column: {}", e);
(0, 0)
}
},
Err(e) => {
error!("Failed to read path_idx column: {}", e);
(0, 0)
}
},
Err(e) => {
error!("Failed to execute xpub index update: {}", e);
(0, 0)
}
Ok(State::Done) => (0, 0),
}
}
pub fn save_new_address(
db: &Connection,
xpub: i64,
address: &String,
path: &String,
remote_addr: &String,
) {
let mut stmt = match db.prepare(
"INSERT INTO tbl_address(address,path,xpub,remote_address) VALUES(?,?,?,?);
",
) {
Ok(s) => s,
Err(e) => {
error!("Failed to prepare address insert statement: {}", e);
return;
}
};
if let Err(e) = stmt.bind((1, Value::String(address.to_string()))) {
error!("Failed to bind address parameter: {}", e);
return;
}
if let Err(e) = stmt.bind((2, Value::String(path.to_string()))) {
error!("Failed to bind path parameter: {}", e);
return;
}
if let Err(e) = stmt.bind((3, Value::Integer(xpub))) {
error!("Failed to bind xpub parameter: {}", e);
return;
}
if let Err(e) = stmt.bind((4, Value::String(remote_addr.to_string()))) {
error!("Failed to bind remote_addr parameter: {}", e);
return;
}
if let Err(e) = stmt.next() {
error!("Failed to insert address: {}", e);
}
}
pub fn execute_insert(
db: &Connection,
sqltxs: String,
ptx: Vec<(usize, Value)>,
sqlinp: String,
pinp: Vec<(usize, Value)>,
sqlout: String,
pout: Vec<(usize, Value)>,
) -> Result<(), Error> {
let _ = db.execute("BEGIN TRANSACTION");
let mut stmt = match db.prepare(sqltxs.as_str()) {
Ok(s) => s,
Err(err) => {
error!("error preparing sqltxs: {}", err);
let _ = db.execute("ROLLBACK");
return Err(err);
}
};
if let Err(err) = stmt.bind::<&[(_, Value)]>(&ptx[..]) {
error!("error binding transaction parameters: {}", err);
let _ = db.execute("ROLLBACK");
return Err(err);
}
if let Err(err) = stmt.next() {
error!("error inserting transactions {}", err);
let _ = db.execute("ROLLBACK");
} else {
let mut stmt = match db.prepare(sqlinp.as_str()) {
Ok(s) => s,
Err(err) => {
error!("error preparing sqlinp: {}", err);
let _ = db.execute("ROLLBACK");
return Err(err);
}
};
if let Err(err) = stmt.bind::<&[(_, Value)]>(&pinp[..]) {
error!("error binding inputs parameters {}", err);
let _ = db.execute("ROLLBACK");
return Err(err);
}
if let Err(err) = stmt.next() {
error!("error inserting inputs {}", err);
let _ = db.execute("ROLLBACK");
return Err(err);
} else {
let mut stmt = match db.prepare(sqlout.as_str()) {
Ok(s) => s,
Err(err) => {
error!("error preparing sqlout: {}", err);
let _ = db.execute("ROLLBACK");
return Err(err);
}
};
if let Err(err) = stmt.bind::<&[(_, Value)]>(&pout[..]) {
error!("error binding outs parameters {}", err);
let _ = db.execute("ROLLBACK");
return Err(err);
}
if let Err(err) = stmt.next() {
error!("error inserting outs {}", err);
let _ = db.execute("ROLLBACK");
return Err(err);
}
}
}
let _ = db.execute("COMMIT");
Ok(())
}
pub fn get_total_transaction_number(db: Connection, network: &String) -> Result<i64, Error> {
let mut stmt = db
.prepare("SELECT COUNT(*) as total_number FROM tbl_tx where network = ?;")
.map_err(|e| {
error!("Failed to prepare statement: {}", e);
e
})?;
if let Err(e) = stmt.bind((1, Value::String(network.to_string()))) {
error!("Failed to bind network parameter: {}", e);
return Err(e);
}
match stmt.next() {
Ok(State::Row) => match stmt.read::<i64, _>("total_number") {
Ok(val) => Ok(val),
Err(e) => {
error!("Failed to read total_number column: {}", e);
Err(e)
}
},
Ok(sqlite::State::Done) => Ok(0),
Err(err) => {
error!("Failed to execute query: {}", err);
Err(err)
}
}
}

868
src/db/mod.rs Normal file
View File

@@ -0,0 +1,868 @@
pub mod schema;
use log::{error, info, trace};
use sqlx::Row;
use std::collections::HashSet;
use std::path::Path;
#[derive(Clone)]
pub enum DatabasePool {
SQLite(sqlx::SqlitePool),
PostgreSQL(sqlx::PgPool),
}
fn validate_sqlite_path(dsn: &str) -> Result<(), String> {
// Extract file path from DSN formats like "sqlite:path" or "file:path?mode=rwc"
let path_str = if let Some(rest) = dsn.strip_prefix("sqlite:") {
rest.split('?').next().unwrap_or(rest)
} else if let Some(rest) = dsn.strip_prefix("file:") {
rest.split('?').next().unwrap_or(rest)
} else {
dsn
};
// Skip validation for in-memory databases
if path_str == ":memory:" || path_str.is_empty() {
return Ok(());
}
let p = Path::new(path_str);
// Prevent directory traversal
for component in p.components() {
if component == std::path::Component::ParentDir {
return Err("Database path may not contain '..'".to_string());
}
}
// If absolute, block known sensitive system directories
if p.is_absolute() {
let forbidden = [
"/etc", "/proc", "/sys", "/dev", "/usr", "/bin", "/sbin", "/lib", "/opt",
];
for prefix in &forbidden {
if path_str.starts_with(prefix) {
return Err(format!(
"Absolute database path under {} is forbidden",
prefix
));
}
}
}
// If file exists, must be a regular file (not a symlink, device, etc.)
if p.exists() {
if p.is_symlink() {
return Err("Database path must not be a symlink".to_string());
}
let metadata = std::fs::metadata(p)
.map_err(|e| format!("Cannot access database file metadata: {}", e))?;
if !metadata.is_file() {
return Err(
"Database path must point to a regular file, not a directory or device".to_string(),
);
}
}
Ok(())
}
pub async fn open_database(backend: &str, connection_string: &str) -> Result<DatabasePool, String> {
match backend {
"sqlite" => {
validate_sqlite_path(connection_string)?;
let pool = sqlx::sqlite::SqlitePoolOptions::new()
.max_connections(1)
.connect(connection_string)
.await
.map_err(|e| format!("Failed to open SQLite database: {}", e))?;
sqlx::query("PRAGMA journal_mode=WAL")
.execute(&pool)
.await
.map_err(|e| format!("Failed to set WAL mode: {}", e))?;
sqlx::query("PRAGMA busy_timeout=5000")
.execute(&pool)
.await
.map_err(|e| format!("Failed to set busy_timeout: {}", e))?;
Ok(DatabasePool::SQLite(pool))
}
"postgresql" => {
let pool = sqlx::postgres::PgPoolOptions::new()
.max_connections(10)
.connect(connection_string)
.await
.map_err(|e| format!("Failed to open PostgreSQL database: {}", e))?;
Ok(DatabasePool::PostgreSQL(pool))
}
other => Err(format!("Unknown database backend: {}", other)),
}
}
pub async fn create_database(pool: &DatabasePool) -> Result<(), sqlx::Error> {
info!("database sanity check");
match pool {
DatabasePool::SQLite(p) => schema::create_sqlite_schema(p).await,
DatabasePool::PostgreSQL(p) => schema::create_pg_schema(p).await,
}
}
pub async fn check_duplicate_txids(
pool: &DatabasePool,
txids: &[String],
) -> Result<HashSet<String>, sqlx::Error> {
if txids.is_empty() {
return Ok(HashSet::new());
}
let mut duplicates = HashSet::new();
let chunk_size = 500;
for chunk in txids.chunks(chunk_size) {
match pool {
DatabasePool::SQLite(p) => {
let placeholders: Vec<String> = chunk.iter().map(|_| "?".to_string()).collect();
let sql = format!(
"SELECT txid FROM tbl_tx WHERE txid IN ({})",
placeholders.join(",")
);
let mut query = sqlx::query(&sql);
for txid in chunk {
query = query.bind(txid);
}
let rows = query.fetch_all(p).await?;
for row in rows {
if let Ok(txid) = row.try_get::<String, _>("txid") {
duplicates.insert(txid);
}
}
}
DatabasePool::PostgreSQL(p) => {
let placeholders: Vec<String> = chunk
.iter()
.enumerate()
.map(|(i, _)| format!("${}", i + 1))
.collect();
let sql = format!(
"SELECT txid FROM tbl_tx WHERE txid IN ({})",
placeholders.join(",")
);
let mut query = sqlx::query(&sql);
for txid in chunk {
query = query.bind(txid);
}
let rows = query.fetch_all(p).await?;
for row in rows {
if let Ok(txid) = row.try_get::<String, _>("txid") {
duplicates.insert(txid);
}
}
}
}
}
Ok(duplicates)
}
pub async fn get_all_addresses_by_xpub(
pool: &DatabasePool,
xpub: &str,
) -> Result<HashSet<String>, sqlx::Error> {
let mut addresses = HashSet::new();
trace!("get_all_addresses_by_xpub: querying for xpub={}", xpub);
match pool {
DatabasePool::SQLite(p) => {
let rows = sqlx::query(
"SELECT a.address FROM tbl_address a JOIN tbl_xpub x ON a.xpub = x.id WHERE x.xpub = ?",
)
.bind(xpub)
.fetch_all(p)
.await?;
trace!(
"get_all_addresses_by_xpub: SQLite returned {} rows",
rows.len()
);
for row in rows {
if let Ok(addr) = row.try_get::<String, _>("address") {
trace!("get_all_addresses_by_xpub: address={}", &addr);
addresses.insert(addr);
}
}
}
DatabasePool::PostgreSQL(p) => {
let rows = sqlx::query(
"SELECT a.address FROM tbl_address a JOIN tbl_xpub x ON a.xpub = x.id WHERE x.xpub = $1",
)
.bind(xpub)
.fetch_all(p)
.await?;
trace!(
"get_all_addresses_by_xpub: PostgreSQL returned {} rows",
rows.len()
);
for row in rows {
if let Ok(addr) = row.try_get::<String, _>("address") {
trace!("get_all_addresses_by_xpub: address={}", &addr);
addresses.insert(addr);
}
}
}
}
trace!(
"get_all_addresses_by_xpub: returning {} addresses",
addresses.len()
);
Ok(addresses)
}
pub async fn insert_xpub(pool: &DatabasePool, network: &str, xpub: &str) {
if xpub.is_empty() {
return;
}
trace!("going to insert: {} xpub:{}", network, xpub);
match pool {
DatabasePool::SQLite(p) => {
if let Err(e) =
sqlx::query("INSERT OR IGNORE INTO tbl_xpub(network, xpub) VALUES(?, ?)")
.bind(network)
.bind(xpub)
.execute(p)
.await
{
error!("Failed to insert xpub: {}", e);
}
}
DatabasePool::PostgreSQL(p) => {
if let Err(e) = sqlx::query(
"INSERT INTO tbl_xpub(network, xpub) VALUES($1, $2) ON CONFLICT DO NOTHING",
)
.bind(network)
.bind(xpub)
.execute(p)
.await
{
error!("Failed to insert xpub: {}", e);
}
}
}
}
pub async fn get_last_used_address_by_ip(
pool: &DatabasePool,
network: &str,
xpub: &str,
address: &str,
) -> Option<String> {
match pool {
DatabasePool::SQLite(p) => {
let result = sqlx::query(
"SELECT tbl_address.address FROM tbl_xpub JOIN tbl_address ON(tbl_xpub.id = tbl_address.xpub) WHERE tbl_xpub.network = ? AND tbl_address.remote_address = ? AND tbl_xpub.xpub = ? ORDER BY tbl_address.date_create DESC LIMIT 1",
)
.bind(network)
.bind(address)
.bind(xpub)
.fetch_optional(p)
.await;
match result {
Ok(Some(row)) => row.try_get::<String, _>("address").ok(),
Ok(None) => None,
Err(e) => {
error!("Failed to query last used address: {}", e);
None
}
}
}
DatabasePool::PostgreSQL(p) => {
let result = sqlx::query(
"SELECT tbl_address.address FROM tbl_xpub JOIN tbl_address ON(tbl_xpub.id = tbl_address.xpub) WHERE tbl_xpub.network = $1 AND tbl_address.remote_address = $2 AND tbl_xpub.xpub = $3 ORDER BY tbl_address.date_create DESC LIMIT 1",
)
.bind(network)
.bind(address)
.bind(xpub)
.fetch_optional(p)
.await;
match result {
Ok(Some(row)) => row.try_get::<String, _>("address").ok(),
Ok(None) => None,
Err(e) => {
error!("Failed to query last used address: {}", e);
None
}
}
}
}
}
pub async fn get_next_address_index(pool: &DatabasePool, network: &str, xpub: &str) -> (i64, i64) {
match pool {
DatabasePool::SQLite(p) => {
let result = sqlx::query(
"UPDATE tbl_xpub SET path_idx = path_idx + 1 WHERE network = ? AND xpub = ? RETURNING path_idx, id",
)
.bind(network)
.bind(xpub)
.fetch_optional(p)
.await;
match result {
Ok(Some(row)) => {
let idx = row.try_get::<i64, _>("path_idx").unwrap_or(0);
let id = row.try_get::<i64, _>("id").unwrap_or(0);
(id, idx)
}
Ok(None) => (0, 0),
Err(e) => {
error!("Failed to get next address index: {}", e);
(0, 0)
}
}
}
DatabasePool::PostgreSQL(p) => {
let result = sqlx::query(
"UPDATE tbl_xpub SET path_idx = path_idx + 1 WHERE network = $1 AND xpub = $2 RETURNING path_idx, id",
)
.bind(network)
.bind(xpub)
.fetch_optional(p)
.await;
match result {
Ok(Some(row)) => {
let idx = row.try_get::<i32, _>("path_idx").unwrap_or(0) as i64;
let id = row.try_get::<i32, _>("id").unwrap_or(0) as i64;
(id, idx)
}
Ok(None) => (0, 0),
Err(e) => {
error!("Failed to get next address index: {}", e);
(0, 0)
}
}
}
}
}
pub async fn save_new_address(
pool: &DatabasePool,
xpub: i64,
address: &str,
path: &str,
remote_addr: &str,
) {
match pool {
DatabasePool::SQLite(p) => {
if let Err(e) = sqlx::query(
"INSERT INTO tbl_address(address, path, xpub, remote_address) VALUES(?, ?, ?, ?)",
)
.bind(address)
.bind(path)
.bind(xpub)
.bind(remote_addr)
.execute(p)
.await
{
error!("Failed to save address: {}", e);
}
}
DatabasePool::PostgreSQL(p) => {
if let Err(e) = sqlx::query(
"INSERT INTO tbl_address(address, path, xpub, remote_address) VALUES($1, $2, $3, $4)",
)
.bind(address)
.bind(path)
.bind(xpub as i32)
.bind(remote_addr)
.execute(p)
.await
{
error!("Failed to save address: {}", e);
}
}
}
}
#[derive(Debug, Clone)]
pub struct InsertTxData {
pub txid: String,
pub wtxid: String,
pub ntxid: String,
pub raw_hex: String,
pub locktime: String,
pub reqid: String,
pub network: String,
pub our_address: String,
pub our_fees: String,
}
#[derive(Debug, Clone)]
pub struct InsertInpData {
pub txid: String,
pub in_txid: String,
pub in_vout: String,
}
#[derive(Debug, Clone)]
pub struct InsertOutData {
pub txid: String,
pub vout: i64,
pub script_pubkey: String,
pub amount: i64,
}
pub async fn execute_insert(
pool: &DatabasePool,
txs: &[InsertTxData],
inps: &[InsertInpData],
outs: &[InsertOutData],
) -> Result<(), sqlx::Error> {
match pool {
DatabasePool::SQLite(p) => {
let mut tx = p.begin().await?;
for item in txs {
sqlx::query(
"INSERT INTO tbl_tx (txid, wtxid, ntxid, tx, locktime, reqid, network, our_address, our_fees) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?)"
)
.bind(&item.txid)
.bind(&item.wtxid)
.bind(&item.ntxid)
.bind(&item.raw_hex)
.bind(&item.locktime)
.bind(&item.reqid)
.bind(&item.network)
.bind(&item.our_address)
.bind(&item.our_fees)
.execute(&mut *tx)
.await?;
}
for item in inps {
sqlx::query("INSERT INTO tbl_inp (txid, in_txid, in_vout) VALUES (?, ?, ?)")
.bind(&item.txid)
.bind(&item.in_txid)
.bind(&item.in_vout)
.execute(&mut *tx)
.await?;
}
for item in outs {
sqlx::query(
"INSERT INTO tbl_out (txid, vout, script_pubkey, amount) VALUES (?, ?, ?, ?)",
)
.bind(&item.txid)
.bind(item.vout)
.bind(&item.script_pubkey)
.bind(item.amount)
.execute(&mut *tx)
.await?;
}
tx.commit().await?;
}
DatabasePool::PostgreSQL(p) => {
let mut tx = p.begin().await?;
for item in txs {
let locktime_i64: i64 = item.locktime.parse().unwrap_or(0);
sqlx::query(
"INSERT INTO tbl_tx (txid, wtxid, ntxid, tx, locktime, reqid, network, our_address, our_fees) VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9)"
)
.bind(&item.txid)
.bind(&item.wtxid)
.bind(&item.ntxid)
.bind(&item.raw_hex)
.bind(locktime_i64)
.bind(&item.reqid)
.bind(&item.network)
.bind(&item.our_address)
.bind(&item.our_fees)
.execute(&mut *tx)
.await?;
}
for item in inps {
let in_vout_i32: i32 = item.in_vout.parse().unwrap_or(0);
sqlx::query("INSERT INTO tbl_inp (txid, in_txid, in_vout) VALUES ($1, $2, $3)")
.bind(&item.txid)
.bind(&item.in_txid)
.bind(in_vout_i32)
.execute(&mut *tx)
.await?;
}
for item in outs {
sqlx::query(
"INSERT INTO tbl_out (txid, vout, script_pubkey, amount) VALUES ($1, $2, $3, $4)",
)
.bind(&item.txid)
.bind(item.vout as i32)
.bind(&item.script_pubkey)
.bind(item.amount.to_string())
.execute(&mut *tx)
.await?;
}
tx.commit().await?;
}
}
Ok(())
}
pub async fn get_total_transaction_number(
pool: &DatabasePool,
network: &str,
) -> Result<i64, sqlx::Error> {
match pool {
DatabasePool::SQLite(p) => {
let row = sqlx::query("SELECT COUNT(*) as total_number FROM tbl_tx WHERE network = ?")
.bind(network)
.fetch_one(p)
.await?;
Ok(row.try_get::<i64, _>("total_number").unwrap_or(0))
}
DatabasePool::PostgreSQL(p) => {
let row = sqlx::query("SELECT COUNT(*) as total_number FROM tbl_tx WHERE network = $1")
.bind(network)
.fetch_one(p)
.await?;
Ok(row.try_get::<i64, _>("total_number").unwrap_or(0))
}
}
}
pub async fn update_tx_status(
pool: &DatabasePool,
txid: &str,
status: i32,
push_err: Option<&str>,
) -> Result<(), sqlx::Error> {
match pool {
DatabasePool::SQLite(p) => {
if let Some(err) = push_err {
sqlx::query("UPDATE tbl_tx SET status = ?, push_err = ? WHERE txid = ?")
.bind(status)
.bind(err)
.bind(txid)
.execute(p)
.await?;
} else {
sqlx::query("UPDATE tbl_tx SET status = ? WHERE txid = ?")
.bind(status)
.bind(txid)
.execute(p)
.await?;
}
}
DatabasePool::PostgreSQL(p) => {
if let Some(err) = push_err {
sqlx::query("UPDATE tbl_tx SET status = $1, push_err = $2 WHERE txid = $3")
.bind(status)
.bind(err)
.bind(txid)
.execute(p)
.await?;
} else {
sqlx::query("UPDATE tbl_tx SET status = $1 WHERE txid = $2")
.bind(status)
.bind(txid)
.execute(p)
.await?;
}
}
}
Ok(())
}
#[derive(Debug, Clone)]
pub struct TxRow {
pub txid: String,
pub tx: String,
pub locktime: i64,
pub network: String,
pub status: i64,
}
pub async fn get_pending_txs(
pool: &DatabasePool,
network: &str,
locktime_threshold: i64,
bestblock_height: i64,
bestblock_time: i64,
) -> Result<Vec<TxRow>, sqlx::Error> {
let mut results = Vec::new();
match pool {
DatabasePool::SQLite(p) => {
let rows = sqlx::query(
"SELECT txid, tx, locktime, network, status FROM tbl_tx WHERE network = ? AND status = 0 AND (locktime < ? OR (locktime > ? AND locktime < ?))",
)
.bind(network)
.bind(bestblock_height)
.bind(locktime_threshold)
.bind(bestblock_time)
.fetch_all(p)
.await?;
for row in rows {
results.push(TxRow {
txid: row.try_get("txid")?,
tx: row.try_get("tx")?,
locktime: row.try_get("locktime")?,
network: row.try_get("network")?,
status: row.try_get("status")?,
});
}
}
DatabasePool::PostgreSQL(p) => {
let rows = sqlx::query(
"SELECT txid, tx, locktime, network, status FROM tbl_tx WHERE network = $1 AND status = 0 AND (locktime < $2 OR (locktime > $3 AND locktime < $4))",
)
.bind(network)
.bind(bestblock_height as i32)
.bind(locktime_threshold as i32)
.bind(bestblock_time as i32)
.fetch_all(p)
.await?;
for row in rows {
results.push(TxRow {
txid: row.try_get("txid")?,
tx: row.try_get("tx")?,
locktime: row.try_get::<i32, _>("locktime").unwrap_or(0) as i64,
network: row.try_get("network")?,
status: row.try_get::<i32, _>("status").unwrap_or(0) as i64,
});
}
}
}
Ok(results)
}
#[derive(Debug, Clone)]
pub struct StatsRow {
pub report_date: String,
pub chain: String,
pub totals: i64,
pub waiting: i64,
pub sent: i64,
pub failed: i64,
pub waiting_profit: i64,
pub sent_profit: i64,
pub missed_profit: i64,
pub unique_inputs: i64,
}
pub async fn get_stats(pool: &DatabasePool, chain: &str) -> Result<Vec<StatsRow>, sqlx::Error> {
let mut results = Vec::new();
match pool {
DatabasePool::SQLite(p) => {
let rows = sqlx::query(
"SELECT report_date, chain, totals, waiting, sent, failed, waiting_profit, sent_profit, missed_profit, unique_inputs FROM tbl_stats WHERE chain = ?",
)
.bind(chain)
.fetch_all(p)
.await?;
for row in rows {
results.push(StatsRow {
report_date: row.try_get("report_date").unwrap_or_default(),
chain: row.try_get("chain").unwrap_or_default(),
totals: row.try_get("totals").unwrap_or(0),
waiting: row.try_get("waiting").unwrap_or(0),
sent: row.try_get("sent").unwrap_or(0),
failed: row.try_get("failed").unwrap_or(0),
waiting_profit: row.try_get("waiting_profit").unwrap_or(0),
sent_profit: row.try_get("sent_profit").unwrap_or(0),
missed_profit: row.try_get("missed_profit").unwrap_or(0),
unique_inputs: row.try_get("unique_inputs").unwrap_or(0),
});
}
}
DatabasePool::PostgreSQL(p) => {
let rows = sqlx::query(
"SELECT report_date, chain, totals, waiting, sent, failed, waiting_profit, sent_profit, missed_profit, unique_inputs FROM tbl_stats WHERE chain = $1",
)
.bind(chain)
.fetch_all(p)
.await?;
for row in rows {
results.push(StatsRow {
report_date: row.try_get("report_date").unwrap_or_default(),
chain: row.try_get("chain").unwrap_or_default(),
totals: row.try_get::<i32, _>("totals").unwrap_or(0) as i64,
waiting: row.try_get::<i32, _>("waiting").unwrap_or(0) as i64,
sent: row.try_get::<i32, _>("sent").unwrap_or(0) as i64,
failed: row.try_get::<i32, _>("failed").unwrap_or(0) as i64,
waiting_profit: row.try_get::<i32, _>("waiting_profit").unwrap_or(0) as i64,
sent_profit: row.try_get::<i32, _>("sent_profit").unwrap_or(0) as i64,
missed_profit: row.try_get::<i32, _>("missed_profit").unwrap_or(0) as i64,
unique_inputs: row.try_get::<i32, _>("unique_inputs").unwrap_or(0) as i64,
});
}
}
}
Ok(results)
}
#[derive(Debug, Clone)]
pub struct SearchTxRow {
pub status: String,
pub tx: String,
pub our_address: String,
pub our_fees: String,
pub reqid: String,
}
pub async fn search_tx(
pool: &DatabasePool,
txid: &str,
) -> Result<Option<SearchTxRow>, sqlx::Error> {
match pool {
DatabasePool::SQLite(p) => {
let result = sqlx::query("SELECT * FROM tbl_tx WHERE txid = ? LIMIT 1")
.bind(txid)
.fetch_optional(p)
.await?;
Ok(result.map(|row| SearchTxRow {
status: row.try_get::<i64, _>("status").unwrap_or(0).to_string(),
tx: row.try_get::<String, _>("tx").unwrap_or_default(),
our_address: row.try_get::<String, _>("our_address").unwrap_or_default(),
our_fees: row.try_get::<String, _>("our_fees").unwrap_or_default(),
reqid: row.try_get::<String, _>("reqid").unwrap_or_default(),
}))
}
DatabasePool::PostgreSQL(p) => {
let result = sqlx::query("SELECT * FROM tbl_tx WHERE txid = $1 LIMIT 1")
.bind(txid)
.fetch_optional(p)
.await?;
Ok(result.map(|row| SearchTxRow {
status: row.try_get::<i32, _>("status").unwrap_or(0).to_string(),
tx: row.try_get::<String, _>("tx").unwrap_or_default(),
our_address: row.try_get::<String, _>("our_address").unwrap_or_default(),
our_fees: row.try_get::<String, _>("our_fees").unwrap_or_default(),
reqid: row.try_get::<String, _>("reqid").unwrap_or_default(),
}))
}
}
}
pub async fn calculate_and_upsert_stats(
pool: &DatabasePool,
chain: &str,
) -> Result<(), sqlx::Error> {
if !chain
.chars()
.all(|c| c.is_alphanumeric() || c == '-' || c == '_')
|| chain.is_empty()
{
error!("Invalid chain name: {chain}");
return Ok(());
}
match pool {
DatabasePool::SQLite(p) => {
sqlx::query("DELETE FROM tbl_stats WHERE chain = ?")
.bind(chain)
.execute(p)
.await?;
sqlx::query(
"INSERT INTO tbl_stats (
report_date, chain, totals, waiting, sent, failed,
waiting_profit, sent_profit, missed_profit, unique_inputs
)
SELECT
CURRENT_TIMESTAMP,
?,
(SELECT COUNT(*) FROM tbl_tx WHERE network = ?),
(SELECT COUNT(*) FROM tbl_tx WHERE status = 0 AND network = ?),
(SELECT COUNT(*) FROM tbl_tx WHERE status = 1 AND network = ?),
(SELECT COUNT(*) FROM tbl_tx WHERE status = 2 AND network = ?),
(SELECT IFNULL(SUM(our_fees),0) FROM tbl_tx WHERE status = 0 AND network = ?),
(SELECT IFNULL(SUM(our_fees),0) FROM tbl_tx WHERE status = 1 AND network = ?),
(SELECT IFNULL(SUM(our_fees),0) FROM tbl_tx WHERE status = 2 AND network = ?),
(SELECT COUNT(DISTINCT tbl_inp.in_txid)
FROM tbl_inp
JOIN tbl_tx ON tbl_inp.txid = tbl_tx.txid
WHERE tbl_tx.status = 0 AND tbl_tx.network = ?)
ON CONFLICT(chain) DO UPDATE SET
report_date = excluded.report_date,
totals = excluded.totals,
waiting = excluded.waiting,
sent = excluded.sent,
failed = excluded.failed,
waiting_profit = excluded.waiting_profit,
sent_profit = excluded.sent_profit,
missed_profit = excluded.missed_profit,
unique_inputs = excluded.unique_inputs",
)
.bind(chain)
.bind(chain)
.bind(chain)
.bind(chain)
.bind(chain)
.bind(chain)
.bind(chain)
.bind(chain)
.bind(chain)
.execute(p)
.await?;
}
DatabasePool::PostgreSQL(p) => {
sqlx::query("DELETE FROM tbl_stats WHERE chain = $1")
.bind(chain)
.execute(p)
.await?;
sqlx::query(
"INSERT INTO tbl_stats (
report_date, chain, totals, waiting, sent, failed,
waiting_profit, sent_profit, missed_profit, unique_inputs
)
SELECT
CURRENT_TIMESTAMP,
$1,
(SELECT COUNT(*) FROM tbl_tx WHERE network = $1),
(SELECT COUNT(*) FROM tbl_tx WHERE status = 0 AND network = $1),
(SELECT COUNT(*) FROM tbl_tx WHERE status = 1 AND network = $1),
(SELECT COUNT(*) FROM tbl_tx WHERE status = 2 AND network = $1),
(SELECT COALESCE(SUM(our_fees),0) FROM tbl_tx WHERE status = 0 AND network = $1),
(SELECT COALESCE(SUM(our_fees),0) FROM tbl_tx WHERE status = 1 AND network = $1),
(SELECT COALESCE(SUM(our_fees),0) FROM tbl_tx WHERE status = 2 AND network = $1),
(SELECT COUNT(DISTINCT tbl_inp.in_txid)
FROM tbl_inp
JOIN tbl_tx ON tbl_inp.txid = tbl_tx.txid
WHERE tbl_tx.status = 0 AND tbl_tx.network = $1)
ON CONFLICT(chain) DO UPDATE SET
report_date = excluded.report_date,
totals = excluded.totals,
waiting = excluded.waiting,
sent = excluded.sent,
failed = excluded.failed,
waiting_profit = excluded.waiting_profit,
sent_profit = excluded.sent_profit,
missed_profit = excluded.missed_profit,
unique_inputs = excluded.unique_inputs",
)
.bind(chain)
.execute(p)
.await?;
}
}
info!("tbl_stats creation success");
Ok(())
}

233
src/db/schema.rs Normal file
View File

@@ -0,0 +1,233 @@
use sqlx::{PgPool, SqlitePool};
pub async fn create_sqlite_schema(pool: &SqlitePool) -> Result<(), sqlx::Error> {
sqlx::query(
"CREATE TABLE IF NOT EXISTS tbl_tx (
txid TEXT PRIMARY KEY,
date_creation TIMESTAMP DEFAULT CURRENT_TIMESTAMP,
date_update TIMESTAMP DEFAULT CURRENT_TIMESTAMP,
wtxid TEXT,
ntxid TEXT,
tx TEXT,
locktime INTEGER,
network TEXT,
network_fees TEXT,
reqid TEXT,
our_fees TEXT,
our_address TEXT,
status INTEGER DEFAULT 0
)",
)
.execute(pool)
.await?;
let _ = sqlx::query("ALTER TABLE tbl_tx ADD COLUMN push_err TEXT")
.execute(pool)
.await;
sqlx::query(
"CREATE TABLE IF NOT EXISTS tbl_inp (
id INTEGER,
txid TEXT,
in_txid TEXT,
in_vout INTEGER
)",
)
.execute(pool)
.await?;
sqlx::query(
"CREATE UNIQUE INDEX IF NOT EXISTS idx_inp_unique ON tbl_inp(txid, in_txid, in_vout)",
)
.execute(pool)
.await?;
sqlx::query(
"CREATE TABLE IF NOT EXISTS tbl_out (
id INTEGER,
txid TEXT,
script_pubkey TEXT,
amount TEXT,
vout INTEGER
)",
)
.execute(pool)
.await?;
sqlx::query(
"CREATE UNIQUE INDEX IF NOT EXISTS idx_out_unique ON tbl_out(txid, script_pubkey, amount, vout)",
)
.execute(pool)
.await?;
sqlx::query(
"CREATE TABLE IF NOT EXISTS tbl_xpub (
id INTEGER PRIMARY KEY,
network TEXT,
xpub TEXT,
date_create TIMESTAMP DEFAULT CURRENT_TIMESTAMP,
path_idx INTEGER DEFAULT -1
)",
)
.execute(pool)
.await?;
sqlx::query("CREATE UNIQUE INDEX IF NOT EXISTS idx_xpub ON tbl_xpub(network, xpub)")
.execute(pool)
.await?;
sqlx::query(
"CREATE TABLE IF NOT EXISTS tbl_address (
address TEXT PRIMARY KEY,
path TEXT NOT NULL,
date_create TIMESTAMP DEFAULT CURRENT_TIMESTAMP,
xpub INTEGER,
remote_address TEXT
)",
)
.execute(pool)
.await?;
sqlx::query(
"CREATE TABLE IF NOT EXISTS tbl_stats (
report_date TEXT,
chain TEXT,
totals INTEGER,
waiting INTEGER,
sent INTEGER,
failed INTEGER,
waiting_profit INTEGER,
sent_profit INTEGER,
missed_profit INTEGER,
unique_inputs INTEGER
)",
)
.execute(pool)
.await?;
sqlx::query("DROP INDEX IF EXISTS idx_stats_chain")
.execute(pool)
.await?;
sqlx::query("CREATE UNIQUE INDEX IF NOT EXISTS idx_stats_chain ON tbl_stats(chain)")
.execute(pool)
.await?;
sqlx::query("UPDATE tbl_tx SET network='bitcoin' WHERE network='mainnet'")
.execute(pool)
.await?;
Ok(())
}
pub async fn create_pg_schema(pool: &PgPool) -> Result<(), sqlx::Error> {
sqlx::query(
"CREATE TABLE IF NOT EXISTS tbl_tx (
txid TEXT PRIMARY KEY,
date_creation TIMESTAMP DEFAULT CURRENT_TIMESTAMP,
date_update TIMESTAMP DEFAULT CURRENT_TIMESTAMP,
wtxid TEXT,
ntxid TEXT,
tx TEXT,
locktime INTEGER,
network TEXT,
network_fees TEXT,
reqid TEXT,
our_fees TEXT,
our_address TEXT,
status INTEGER DEFAULT 0,
push_err TEXT
)",
)
.execute(pool)
.await?;
sqlx::query(
"CREATE TABLE IF NOT EXISTS tbl_inp (
id SERIAL PRIMARY KEY,
txid TEXT,
in_txid TEXT,
in_vout INTEGER
)",
)
.execute(pool)
.await?;
sqlx::query(
"CREATE UNIQUE INDEX IF NOT EXISTS idx_inp_unique ON tbl_inp(txid, in_txid, in_vout)",
)
.execute(pool)
.await?;
sqlx::query(
"CREATE TABLE IF NOT EXISTS tbl_out (
id SERIAL PRIMARY KEY,
txid TEXT,
script_pubkey TEXT,
amount TEXT,
vout INTEGER
)",
)
.execute(pool)
.await?;
sqlx::query(
"CREATE UNIQUE INDEX IF NOT EXISTS idx_out_unique ON tbl_out(txid, script_pubkey, amount, vout)",
)
.execute(pool)
.await?;
sqlx::query(
"CREATE TABLE IF NOT EXISTS tbl_xpub (
id SERIAL PRIMARY KEY,
network TEXT,
xpub TEXT,
date_create TIMESTAMP DEFAULT CURRENT_TIMESTAMP,
path_idx INTEGER DEFAULT -1
)",
)
.execute(pool)
.await?;
sqlx::query("CREATE UNIQUE INDEX IF NOT EXISTS idx_xpub ON tbl_xpub(network, xpub)")
.execute(pool)
.await?;
sqlx::query(
"CREATE TABLE IF NOT EXISTS tbl_address (
address TEXT PRIMARY KEY,
path TEXT NOT NULL,
date_create TIMESTAMP DEFAULT CURRENT_TIMESTAMP,
xpub INTEGER,
remote_address TEXT
)",
)
.execute(pool)
.await?;
sqlx::query(
"CREATE TABLE IF NOT EXISTS tbl_stats (
report_date TEXT,
chain TEXT,
totals INTEGER,
waiting INTEGER,
sent INTEGER,
failed INTEGER,
waiting_profit INTEGER,
sent_profit INTEGER,
missed_profit INTEGER,
unique_inputs INTEGER
)",
)
.execute(pool)
.await?;
sqlx::query("CREATE UNIQUE INDEX IF NOT EXISTS idx_stats_chain ON tbl_stats(chain)")
.execute(pool)
.await?;
sqlx::query("UPDATE tbl_tx SET network='bitcoin' WHERE network='mainnet'")
.execute(pool)
.await?;
Ok(())
}

242
src/xpub.rs2 Normal file
View File

@@ -0,0 +1,242 @@
//use bs58;
use bitcoin::Address;
use bitcoin::Network;
use bitcoin::ScriptBuf;
use bitcoin::WPubkeyHash;
use bitcoin::bip32::DerivationPath;
use bitcoin::bip32::Xpub;
use bitcoin::hashes::Hash;
use bitcoin::key::Secp256k1;
use sha2::{Digest, Sha256};
use std::str::FromStr;
// Mainnet (BIP44/BIP49/BIP84)
enum BS58Prefix {
Xpub,
Ypub,
Zpub,
Tpub,
Vpub,
Upub,
}
const XPUB_PREFIX: [u8; 4] = [0x04, 0x88, 0xB2, 0x1E]; // xpub (Legacy P2PKH)
const YPUB_PREFIX: [u8; 4] = [0x04, 0x9D, 0x7C, 0xB2]; // ypub (Nested SegWit P2SH-P2WPKH)
const ZPUB_PREFIX: [u8; 4] = [0x04, 0xB2, 0x47, 0x46]; // zpub (Native SegWit P2WPKH)
const TPUB_PREFIX: [u8; 4] = [0x04, 0x35, 0x87, 0xCF]; // tpub (Testnet Legacy P2PKH)
const VPUB_PREFIX: [u8; 4] = [0x04, 0x5F, 0x1C, 0xF6]; // vpub (Testnet Nested SegWit)
const UPUB_PREFIX: [u8; 4] = [0x04, 0x4A, 0x52, 0x62]; // upub (RegTest Nested SegWit)
// Constants from Bitcoin Core's checksum algorithm
const INPUT_CHARSET: &[u8] = b"0123456789()[],'/*abcdefgh@:$%{}IJKLMNOPQRSTUVWXYZ&+-.;<=>?!^_|~ijklmnopqrstuvwxyzABCDEFGH`#\"\\ ";
const CHECKSUM_CHARSET: &[u8] = b"qpzry9x8gf2tvdw0s3jn54khce6mua7l";
// Polynomial modulo function used in checksum calculation (same as in Bitcoin Core)
fn poly_mod(mut c: u64, val: u64) -> u64 {
let c0 = c >> 35;
c = ((c & 0x7ffffffff) << 5) ^ val;
if c0 & 1 > 0 {
c ^= 0xf5dee51989
};
if c0 & 2 > 0 {
c ^= 0xa9fdca3312
};
if c0 & 4 > 0 {
c ^= 0x1bab10e32d
};
if c0 & 8 > 0 {
c ^= 0x3706b1677a
};
if c0 & 16 > 0 {
c ^= 0x644d626ffd
};
c
}
// Calculate checksum for a descriptor string
fn calc_checksum(desc: &str) -> Result<String, String> {
// Separate descriptor from any existing checksum
let desc = match desc.split_once('#') {
Some((d, _)) => d,
None => desc,
};
let mut c: u64 = 1;
let mut cls: u64 = 0;
let mut clscount: u64 = 0;
// Process each character in the descriptor
for ch in desc.as_bytes() {
let pos = match INPUT_CHARSET.iter().position(|b| b == ch) {
Some(p) => p as u64,
None => return Err(format!("Invalid character in descriptor: {}", *ch as char)),
};
c = poly_mod(c, pos & 31);
cls = cls * 3 + (pos >> 5);
clscount += 1;
if clscount == 3 {
c = poly_mod(c, cls);
cls = 0;
clscount = 0;
}
}
if clscount > 0 {
c = poly_mod(c, cls);
}
// Final steps in checksum calculation
for _ in 0..8 {
c = poly_mod(c, 0);
}
c ^= 1;
// Convert checksum to characters
let mut checksum = String::with_capacity(8);
for j in 0..8 {
let idx = ((c >> (5 * (7 - j))) & 31) as usize;
checksum.push(CHECKSUM_CHARSET[idx] as char);
}
Ok(checksum)
}
pub fn get_bitcoincore_descriptor(xpub: &String) -> String {
let fingerprint = calculate_fingerprint(xpub);
let mut bip = 84;
let cpub = xpub.to_string();
match &xpub[0..4] {
"vpub" => {
bip = 84;
}
"zpub" => {
bip = 84;
}
&_ => {
bip = 84;
}
};
let descriptor = format!(
"wpkh([{}/84h/0h/0h]{}/0/*)",
fingerprint,
convert_xpub(xpub)
);
let 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()
}
};
descriptor
//format!("{}#{}",descriptor,checksum)
}
fn convert_xpub(xpub: &String) -> String {
if xpub[0..4] == *"xpub" || xpub[0..4] == *"ypub" || xpub[0..4] == *"zpub" {
return convert_to(xpub, BS58Prefix::Xpub).unwrap();
} else {
return convert_to(xpub, BS58Prefix::Tpub).unwrap();
}
}
pub fn calculate_fingerprint(tpub: &str) -> String {
let xpub = Xpub::from_str(&convert_to(tpub, BS58Prefix::Xpub).unwrap()).unwrap();
let fp = xpub.fingerprint();
let pp = xpub.parent_fingerprint;
format!("{}", fp)
}
fn base58check_decode(s: &str) -> Result<Vec<u8>, String> {
let data = bs58::decode(s).into_vec().map_err(|e| e.to_string())?;
if data.len() < 4 {
return Err("Data troppo corta".to_string());
}
let (payload, checksum) = data.split_at(data.len() - 4);
let hash = Sha256::digest(&Sha256::digest(payload));
if hash[0..4] != checksum[..] {
return Err("Checksum invalido".to_string());
}
Ok(payload.to_vec())
}
fn base58check_encode(data: &[u8]) -> String {
let checksum = &Sha256::digest(&Sha256::digest(data))[0..4];
let full = [data, checksum].concat();
bs58::encode(full).into_string()
}
fn convert_to(zpub: &str, prefix: BS58Prefix) -> Result<String, String> {
let mut data = base58check_decode(zpub)?;
if data.len() < 4 {
return Err("Non è una zpub valida.".to_string());
}
data.splice(
0..4,
match prefix {
BS58Prefix::Xpub => XPUB_PREFIX,
BS58Prefix::Ypub => YPUB_PREFIX,
BS58Prefix::Zpub => ZPUB_PREFIX,
BS58Prefix::Vpub => VPUB_PREFIX,
BS58Prefix::Tpub => TPUB_PREFIX,
BS58Prefix::Upub => UPUB_PREFIX,
},
);
Ok(base58check_encode(&data))
}
pub fn new_address_from_xpub(
zpub: &str,
index: i64,
network: Network,
) -> Result<(String, String), Box<dyn std::error::Error>> {
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 secp = Secp256k1::new();
let derived_xpub = xpub.derive_pub(&secp, &derivation_path)?;
let public_key = derived_xpub.public_key;
let pubkey_bytes = public_key.serialize();
let witness_program = WPubkeyHash::hash(&pubkey_bytes);
let redeem_script = ScriptBuf::new_p2wpkh(&witness_program);
//let script_pubkey = ScriptBuf::new_p2sh(&redeem_script.script_hash());
let address = Address::from_script(&redeem_script, network)?;
//let address = Address::from_script(&script_pubkey, network)?;
Ok((address.to_string(), path.to_string()))
}
/*
fn main() -> Result<(), Box<dyn std::error::Error>>{
//let zpub = "xpub6C29v8gxCXREHUzoGNfqqFqZWxTVEmYtmZshuzfSwBKNmfYQxoizRziCkkUUA4WwJZkJs2i7nttRiC6MQG7mxZpouXeYkTZe3U52RyPAeo2";
//let zpub = "vpub5Ut36m34VebUUjdhYaxJCjSPqk3ZR8bA2MXLmbHRQCycAxy5Q1GFPJspLkJywJjBgQnvU3rmwPKTPp1ELLWeXrve3zBufpZR4MRCCTNHzsn";
let zpub = "zpub6qdfveGrxBQN3z8paZ88EHpCn5MGXpUoHwQmHhPbj4rPQtUjbWyCHrJFYZGVY7MsmVbDaeu4JYqRqcdLzMx78wZFEWbLrF9FG3gr2MPQC5H";
match convert_to(zpub,BS58Prefix::Tpub) {
Ok(tpub) => println!("XPUB: {}", tpub),
Err(e) => eprintln!("Errore: {}", e),
}
let fingerprint = base58check_encode(&calculate_fingerprint(zpub));
println!("ZPUB: {}, FINGERPRINT: {}",zpub,fingerprint);
let xpub = Xpub::from_str(&convert_to(zpub,BS58Prefix::Xpub)?)?;
let tpub = convert_to(zpub,BS58Prefix::Tpub)?;
let fingerprint = base58check_encode(&calculate_fingerprint(&tpub));
println!("TPUB: {}, FINGERPRINT: {}",tpub,fingerprint);
let derivation_path = DerivationPath::from_str("m/0/0")?;
let secp = Secp256k1::new();
let derived_xpub = xpub.derive_pub(&secp, &derivation_path)?;
let public_key = derived_xpub.public_key;
let pubkey_bytes = public_key.serialize();
let witness_program = WPubkeyHash::hash(&pubkey_bytes);
let redeem_script = ScriptBuf::new_p2wpkh(&witness_program);
let script_pubkey = ScriptBuf::new_p2sh(&redeem_script.script_hash());
// Generate the Bitcoin SegWit (BIP49) address
let network = Network::Bitcoin;
let address = Address::from_script(&redeem_script, network)?;
let address = Address::from_script(&script_pubkey, network)?;
Ok(())
}*/