Compare commits
5 Commits
c371a4f478
...
v0.3.2
| Author | SHA1 | Date | |
|---|---|---|---|
|
7999902cc0
|
|||
|
d6b888e403
|
|||
|
36219c49a0
|
|||
|
ca530bf987
|
|||
|
1c76755ea6
|
1
.gitignore
vendored
1
.gitignore
vendored
@@ -36,3 +36,4 @@ Cargo.lock
|
||||
!lib/
|
||||
!contrib/
|
||||
!src/
|
||||
make_release.sh
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
[package]
|
||||
name = "bal_server"
|
||||
version = "0.3.1"
|
||||
version = "0.3.2"
|
||||
edition = "2024"
|
||||
|
||||
# See more keys and their definitions at https://doc.rust-lang.org/cargo/reference/manifest.html
|
||||
|
||||
@@ -12,4 +12,4 @@ export WELIST_SERVER_URL=http://localhost:8086
|
||||
export WELIST_SKIP_URL_VALIDATION=true
|
||||
export BAL_SERVER_URL="http://127.0.0.1:9133"
|
||||
export SSL_KEY_PATH="$(pwd)/private_key.pem"
|
||||
cargo run --bin=bal-pusher regtest
|
||||
cargo run --bin=bal-pusher regtest --features=pusher
|
||||
|
||||
203
make_release.sh
203
make_release.sh
@@ -1,203 +0,0 @@
|
||||
#!/bin/bash
|
||||
#author: <your-name>
|
||||
|
||||
source lib.sh
|
||||
usage() {
|
||||
echo_w "./make_release <version> <message>"
|
||||
}
|
||||
|
||||
if [ -n "$1" ]; then release=$1; else usage; exit; fi
|
||||
|
||||
if [ -n "$2" ]; then message=$2; else
|
||||
# Create temporary file using mktemp
|
||||
TEMPFILE=$(mktemp)
|
||||
vi $TEMPFILE
|
||||
message=$(cat $TEMPFILE)
|
||||
rm $TEMPFILE
|
||||
fi
|
||||
|
||||
echo_i $message
|
||||
|
||||
# Load secrets from .env file (not committed to git)
|
||||
if [ -f .env ]; then
|
||||
export $(grep -v '^#' .env | xargs)
|
||||
fi
|
||||
|
||||
TOKEN="${GITEA_API_TOKEN:-${TOKEN_GITEA}}"
|
||||
if [ -z "$TOKEN" ]; then
|
||||
echo_e "Error: GITEA_API_TOKEN (or TOKEN_GITEA) is not set in .env file."
|
||||
echo_e "Please create a .env file with: GITEA_API_TOKEN=your_token_here"
|
||||
exit 1
|
||||
fi
|
||||
|
||||
for cmd in gpg sha256sum jq; do
|
||||
if ! command -v "$cmd" >/dev/null 2>&1; then
|
||||
echo_e "Error: '$cmd' is required but not installed."
|
||||
exit 1
|
||||
fi
|
||||
done
|
||||
|
||||
OWNER="bitcoinafterlife"
|
||||
basename=$(basename $(pwd))
|
||||
REPO=$basename
|
||||
TAG="v$release"
|
||||
binpath="target/release/$basename"
|
||||
release_name="$basename-$release"
|
||||
dest="releases/$release"
|
||||
arch=$(uname -m)
|
||||
platform="linux-gnu"
|
||||
destbin="$dest/$arch"
|
||||
destsrc="$dest/src"
|
||||
assetname="$release_name""_$arch""_$platform"
|
||||
|
||||
asset_tar_gz="$assetname.tar.gz"
|
||||
ASSET_PATH="$destbin/$assetname.tar.gz"
|
||||
|
||||
SIGNER_KEY="svatantrya@bitcoin-after.life"
|
||||
ASSET_SHA256="$ASSET_PATH.sha256"
|
||||
ASSET_SIG="$ASSET_PATH.sig"
|
||||
ASSET_ASC="$ASSET_PATH.asc"
|
||||
|
||||
giteahost="https://bitcoin-after.life/gitea"
|
||||
url_releases="$giteahost/api/v1/repos/$OWNER/$REPO/releases"
|
||||
|
||||
echo_i() {
|
||||
echo -e "\033[1m==> $1\033[0m"
|
||||
}
|
||||
|
||||
echo_e() {
|
||||
echo -e "\033[31;1m$1\033[0m"
|
||||
}
|
||||
|
||||
echo_s() {
|
||||
echo -e "\033[32;1m$1\033[0m"
|
||||
}
|
||||
echo_w() {
|
||||
echo -e "\033[33;1m$1\033[0m"
|
||||
}
|
||||
|
||||
|
||||
prepare_release(){
|
||||
mkdir -p "$destbin/$assetname"
|
||||
if ! cargo build --release --no-default-features --features server --bin bal-server; then
|
||||
echo_w "error building bal-server"
|
||||
exit 1
|
||||
fi
|
||||
if ! cargo build --release --no-default-features --features pusher --bin bal-pusher; then
|
||||
echo_w "error building bal-pusher"
|
||||
exit 1
|
||||
fi
|
||||
ls -l target/release/bal-server target/release/bal-pusher
|
||||
cp target/release/bal-server \
|
||||
target/release/bal-pusher \
|
||||
README.md \
|
||||
"$destbin/$assetname"
|
||||
(
|
||||
cd "$destbin"
|
||||
echo_w $ASSET_PATH
|
||||
echo "ls $(pwd)"
|
||||
ls
|
||||
echo "ls $(pwd)/$assetname"
|
||||
ls "$(pwd)/$assetname"
|
||||
ls $assetname
|
||||
tar -czf "$asset_tar_gz" "$assetname"
|
||||
|
||||
sha256sum "$asset_tar_gz" > "$asset_tar_gz.sha256"
|
||||
|
||||
if ! gpg --batch --yes --detach-sign --local-user "$SIGNER_KEY" "$asset_tar_gz"; then
|
||||
echo_e "error signing release tarball (binary)"
|
||||
exit 1
|
||||
fi
|
||||
|
||||
if ! gpg --batch --yes --detach-sign --armor --output "$asset_tar_gz.asc" --local-user "$SIGNER_KEY" "$asset_tar_gz"; then
|
||||
echo_e "error signing release tarball (ascii-armored)"
|
||||
exit 1
|
||||
fi
|
||||
)
|
||||
}
|
||||
|
||||
push_tag() {
|
||||
git commit -am"release: $release_name"
|
||||
git push
|
||||
#git tag -a "$TAG" -m"release: $release_name"
|
||||
#git push origin --tags
|
||||
}
|
||||
|
||||
# Configurazioni
|
||||
post_release() {
|
||||
if [ -z "$1" ]; then
|
||||
echo_e "no data to release"
|
||||
exit 1
|
||||
else
|
||||
echo "data: $1"
|
||||
fi
|
||||
echo "token:$TOKEN"
|
||||
echo url_releases: $url_releases
|
||||
RELEASE="$(curl -s -X POST \
|
||||
-H "accept: application/json" \
|
||||
-H "Authorization: token $TOKEN" \
|
||||
-H "Content-Type: application/json" \
|
||||
-d "$1" \
|
||||
$url_releases
|
||||
)"
|
||||
echo $RELEASE
|
||||
}
|
||||
|
||||
add_asset_release() {
|
||||
if [ -z "$1" ]; then
|
||||
echo_e "error add_asset_release"
|
||||
exit 1
|
||||
fi
|
||||
|
||||
local assets=("$ASSET_PATH" "$ASSET_SHA256" "$ASSET_SIG" "$ASSET_ASC")
|
||||
for asset in "${assets[@]}"; do
|
||||
if [ ! -f "$asset" ]; then
|
||||
echo_e "missing asset: $asset"
|
||||
exit 1
|
||||
fi
|
||||
echo "Uploading: $asset"
|
||||
ls -l "$asset"
|
||||
curl -X POST \
|
||||
-H "accept: application/json" \
|
||||
-H "Authorization: token $TOKEN" \
|
||||
-H "Content-Type: multipart/form-data" \
|
||||
-F "attachment=@$asset" \
|
||||
"$url_releases/$1/assets"
|
||||
done
|
||||
}
|
||||
|
||||
release_body=$(printf '%s' "Release: $release_name enjoy
|
||||
|
||||
$message
|
||||
|
||||
---
|
||||
Verification Instructions:
|
||||
|
||||
SHA256 Checksum:
|
||||
sha256sum -c $assetname.tar.gz.sha256
|
||||
|
||||
GPG Signature (binary):
|
||||
gpg --verify $assetname.tar.gz.sig $assetname.tar.gz
|
||||
|
||||
GPG Signature (ASCII-armored):
|
||||
gpg --verify $assetname.tar.gz.asc $assetname.tar.gz
|
||||
|
||||
Signed by Svātantrya (svatantrya@bitcoin-after.life)")
|
||||
|
||||
release_data=$(jq -n -c \
|
||||
--arg tag_name "$TAG" \
|
||||
--arg name "$release_name" \
|
||||
--arg body "$release_body" \
|
||||
'{tag_name: $tag_name, name: $name, body: $body}')
|
||||
|
||||
prepare_release
|
||||
echo_s "prepare release done"
|
||||
push_tag
|
||||
echo_s "push tag done"
|
||||
echo "$release_data"
|
||||
post_release "$release_data"
|
||||
echo_s "prepare release done"
|
||||
echo $RELEASE
|
||||
id_release=
|
||||
add_asset_release $(echo $RELEASE | jq .id)
|
||||
echo_s "done"
|
||||
@@ -5,7 +5,6 @@ use bitcoin::Network;
|
||||
use bitcoincore_rpc::{Auth, Client, Error, RpcApi, bitcoin};
|
||||
use bitcoincore_rpc_json::GetBlockchainInfoResult;
|
||||
|
||||
use byteorder::{LittleEndian, ReadBytesExt};
|
||||
use ed25519_dalek::{Signer as _, SigningKey, pkcs8::DecodePrivateKey};
|
||||
use log::{debug, error, info, trace, warn};
|
||||
use serde::Deserialize;
|
||||
@@ -15,21 +14,19 @@ use sqlite::{Connection, Value};
|
||||
use std::collections::HashMap;
|
||||
use std::env;
|
||||
use std::error::Error as StdError;
|
||||
use std::io::Cursor;
|
||||
use std::str;
|
||||
use std::{thread, time::Duration};
|
||||
use zmq::{Context, DEALER, DONTWAIT, Socket};
|
||||
use zmq::{Context, Socket};
|
||||
|
||||
use bal_server::db::open_db;
|
||||
use bal_server::validation::is_valid_welist_url;
|
||||
use base64::{Engine as _, engine::general_purpose};
|
||||
use reqwest::Client as rClient;
|
||||
use std::net::SocketAddr;
|
||||
use std::time::Instant;
|
||||
use url::Url;
|
||||
|
||||
const LOCKTIME_THRESHOLD: i64 = 5000000;
|
||||
const VERSION: &str = "0.0.2";
|
||||
const VERSION: &str = env!("CARGO_PKG_VERSION");
|
||||
#[derive(Debug, Clone, Serialize, Deserialize)]
|
||||
struct MyConfig {
|
||||
db_file: String,
|
||||
@@ -87,7 +84,7 @@ fn get_network_params(cfg: &MyConfig, network: Network) -> &NetworkParams {
|
||||
fn get_network_params_default(network: Network) -> NetworkParams {
|
||||
match network {
|
||||
Network::Testnet => NetworkParams {
|
||||
host: "http://i27.0.0.1".to_string(),
|
||||
host: "http://127.0.0.1".to_string(),
|
||||
port: 18332,
|
||||
dir_path: "testnet3/".to_string(),
|
||||
db_field: "testnet".to_string(),
|
||||
@@ -97,7 +94,7 @@ fn get_network_params_default(network: Network) -> NetworkParams {
|
||||
zmq_listener: "tcp://127.0.0.1:23332".to_string(),
|
||||
},
|
||||
Network::Testnet4 => NetworkParams {
|
||||
host: "http://i27.0.0.1".to_string(),
|
||||
host: "http://127.0.0.1".to_string(),
|
||||
port: 48332,
|
||||
dir_path: "testnet4/".to_string(),
|
||||
db_field: "testnet4".to_string(),
|
||||
@@ -205,29 +202,9 @@ fn get_client(
|
||||
}
|
||||
}
|
||||
async fn main_result(cfg: &MyConfig, network_params: &NetworkParams) -> Result<(), Error> {
|
||||
/*let url = args.next().expect("Usage: <rpc_url> <username> <password>");
|
||||
let user = args.next().expect("no user given");
|
||||
let pass = args.next().expect("no pass given");
|
||||
*/
|
||||
//let network = Network::Regtest
|
||||
match get_client(network_params) {
|
||||
Ok((rpc, bcinfo)) => {
|
||||
info!("connected");
|
||||
//let best_block_hash = rpc.get_best_block_hash()?;
|
||||
//info!("best block hash: {}", best_block_hash);
|
||||
//let bestblockcount = rpc.get_block_count()?;
|
||||
//info!("best block height: {}", bestblockcount);
|
||||
//let best_block_hash_by_height = rpc.get_block_hash(bestblockcount)?;
|
||||
//info!("best block hash by height: {}", best_block_hash_by_height);
|
||||
//assert_eq!(best_block_hash_by_height, best_block_hash);
|
||||
//let from_block= std::cmp::max(0, bestblockcount - 11);
|
||||
//let mut time_sum:u64=0;
|
||||
//for i in from_block..bestblockcount{
|
||||
// let hash = rpc.get_block_hash(i).unwrap();
|
||||
// let block: bitcoin::Block = rpc.get_by_id(&hash).unwrap();
|
||||
// time_sum += <u32 as Into<u64>>::into(block.header.time);
|
||||
//}
|
||||
//let average_time = time_sum/11;
|
||||
info!("median time: {}", bcinfo.median_time);
|
||||
//info!("height time: {}",bcinfo.median_time);
|
||||
info!("blocks: {}", bcinfo.blocks);
|
||||
@@ -288,26 +265,10 @@ async fn main_result(cfg: &MyConfig, network_params: &NetworkParams) -> Result<(
|
||||
info!("to be pushed: {}: {}", txid, locktime);
|
||||
match rpc.send_raw_transaction(tx) {
|
||||
Ok(o) => {
|
||||
/*let mut file = OpenOptions::new()
|
||||
.append(true) // Set the append option
|
||||
.create(true) // Create the file if it doesn't exist
|
||||
.open("valid_txs")?;
|
||||
let data = format!("{}\t:\t{}\t:\t{}\n",txid,average_time,locktime);
|
||||
file.write_all(data.as_bytes())?;
|
||||
drop(file);
|
||||
*/
|
||||
info!("tx: {} pusshata PUSHED\n{}", txid, o);
|
||||
pushed_txs.push(txid.to_string());
|
||||
}
|
||||
Err(err) => {
|
||||
/*let mut file = OpenOptions::new()
|
||||
.append(true) // Set the append option
|
||||
.create(true) // Create the file if it doesn't exist
|
||||
.open("/home/bal/invalid_txs")?;
|
||||
let data = format!("{}:\t{}\t:\t{}\t:\t{}\n",txid,err,average_time,locktime);
|
||||
file.write_all(data.as_bytes())?;
|
||||
drop(file);
|
||||
*/
|
||||
warn!("Error: {}\n{}", err, txid);
|
||||
//store err in invalid_txs
|
||||
invalid_txs.insert(txid.to_string(), err.to_string());
|
||||
@@ -317,16 +278,37 @@ async fn main_result(cfg: &MyConfig, network_params: &NetworkParams) -> Result<(
|
||||
|
||||
for txid in &pushed_txs {
|
||||
let sql = "UPDATE tbl_tx SET status = 1 WHERE txid = ?";
|
||||
let mut stmt = db.prepare(sql).unwrap();
|
||||
stmt.bind((1, Value::String(txid.clone()))).unwrap();
|
||||
let _ = stmt.next();
|
||||
match db.prepare(sql) {
|
||||
Ok(mut stmt) => {
|
||||
if let Err(e) = stmt.bind((1, Value::String(txid.clone()))) {
|
||||
error!("Failed to bind txid for status update: {}", e);
|
||||
continue;
|
||||
}
|
||||
let _ = stmt.next();
|
||||
}
|
||||
Err(e) => {
|
||||
error!("Failed to prepare status update: {}", e);
|
||||
}
|
||||
}
|
||||
}
|
||||
for (txid, txerr) in &invalid_txs {
|
||||
let sql = "UPDATE tbl_tx SET status = 2, push_err = ? WHERE txid = ?";
|
||||
let mut stmt = db.prepare(sql).unwrap();
|
||||
stmt.bind((1, Value::String(txerr.clone()))).unwrap();
|
||||
stmt.bind((2, Value::String(txid.clone()))).unwrap();
|
||||
let _ = stmt.next();
|
||||
match db.prepare(sql) {
|
||||
Ok(mut stmt) => {
|
||||
if let Err(e) = stmt.bind((1, Value::String(txerr.clone()))) {
|
||||
error!("Failed to bind txerr for error update: {}", e);
|
||||
continue;
|
||||
}
|
||||
if let Err(e) = stmt.bind((2, Value::String(txid.clone()))) {
|
||||
error!("Failed to bind txid for error update: {}", e);
|
||||
continue;
|
||||
}
|
||||
let _ = stmt.next();
|
||||
}
|
||||
Err(e) => {
|
||||
error!("Failed to prepare error update: {}", e);
|
||||
}
|
||||
}
|
||||
}
|
||||
if let Err(e) = send_stats_report(cfg, bcinfo).await {
|
||||
error!("send_stats_report failed: {}", e);
|
||||
@@ -391,31 +373,6 @@ ON CONFLICT(chain) DO UPDATE SET
|
||||
"
|
||||
);
|
||||
|
||||
/*
|
||||
let sql = format!("CREATE TABLE tbl_stats AS
|
||||
SELECT
|
||||
CURRENT_TIMESTAMP AS report_date,
|
||||
'{chain}' as chain,
|
||||
(SELECT COUNT(*) FROM tbl_tx WHERE network ='{chain}') AS totals,
|
||||
(SELECT COUNT(*) FROM tbl_tx WHERE status = 0 AND network ='{chain}') AS waiting,
|
||||
(SELECT COUNT(*) FROM tbl_tx WHERE status = 1 AND network ='{chain}') AS sent,
|
||||
(SELECT COUNT(*) FROM tbl_tx WHERE status = 2 AND network ='{chain}') AS failed,
|
||||
(SELECT SUM(our_fees) FROM tbl_tx WHERE status = 0 AND network ='{chain}') AS waiting_profit,
|
||||
(SELECT SUM(our_fees) OR 0 FROM tbl_tx WHERE status = 1 AND network ='{chain}') AS sent_profit,
|
||||
(SELECT SUM(our_fees) FROM tbl_tx WHERE status = 2 AND network ='{chain}') AS missed_profit,
|
||||
(SELECT COUNT(*) FROM tbl_inp JOIN tbl_tx ON(tbl_inp.txid = tbl_tx.txid) WHERE tbl_tx.status=0 AND tbl_tx.network ='{chain}') AS unique_inputs;
|
||||
");
|
||||
let sql = "UPDATE tbl_stats set
|
||||
totals = (SELECT COUNT(*) FROM tbl_tx WHERE network ='{chain}'),
|
||||
waiting = (SELECT COUNT(*) FROM tbl_tx WHERE status = 0 AND network ='{chain}'),
|
||||
sent = (SELECT COUNT(*) FROM tbl_tx WHERE status = 1 AND network ='{chain}'),
|
||||
failed = (SELECT COUNT(*) FROM tbl_tx WHERE status = 1 AND network ='{chain}'),
|
||||
waiting_profit = (SELECT SUM(our_fees) FROM tbl_tx WHERE status = 0 AND network ='{chain}'),
|
||||
sent_profit = (SELECT SUM(our_fees) FROM tbl_tx WHERE status = 0 AND network ='{chain}'),
|
||||
missed_profit = (SELECT SUM(our_fees) FROM tbl_tx WHERE status = 0 AND network ='{chain}')
|
||||
unique_inputs = (SELECT COUNT(*) FROM tbl_inp JOIN tbl_tx ON(tbl_inp.txid = tbl_tx.txid) WHERE tbl_tx.status=0 AND tbl_tx.network ='{chain}')
|
||||
WHERE chain = '{chain}'
|
||||
*/
|
||||
if let Err(err) = db.execute(&sql) {
|
||||
error!("error inserting creating stats table {err}");
|
||||
} else {
|
||||
@@ -609,85 +566,12 @@ fn parse_env_netconfig(cfg_lock: &mut MyConfig, chain: &str) -> NetworkParams {
|
||||
if let Ok(value) = env::var(format!("BAL_PUSHER_{}_RPC_PASSWORD", chain.to_uppercase())) {
|
||||
cfg.rpc_pass = value;
|
||||
}
|
||||
println!("BAL_PUSHER_{}_ZMQ_HASHBLOCK", chain.to_uppercase());
|
||||
if let Ok(value) = env::var(format!("BAL_PUSHER_{}_ZMQ_HASHBLOCK", chain.to_uppercase())) {
|
||||
println!("value:{}", value);
|
||||
cfg.zmq_listener = value;
|
||||
}
|
||||
cfg.clone()
|
||||
}
|
||||
|
||||
#[allow(dead_code)]
|
||||
fn check_zmq_connection(endpoint: &str) -> bool {
|
||||
trace!("check zmq connection");
|
||||
let context = Context::new();
|
||||
let socket = match context.socket(DEALER) {
|
||||
Ok(sock) => sock,
|
||||
Err(_) => return false,
|
||||
};
|
||||
|
||||
if socket.connect(endpoint).is_err() {
|
||||
return false;
|
||||
}
|
||||
|
||||
// Try to send an empty message non-blocking
|
||||
socket.send("", DONTWAIT).is_ok()
|
||||
}
|
||||
|
||||
// Add this struct to monitor connection health
|
||||
#[allow(dead_code)]
|
||||
struct ConnectionMonitor {
|
||||
last_message_time: Instant,
|
||||
timeout: Duration,
|
||||
consecutive_timeouts: u32,
|
||||
max_consecutive_timeouts: u32,
|
||||
}
|
||||
|
||||
#[allow(dead_code)]
|
||||
impl ConnectionMonitor {
|
||||
fn new(timeout_secs: u64, max_timeouts: u32) -> Self {
|
||||
Self {
|
||||
last_message_time: Instant::now(),
|
||||
timeout: Duration::from_secs(timeout_secs),
|
||||
consecutive_timeouts: 0,
|
||||
max_consecutive_timeouts: max_timeouts,
|
||||
}
|
||||
}
|
||||
|
||||
fn update(&mut self) {
|
||||
self.last_message_time = Instant::now();
|
||||
self.consecutive_timeouts = 0;
|
||||
}
|
||||
|
||||
fn check_connection(&mut self) -> ConnectionStatus {
|
||||
let elapsed = self.last_message_time.elapsed();
|
||||
|
||||
if elapsed > self.timeout {
|
||||
self.consecutive_timeouts += 1;
|
||||
|
||||
if self.consecutive_timeouts >= self.max_consecutive_timeouts {
|
||||
ConnectionStatus::Lost(elapsed)
|
||||
} else {
|
||||
ConnectionStatus::Warning(elapsed)
|
||||
}
|
||||
} else {
|
||||
ConnectionStatus::Healthy
|
||||
}
|
||||
}
|
||||
|
||||
fn reset(&mut self) {
|
||||
self.consecutive_timeouts = 0;
|
||||
self.last_message_time = Instant::now();
|
||||
}
|
||||
}
|
||||
|
||||
#[allow(dead_code)]
|
||||
enum ConnectionStatus {
|
||||
Healthy,
|
||||
Warning(Duration),
|
||||
Lost(Duration),
|
||||
}
|
||||
|
||||
#[tokio::main]
|
||||
async fn main() -> std::io::Result<()> {
|
||||
env_logger::init();
|
||||
@@ -747,15 +631,15 @@ async fn main() -> std::io::Result<()> {
|
||||
Ok(m) => m,
|
||||
Err(e) => {
|
||||
consecutive_timeouts += 1;
|
||||
if consecutive_timeouts == 1 {
|
||||
warn!("ZMQ recv timeout or error: {}, retrying...", e);
|
||||
} else if consecutive_timeouts.is_multiple_of(12) {
|
||||
warn!(
|
||||
if consecutive_timeouts.is_multiple_of(720) {
|
||||
error!(
|
||||
"No ZMQ messages for {}s ({} consecutive timeouts), is bitcoind ZMQ active on {}?",
|
||||
consecutive_timeouts * 5,
|
||||
consecutive_timeouts,
|
||||
zmq_address
|
||||
);
|
||||
} else {
|
||||
trace!("ZMQ recv timeout or error: {}, retrying...", e);
|
||||
}
|
||||
continue;
|
||||
}
|
||||
@@ -783,17 +667,6 @@ async fn main() -> std::io::Result<()> {
|
||||
thread::sleep(Duration::from_millis(100)); // Sleep for 100ms
|
||||
}
|
||||
}
|
||||
#[allow(dead_code)]
|
||||
fn seq_to_str(seq: &[u8]) -> String {
|
||||
if seq.len() == 4 {
|
||||
let mut rdr = Cursor::new(seq);
|
||||
let sequence = rdr
|
||||
.read_u32::<LittleEndian>()
|
||||
.expect("Failed to read integer");
|
||||
return sequence.to_string();
|
||||
}
|
||||
"Unknown".to_string()
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
|
||||
@@ -118,7 +118,7 @@ pub struct StatsResponse {
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone)]
|
||||
#[allow(dead_code)]
|
||||
#[expect(dead_code)]
|
||||
struct ActixConfig {
|
||||
max_body_size: usize,
|
||||
timeout_secs: u64,
|
||||
@@ -217,6 +217,38 @@ async fn echo_version() -> impl Responder {
|
||||
HttpResponse::Ok().body(VERSION)
|
||||
}
|
||||
|
||||
fn is_valid_ip(ip: &str) -> bool {
|
||||
ip.parse::<std::net::IpAddr>().is_ok()
|
||||
}
|
||||
|
||||
fn extract_client_ip(req: &actix_web::HttpRequest) -> String {
|
||||
if let Some(val) = req.headers().get("X-Real-IP")
|
||||
&& let Ok(s) = val.to_str()
|
||||
{
|
||||
let ip = s.split(',').next().unwrap_or(s).trim();
|
||||
if is_valid_ip(ip) {
|
||||
debug!("client IP from X-Real-IP: {}", ip);
|
||||
return ip.to_string();
|
||||
}
|
||||
}
|
||||
if let Some(val) = req.headers().get("X-Forwarded-For")
|
||||
&& let Ok(s) = val.to_str()
|
||||
{
|
||||
let ip = s.split(',').next().unwrap_or(s).trim();
|
||||
if is_valid_ip(ip) {
|
||||
debug!("client IP from X-Forwarded-For: {}", ip);
|
||||
return ip.to_string();
|
||||
}
|
||||
}
|
||||
let fallback = req
|
||||
.connection_info()
|
||||
.peer_addr()
|
||||
.unwrap_or("unknown")
|
||||
.to_string();
|
||||
debug!("client IP from peer_addr fallback: {}", fallback);
|
||||
fallback
|
||||
}
|
||||
|
||||
async fn echo_info(
|
||||
path: web::Path<String>,
|
||||
data: web::Data<AppState>,
|
||||
@@ -232,18 +264,7 @@ async fn echo_info(
|
||||
debug!("network disabled {}", param);
|
||||
return HttpResponse::BadRequest().body("error");
|
||||
}
|
||||
let remote_addr = req
|
||||
.headers()
|
||||
.get("X-Real-IP")
|
||||
.and_then(|value| value.to_str().ok())
|
||||
.and_then(|xff| xff.split(',').next())
|
||||
.map(|ip| ip.trim().to_string())
|
||||
.unwrap_or_else(|| {
|
||||
req.connection_info()
|
||||
.peer_addr()
|
||||
.unwrap_or("unknown")
|
||||
.to_string()
|
||||
});
|
||||
let remote_addr = extract_client_ip(&req);
|
||||
let address = match netconfig.xpub {
|
||||
false => {
|
||||
let address = netconfig.address.to_string();
|
||||
@@ -413,13 +434,8 @@ async fn echo_stats(path: web::Path<String>, data: web::Data<AppState>) -> impl
|
||||
unique_inputs,
|
||||
});
|
||||
}
|
||||
match serde_json::to_string(&stats) {
|
||||
Ok(json_data) => {
|
||||
debug!("echo info reply: {}", json_data);
|
||||
HttpResponse::Ok().json(stats)
|
||||
}
|
||||
Err(_err) => HttpResponse::InternalServerError().body("error"),
|
||||
}
|
||||
debug!("echo stats reply for chain: {}", netconfig.name);
|
||||
HttpResponse::Ok().json(stats)
|
||||
}
|
||||
|
||||
async fn echo_search(body: Bytes, data: web::Data<AppState>) -> impl Responder {
|
||||
@@ -531,7 +547,6 @@ fn parse_request_transactions(
|
||||
known_addresses: &HashSet<String>,
|
||||
) -> Vec<(ParsedTx, String, u64)> {
|
||||
let mut result: Vec<(ParsedTx, String, u64)> = Vec::new();
|
||||
let mut union_tx = true;
|
||||
|
||||
for line in strbody.split('\n') {
|
||||
if line.is_empty() {
|
||||
@@ -615,11 +630,6 @@ fn parse_request_transactions(
|
||||
trace!("willexecutor output not found for tx {}, skipping", txid);
|
||||
continue;
|
||||
}
|
||||
if !union_tx {
|
||||
// This is only used for SQL building later; we track it in the caller
|
||||
} else {
|
||||
union_tx = false;
|
||||
}
|
||||
result.push((
|
||||
ParsedTx {
|
||||
txid,
|
||||
@@ -911,15 +921,6 @@ async fn main() -> std::io::Result<()> {
|
||||
cfg: cfg.clone(),
|
||||
});
|
||||
|
||||
// Initialize networks
|
||||
{
|
||||
let db = data.db.lock().unwrap();
|
||||
for network in NETWORKS {
|
||||
let netconfig = data.cfg.get_net_config(network);
|
||||
insert_xpub(&db, &netconfig.name.to_string(), &netconfig.address);
|
||||
}
|
||||
}
|
||||
|
||||
let bind_address = data.cfg.bind_address.clone();
|
||||
let bind_port = data.cfg.bind_port;
|
||||
|
||||
|
||||
17
src/db.rs
17
src/db.rs
@@ -164,7 +164,7 @@ pub fn create_database(db: &Connection) {
|
||||
let _ = db.execute("DROP INDEX IF EXISTS idx_stats_chain;");
|
||||
let _ = db.execute("CREATE UNIQUE INDEX IF NOT EXISTS idx_stats_chain ON tbl_stats(chain);");
|
||||
|
||||
let _ = db.execute("UPDATE tbl_tx set network='bitcoin' where network='mainnet');");
|
||||
let _ = db.execute("UPDATE tbl_tx set network='bitcoin' where network='mainnet';");
|
||||
}
|
||||
/*
|
||||
pub fn get_xpub_id(db: &Connection, network: &String, xpub: &String) -> Option<i64>{
|
||||
@@ -181,13 +181,14 @@ pub fn create_database(db: &Connection) {
|
||||
pub fn insert_xpub(db: &Connection, network: &str, xpub: &str) {
|
||||
if !xpub.is_empty() {
|
||||
trace!("going to insert: {} xpub:{}", network, xpub);
|
||||
let mut stmt = match db.prepare("INSERT INTO tbl_xpub(network,xpub) VALUES(?, ?);") {
|
||||
Ok(s) => s,
|
||||
Err(e) => {
|
||||
error!("Failed to prepare xpub insert statement: {}", e);
|
||||
return;
|
||||
}
|
||||
};
|
||||
let mut stmt =
|
||||
match db.prepare("INSERT OR IGNORE INTO tbl_xpub(network,xpub) VALUES(?, ?);") {
|
||||
Ok(s) => s,
|
||||
Err(e) => {
|
||||
error!("Failed to prepare xpub insert statement: {}", e);
|
||||
return;
|
||||
}
|
||||
};
|
||||
if let Err(e) = stmt.bind((1, Value::String(network.to_string()))) {
|
||||
error!("Failed to bind network parameter for xpub insert: {}", e);
|
||||
return;
|
||||
|
||||
Reference in New Issue
Block a user