formatting code

This commit is contained in:
2026-02-06 12:03:10 -04:00
parent efa8c470cc
commit dc02d9248c
3 changed files with 217 additions and 69 deletions

View File

@@ -11,7 +11,7 @@ use log::{debug, error, info, trace, warn};
use serde::Deserialize;
use serde::Serialize;
use serde_json::json;
use sqlite::{Value, Connection};
use sqlite::{Connection, Value};
use std::collections::HashMap;
use std::env;
use std::error::Error as StdError;
@@ -48,14 +48,18 @@ struct MyConfig {
impl Default for MyConfig {
fn default() -> Self {
MyConfig {
zmq_listener: env::var("BAL_PUSHER_ZMQ_LISTENER").unwrap_or("tcp://127.0.0.1:28332".to_string()),
zmq_listener: env::var("BAL_PUSHER_ZMQ_LISTENER")
.unwrap_or("tcp://127.0.0.1:28332".to_string()),
db_file: env::var("BAL_PUSHER_DB_FILE").unwrap_or("bal.db".to_string()),
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),
signet: get_network_params_default(Network::Signet),
mainnet: get_network_params_default(Network::Bitcoin),
send_stats: env::var("BAL_PUSHER_SEND_STATS").unwrap_or("false".to_string()).parse::<bool>().unwrap(),
send_stats: env::var("BAL_PUSHER_SEND_STATS")
.unwrap_or("false".to_string())
.parse::<bool>()
.unwrap(),
url: env::var("BAL_SERVER_URL").unwrap_or("http://localhost/".to_string()),
ssl_key_path: env::var("SSL_KEY_PATH").unwrap_or("privkey.pem".to_string()),
}
@@ -216,7 +220,7 @@ async fn main_result(cfg: &MyConfig, network_params: &NetworkParams) -> Result<(
let average_time = bcinfo.median_time;
let db = sqlite::open(&cfg.db_file).unwrap();
info!("db open {}",&cfg.db_file);
info!("db open {}", &cfg.db_file);
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 = db.prepare(sqlquery).unwrap().into_iter();
@@ -294,7 +298,7 @@ async fn main_result(cfg: &MyConfig, network_params: &NetworkParams) -> Result<(
}
}
let _ = send_stats_report(cfg, bcinfo).await;
let _ = calculate_stats(&db,network_params.db_field.clone()).await;
let _ = calculate_stats(&db, network_params.db_field.clone()).await;
}
Err(erx) => {
panic!("impossible to get client {}", erx)
@@ -302,13 +306,14 @@ async fn main_result(cfg: &MyConfig, network_params: &NetworkParams) -> Result<(
}
Ok(())
}
async fn calculate_stats(db: &Connection,chain: String) -> Result<(), reqwest::Error> {
async fn calculate_stats(db: &Connection, chain: String) -> Result<(), reqwest::Error> {
//let sql = "drop table if exists tbl_stats;";
let sql = "DELETE FROM tbl_stats WHERE chain = '{chain}';";
if let Err(err) = db.execute(&sql){
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 (
let sql = format!(
"INSERT INTO tbl_stats (
report_date, chain, totals, waiting, sent, failed,
waiting_profit, sent_profit, missed_profit, unique_inputs
)
@@ -322,7 +327,7 @@ VALUES (
(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) -- or appropriate input identifier
(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}')
@@ -337,10 +342,11 @@ ON CONFLICT(chain) DO UPDATE SET
sent_profit = excluded.sent_profit,
missed_profit = excluded.missed_profit,
unique_inputs = excluded.unique_inputs;
");
"
);
/*
let sql = format!("CREATE TABLE tbl_stats AS
let sql = format!("CREATE TABLE tbl_stats AS
SELECT
CURRENT_TIMESTAMP AS report_date,
'{chain}' as chain,
@@ -364,10 +370,9 @@ ON CONFLICT(chain) DO UPDATE SET
unique_inputs = (SELECT COUNT(*) FROM tbl_inp JOIN tbl_tx ON(tbl_inp.txid = tbl_tx.txid) WHERE tbl_tx.status=0 AND tbl_tx.network ='{chain}')
WHERE chain = '{chain}'
*/
if let Err(err) = db.execute(&sql){
if let Err(err) = db.execute(&sql) {
error!("error inserting creating stats table {err}");
}
else{
} else {
info!("tbl_stats creation success");
}
Ok(())
@@ -495,7 +500,6 @@ fn parse_env_netconfig(cfg_lock: &mut MyConfig, chain: &str) -> NetworkParams {
cfg.clone()
}
fn check_zmq_connection(endpoint: &str) -> bool {
trace!("check zmq connection");
let context = Context::new();

View File

@@ -71,7 +71,7 @@ struct MyConfig {
}
#[derive(Debug, Serialize, Deserialize)]
pub struct InfoResponse{
pub struct InfoResponse {
pub address: String,
pub base_fee: u64,
pub chain: String,
@@ -82,14 +82,14 @@ pub struct InfoResponse{
pub struct StatsResponse {
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 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,
}
impl Default for MyConfig {
@@ -104,7 +104,10 @@ impl Default for MyConfig {
db_file: "bal.db".to_string(),
info: "Will Executor Server".to_string(),
pub_key_path: "public_key.pem".to_string(),
expose_stats:env::var("BAL_SERVER_EXPOSE_STATS").unwrap_or("false".to_string()).parse::<bool>().unwrap(),
expose_stats: env::var("BAL_SERVER_EXPOSE_STATS")
.unwrap_or("false".to_string())
.parse::<bool>()
.unwrap(),
}
}
}
@@ -136,15 +139,15 @@ async fn echo_pub_key(
async fn echo_stats(
param: &str,
cfg: &MyConfig,
remote_addr: &String,
) -> Result<Response<BoxBody<Bytes, hyper::Error>>, hyper::Error> {
info!("echo stats!!! {} - {}", param,cfg.expose_stats);
info!("echo stats!!! {} - {}", param, cfg.expose_stats);
let netconfig = MyConfig::get_net_config(cfg, param);
if !netconfig.enabled {
debug!("network disabled {}", param);
return Ok(Response::new(full("network disabled")));
}
let sql = format!("SELECT
let sql = format!(
"SELECT
report_date,
chain,
totals,
@@ -155,26 +158,43 @@ async fn echo_stats(
sent_profit,
missed_profit,
unique_inputs FROM tbl_stats where chain = '{}'
", netconfig.name);
let mut stats:Vec<StatsResponse>=vec![];
",
netconfig.name
);
let mut stats: Vec<StatsResponse> = vec![];
let db = sqlite::open(&cfg.db_file).unwrap();
db.iterate(&sql,|pairs|{
let row: HashMap<_, _> = pairs.into_iter().map(|(k,v)| (k.to_string(), v.map(|s| s))).collect();
let _ = db.iterate(&sql, |pairs| {
let row: HashMap<_, _> = pairs
.into_iter()
.map(|(k, v)| (k.to_string(), v.map(|s| s)))
.collect();
//let row:HashMap<_,_>= pairs.into_iter().collect();
println!("row report date {}",row["report_date"].clone().unwrap());
println!("row report date {}", row["report_date"].clone().unwrap());
dbg!(&row);
stats.push(StatsResponse{
stats.push(StatsResponse {
report_date: row["report_date"].clone().unwrap().to_string(),
chain: row["chain"].clone().unwrap().to_string(),
totals: row["totals"].clone().unwrap().parse::<i64>().unwrap(),
waiting: row["waiting"].clone().unwrap().parse::<i64>().unwrap(),
sent: row["sent"].clone().unwrap().parse::<i64>().unwrap(),
failed: row["failed"].clone().unwrap().parse::<i64>().unwrap(),
waiting_profit: row["waiting_profit"].clone().unwrap().parse::<i64>().unwrap(),
waiting_profit: row["waiting_profit"]
.clone()
.unwrap()
.parse::<i64>()
.unwrap(),
sent_profit: row["sent_profit"].clone().unwrap().parse::<i64>().unwrap(),
missed_profit:row["missed_profit"].clone().unwrap().parse::<i64>().unwrap(),
unique_inputs: row["unique_inputs"].clone().unwrap().parse::<i64>().unwrap(),
missed_profit: row["missed_profit"]
.clone()
.unwrap()
.parse::<i64>()
.unwrap(),
unique_inputs: row["unique_inputs"]
.clone()
.unwrap()
.parse::<i64>()
.unwrap(),
});
true
});
@@ -185,8 +205,7 @@ async fn echo_stats(
}
Err(err) => Ok(Response::new(full(format!("error:{}", err)))),
}
}
}
async fn echo_info(
param: &str,
@@ -332,7 +351,7 @@ async fn echo_push(
*response_not_enable.status_mut() = StatusCode::BAD_REQUEST;
let netconfig = MyConfig::get_net_config(cfg, param);
if !netconfig.enabled {
trace!("network not enabled {}",&netconfig.name);
trace!("network not enabled {}", &netconfig.name);
return Ok(response_not_enable);
}
let req_time = Utc::now().timestamp_nanos_opt().unwrap(); // Returns i64
@@ -438,7 +457,7 @@ async fn echo_push(
}
if address == our_address && amount.to_sat() >= netconfig.fixed_fee {
our_fees = amount.to_sat();
our_address = netconfig.address.to_string();
//our_address = netconfig.address.to_string();
found = true;
trace!("address and fees are correct {}: {}", our_address, our_fees);
}
@@ -541,7 +560,7 @@ async fn echo(
}
Method::GET => {
if let Some(param) = match_uri(r"^?/?(?P<param>[^/]?+)?/stats$", uri.as_str()) {
ret = echo_stats(param, cfg, &remote_addr).await;
ret = echo_stats(param, cfg).await;
}
if let Some(param) = match_uri(r"^?/?(?P<param>[^/]?+)?/info$", uri.as_str()) {
ret = echo_info(param, cfg, &remote_addr).await;