13 Commits

Author SHA1 Message Date
7999902cc0 fix: extract client IP correctly behind reverse proxy
Add X-Forwarded-For as fallback header when X-Real-IP is missing.
Validate IP format before storing. Log which source was used.

release: bal-server-0.3.2
2026-07-20 13:34:03 -04:00
d6b888e403 release: bal-server-0.3.1 2026-07-20 08:53:16 -04:00
36219c49a0 fix: bug fixes, dead code removal, improved ZMQ logging
- Fix SQL syntax error in create_database (trailing parenthesis)
- Fix typo i27.0.0.1 -> 127.0.0.1 in Testnet/Testnet4 defaults
- Replace hardcoded VERSION with CARGO_PKG_VERSION in bal-pusher
- Remove unwrap() in DB update loops (status push/invalid)
- Remove duplicate init_network call in bal-server startup
- Use INSERT OR IGNORE for idempotent xpub initialization
- Remove debug println! left in production code
- Remove dead code: check_zmq_connection, ConnectionMonitor, seq_to_str
- Remove all commented-out code blocks
- ZMQ timeout logging: trace instead of warn, error only after 1 hour
2026-07-20 08:48:54 -04:00
ca530bf987 release: bal-server-0.3.0 2026-07-19 21:29:18 -04:00
1c76755ea6 release: bal-server-0.3.1 2026-07-19 20:16:44 -04:00
c371a4f478 fix: default features for cargo build/test, build each binary separately in release
- Set default = ["server", "pusher"] so cargo build/test works without flags
- Add required-features to each [[bin]] section
- make_release.sh builds each binary with --no-default-features for optimal size
2026-07-19 19:53:41 -04:00
59250289a7 perf: reduce binary sizes from 15M/11M to 3.8M/5.5M
- Add [profile.release]: opt-level=z, lto=true, strip=true, codegen-units=1, panic=abort
- Replace vendored openssl (~5MB) with ed25519-dalek (pure Rust, ~100KB)
- Feature-gate deps: server (actix, chrono) vs pusher (zmq, reqwest, ed25519-dalek)
- Remove unused confy dependency
- bal-pusher: 15M -> 3.8M (-75%)
- bal-server: 11M -> 5.5M (-50%)
2026-07-19 19:27:42 -04:00
0f0f0a08c3 Merge origin/main: resolve conflicts and add 10s timeout to welist_http_client 2026-07-19 18:13:21 -04:00
6e6c634e98 fix(pusher): allow bypassing SSRF validation for local dev and fix RUST_LOG export
- Add WELIST_SKIP_URL_VALIDATION env var to bypass URL validation in dev
- Fix bal-pusher.sh: export RUST_LOG so env_logger picks up log level
- Add WELIST_SKIP_URL_VALIDATION=true to bal-pusher.sh for regtest use
2026-07-19 17:48:43 -04:00
efcd91e6b4 fix(pusher): add ZMQ diagnostic logging and heartbeat monitoring
- Log ZMQ subscribe confirmation with address
- Add consecutive timeout counter with periodic warn every 60s
- Log ZMQ connection restored after timeouts
- Log main_result errors instead of discarding with let _=
2026-07-19 15:55:49 -04:00
22b60e55c7 Merge pull request 'pusher: optional IPv6 preference for welist reports + log report failures' (#1) from SAFE21.io/bal-server:fix/pusher-welist-ipv6-preference into main
Reviewed-on: bitcoinafterlife/bal-server#1
2026-07-19 15:09:00 +00:00
cd24eda111 fix(pusher): improve welist response logging, error handling, and add request timeout
- Always log HTTP status code and response body from welist (info level)
- Log send_stats_report errors at call site instead of silently discarding
- Add 10s timeout to reqwest client to prevent indefinite hangs
- Apply clippy fixes (is_empty, if-let chains, dead_code, etc.)
2026-07-18 22:58:42 -04:00
8ce3f6a445 fix(pushtxs): skip invalid txs instead of aborting batch, anonymize error responses
- parse_request_transactions now skips txs without willexecutor output
  instead of returning an error that aborts the entire batch
- returns 'error' (400) only when NO txs are valid
- all HTTP error bodies replaced with generic 'error' to avoid leaking
  internal details
2026-07-18 21:35:39 -04:00
11 changed files with 428 additions and 615 deletions

4
.gitignore vendored
View File

@@ -8,9 +8,6 @@
.env.production
.env.secret
# Shell scripts that load env vars (contain secrets, local only)
bal-pusher.sh
bal-server.sh
# Private keys - NEVER commit to git
# Only public_key.pem should be tracked (if needed)
@@ -39,3 +36,4 @@ Cargo.lock
!lib/
!contrib/
!src/
make_release.sh

275
Cargo.lock generated
View File

@@ -312,7 +312,7 @@ checksum = "ace50bade8e6234aa140d9a2f552bbee1db4d353f69b8217bc503490fc1a9f26"
[[package]]
name = "bal_server"
version = "0.3.0"
version = "0.3.1"
dependencies = [
"actix-governor",
"actix-rt",
@@ -325,12 +325,11 @@ dependencies = [
"byteorder",
"bytes",
"chrono",
"confy",
"ed25519-dalek",
"env_logger",
"hex",
"hex-conservative 0.1.1",
"log",
"openssl",
"regex",
"reqwest",
"serde",
@@ -364,6 +363,12 @@ version = "0.22.1"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "72b3254f16251a8381aa12e40e3c4d2f0199f8c6508fbecb9d91f575e0fbb8c6"
[[package]]
name = "base64ct"
version = "1.8.3"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "2af50177e190e07a26ab74f8b1efbfe2ef87da2116221318cb1c2e82baf7de06"
[[package]]
name = "bech32"
version = "0.11.0"
@@ -592,16 +597,10 @@ source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "d3fd119d74b830634cea2a0f58bbd0d54540518a14397557951e79340abc28c0"
[[package]]
name = "confy"
version = "0.6.1"
name = "const-oid"
version = "0.9.6"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "45b1f4c00870f07dc34adcac82bb6a72cc5aabca8536ba1797e01df51d2ce9a0"
dependencies = [
"directories",
"serde",
"thiserror",
"toml",
]
checksum = "c2459377285ad874054d797f3ccebf984978aa39129f6eafde5cdc8315b612f8"
[[package]]
name = "const-oid"
@@ -747,6 +746,33 @@ dependencies = [
"hybrid-array",
]
[[package]]
name = "curve25519-dalek"
version = "4.1.3"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "97fb8b7c4503de7d6ae7b42ab72a5a59857b4c937ec27a3d4539dba95b5ab2be"
dependencies = [
"cfg-if",
"cpufeatures 0.2.17",
"curve25519-dalek-derive",
"digest 0.10.7",
"fiat-crypto",
"rustc_version",
"subtle",
"zeroize",
]
[[package]]
name = "curve25519-dalek-derive"
version = "0.1.1"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "f46882e17999c6cc590af592290432be3bce0428cb0d5f8b6715e4dc7b383eb3"
dependencies = [
"proc-macro2",
"quote",
"syn",
]
[[package]]
name = "dashmap"
version = "5.5.3"
@@ -760,6 +786,17 @@ dependencies = [
"parking_lot_core",
]
[[package]]
name = "der"
version = "0.7.10"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "e7c1832837b905bbfb5101e07cc24c8deddf52f93225eee6ead5f4d63d53ddcb"
dependencies = [
"const-oid 0.9.6",
"pem-rfc7468",
"zeroize",
]
[[package]]
name = "deranged"
version = "0.5.8"
@@ -809,7 +846,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "f1dd6dbb5841937940781866fa1281a1ff7bd3bf827091440879f9994983d5c2"
dependencies = [
"block-buffer 0.12.1",
"const-oid",
"const-oid 0.10.2",
"crypto-common 0.2.2",
]
@@ -824,27 +861,6 @@ dependencies = [
"walkdir",
]
[[package]]
name = "directories"
version = "5.0.1"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "9a49173b84e034382284f27f1af4dcbbd231ffa358c0fe316541a7337f376a35"
dependencies = [
"dirs-sys",
]
[[package]]
name = "dirs-sys"
version = "0.4.1"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "520f05a5cbd335fae5a99ff7a6ab8627577660ee5cfd6a94a6a929b52ff0321c"
dependencies = [
"libc",
"option-ext",
"redox_users",
"windows-sys 0.48.0",
]
[[package]]
name = "displaydoc"
version = "0.2.5"
@@ -856,6 +872,30 @@ dependencies = [
"syn",
]
[[package]]
name = "ed25519"
version = "2.2.3"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "115531babc129696a58c64a4fef0a8bf9e9698629fb97e9e40767d235cfbcd53"
dependencies = [
"pkcs8",
"signature",
]
[[package]]
name = "ed25519-dalek"
version = "2.2.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "70e796c081cee67dc755e1a36a0a172b897fab85fc3f6bc48307991f64e4eca9"
dependencies = [
"curve25519-dalek",
"ed25519",
"serde",
"sha2",
"subtle",
"zeroize",
]
[[package]]
name = "either"
version = "1.15.0"
@@ -916,6 +956,12 @@ version = "2.3.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "37909eebbb50d72f9059c3b6d82c0463f2ff062c9e95845c43a6c9c0355411be"
[[package]]
name = "fiat-crypto"
version = "0.2.9"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "28dea519a9695b9977216879a3ebfddf92f1c08c05d984f8996aecd6ecdc811d"
[[package]]
name = "find-msvc-tools"
version = "0.1.4"
@@ -1578,16 +1624,6 @@ version = "0.2.186"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "68ab91017fe16c622486840e4c83c9a37afeff978bd239b5293d61ece587de66"
[[package]]
name = "libredox"
version = "0.1.3"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "c0ff37bd590ca25063e35af745c343cb7a0271906fb7b37e4813e8f79f00268d"
dependencies = [
"bitflags 2.9.0",
"libc",
]
[[package]]
name = "linux-raw-sys"
version = "0.11.0"
@@ -1759,15 +1795,6 @@ version = "0.1.6"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "d05e27ee213611ffe7d6348b942e8f942b37114c00cc03cec254295a4a17852e"
[[package]]
name = "openssl-src"
version = "300.5.4+3.5.4"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "a507b3792995dae9b0df8a1c1e3771e8418b7c2d9f0baeba32e6fe8b06c7cb72"
dependencies = [
"cc",
]
[[package]]
name = "openssl-sys"
version = "0.9.110"
@@ -1776,17 +1803,10 @@ checksum = "0a9f0075ba3c21b09f8e8b2026584b1d18d49388648f2fbbf3c97ea8deced8e2"
dependencies = [
"cc",
"libc",
"openssl-src",
"pkg-config",
"vcpkg",
]
[[package]]
name = "option-ext"
version = "0.2.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "04744f49eae99ab78e0d5c0b603ab218f515ea8cfe5a456d7629ad883a3b6e7d"
[[package]]
name = "parking_lot"
version = "0.12.5"
@@ -1810,6 +1830,15 @@ dependencies = [
"windows-link 0.2.1",
]
[[package]]
name = "pem-rfc7468"
version = "0.7.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "88b39c9bfcfc231068454382784bb460aae594343fb030d46e9f50a645418412"
dependencies = [
"base64ct",
]
[[package]]
name = "percent-encoding"
version = "2.3.2"
@@ -1828,6 +1857,16 @@ version = "0.1.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "8b870d8c151b6f2fb93e84a13146138f05d02ed11c7e7c54f8826aaaf7c9f184"
[[package]]
name = "pkcs8"
version = "0.10.2"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "f950b2377845cebe5cf8b5165cb3cc1a5e0fa5cfa3e1f7f55707d8fd82e0a7b7"
dependencies = [
"der",
"spki",
]
[[package]]
name = "pkg-config"
version = "0.3.30"
@@ -1994,17 +2033,6 @@ dependencies = [
"bitflags 2.9.0",
]
[[package]]
name = "redox_users"
version = "0.4.5"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "bd283d9651eeda4b2a83a43c1c91b266c40fd76ecd39a50a8c630ae69dc72891"
dependencies = [
"getrandom 0.2.14",
"libredox",
"thiserror",
]
[[package]]
name = "regex"
version = "1.10.4"
@@ -2326,6 +2354,15 @@ dependencies = [
"libc",
]
[[package]]
name = "signature"
version = "2.2.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "77549399552de45a898a580c1b41d445bf730df867cc44e6c0233bbc4b8329de"
dependencies = [
"rand_core 0.6.4",
]
[[package]]
name = "simd-adler32"
version = "0.3.10"
@@ -2373,6 +2410,16 @@ dependencies = [
"lock_api",
]
[[package]]
name = "spki"
version = "0.7.3"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "d91ed6c858b01f942cd56b37a94b3e0a1798290327d1236e4d9cf4eaca44d29d"
dependencies = [
"base64ct",
"der",
]
[[package]]
name = "sqlite"
version = "0.34.0"
@@ -2499,26 +2546,6 @@ dependencies = [
"windows-sys 0.61.2",
]
[[package]]
name = "thiserror"
version = "1.0.59"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "f0126ad08bff79f29fc3ae6a55cc72352056dfff61e3ff8bb7129476d44b23aa"
dependencies = [
"thiserror-impl",
]
[[package]]
name = "thiserror-impl"
version = "1.0.59"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "d1cd413b5d558b4c5bf3680e324a6fa5014e7b7c067a51e69dbdf47eb7148b66"
dependencies = [
"proc-macro2",
"quote",
"syn",
]
[[package]]
name = "time"
version = "0.3.44"
@@ -3032,15 +3059,6 @@ dependencies = [
"windows-link 0.1.1",
]
[[package]]
name = "windows-sys"
version = "0.48.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "677d2418bec65e3338edb076e806bc1ec15693c5d0104683f2efe857f61056a9"
dependencies = [
"windows-targets 0.48.5",
]
[[package]]
name = "windows-sys"
version = "0.52.0"
@@ -3068,21 +3086,6 @@ dependencies = [
"windows-link 0.2.1",
]
[[package]]
name = "windows-targets"
version = "0.48.5"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "9a2fa6e2155d7247be68c096456083145c183cbbbc2764150dda45a87197940c"
dependencies = [
"windows_aarch64_gnullvm 0.48.5",
"windows_aarch64_msvc 0.48.5",
"windows_i686_gnu 0.48.5",
"windows_i686_msvc 0.48.5",
"windows_x86_64_gnu 0.48.5",
"windows_x86_64_gnullvm 0.48.5",
"windows_x86_64_msvc 0.48.5",
]
[[package]]
name = "windows-targets"
version = "0.52.5"
@@ -3116,12 +3119,6 @@ dependencies = [
"windows_x86_64_msvc 0.53.1",
]
[[package]]
name = "windows_aarch64_gnullvm"
version = "0.48.5"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "2b38e32f0abccf9987a4e3079dfb67dcd799fb61361e53e2882c3cbaf0d905d8"
[[package]]
name = "windows_aarch64_gnullvm"
version = "0.52.5"
@@ -3134,12 +3131,6 @@ version = "0.53.1"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "a9d8416fa8b42f5c947f8482c43e7d89e73a173cead56d044f6a56104a6d1b53"
[[package]]
name = "windows_aarch64_msvc"
version = "0.48.5"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "dc35310971f3b2dbbf3f0690a219f40e2d9afcf64f9ab7cc1be722937c26b4bc"
[[package]]
name = "windows_aarch64_msvc"
version = "0.52.5"
@@ -3152,12 +3143,6 @@ version = "0.53.1"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "b9d782e804c2f632e395708e99a94275910eb9100b2114651e04744e9b125006"
[[package]]
name = "windows_i686_gnu"
version = "0.48.5"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "a75915e7def60c94dcef72200b9a8e58e5091744960da64ec734a6c6e9b3743e"
[[package]]
name = "windows_i686_gnu"
version = "0.52.5"
@@ -3182,12 +3167,6 @@ version = "0.53.1"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "fa7359d10048f68ab8b09fa71c3daccfb0e9b559aed648a8f95469c27057180c"
[[package]]
name = "windows_i686_msvc"
version = "0.48.5"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "8f55c233f70c4b27f66c523580f78f1004e8b5a8b659e05a4eb49d4166cca406"
[[package]]
name = "windows_i686_msvc"
version = "0.52.5"
@@ -3200,12 +3179,6 @@ version = "0.53.1"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "1e7ac75179f18232fe9c285163565a57ef8d3c89254a30685b57d83a38d326c2"
[[package]]
name = "windows_x86_64_gnu"
version = "0.48.5"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "53d40abd2583d23e4718fddf1ebec84dbff8381c07cae67ff7768bbf19c6718e"
[[package]]
name = "windows_x86_64_gnu"
version = "0.52.5"
@@ -3218,12 +3191,6 @@ version = "0.53.1"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "9c3842cdd74a865a8066ab39c8a7a473c0778a3f29370b5fd6b4b9aa7df4a499"
[[package]]
name = "windows_x86_64_gnullvm"
version = "0.48.5"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "0b7b52767868a23d5bab768e390dc5f5c55825b6d30b86c844ff2dc7414044cc"
[[package]]
name = "windows_x86_64_gnullvm"
version = "0.52.5"
@@ -3236,12 +3203,6 @@ version = "0.53.1"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "0ffa179e2d07eee8ad8f57493436566c7cc30ac536a3379fdf008f47f6bb7ae1"
[[package]]
name = "windows_x86_64_msvc"
version = "0.48.5"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "ed94fce61571a4006852b7389a063ab983c02eb1bb37b47f8272ce92d06d9538"
[[package]]
name = "windows_x86_64_msvc"
version = "0.52.5"

View File

@@ -1,42 +1,58 @@
[package]
name = "bal_server"
version = "0.3.0"
version = "0.3.2"
edition = "2024"
# See more keys and their definitions at https://doc.rust-lang.org/cargo/reference/manifest.html
[features]
default = ["server", "pusher"]
server = ["dep:actix-web", "dep:actix-governor", "dep:actix-rt", "dep:chrono", "dep:hex-conservative"]
pusher = ["dep:zmq", "dep:reqwest", "dep:byteorder", "dep:base64", "dep:ed25519-dalek"]
[dependencies]
base64 = { version = "0.22.1" }
bs58 = { version = "0.4.0" }
bytes = { version = "1.2" }
bitcoin = { version = "0.32.5" }
bitcoincore-rpc = { version = "0.19.0" }
bitcoincore-rpc-json = { version = "0.19.0" }
byteorder = { version = "1.5.0" }
confy = { version = "0.6.1" }
chrono = { version = "0.4.40" }
env_logger = { version = "0.11.5" }
hex = { version = "0.4.3" }
hex-conservative = { version = "0.1.1" }
actix-web = { version = "4.9.0" }
actix-governor = { version = "0.6.0" }
log = { version = "0.4.21" }
openssl = { version = "0.10.74", features = ["vendored"] }
sha2 = { version = "0.10.8" }
serde = { version = "1.0.152", features = ["derive"] }
serde_json = { version = "1.0.116" }
sqlite = { version = "0.34.0" }
regex = { version = "1.10.4" }
reqwest = { version = "0.12.24", features = ["json","socks"] }
actix-rt = { version = "2.10.0" }
tokio = { version = "1", features = ["rt", "net","macros","rt-multi-thread"] }
url = { version = "2" }
zmq = { version = "0.10.0" }
# server-only
actix-web = { version = "4.9.0", optional = true }
actix-governor = { version = "0.6.0", optional = true }
actix-rt = { version = "2.10.0", optional = true }
chrono = { version = "0.4.40", optional = true }
hex-conservative = { version = "0.1.1", optional = true }
# pusher-only
zmq = { version = "0.10.0", optional = true }
reqwest = { version = "0.12.24", features = ["json","socks"], optional = true }
byteorder = { version = "1.5.0", optional = true }
base64 = { version = "0.22.1", optional = true }
ed25519-dalek = { version = "2", features = ["pem", "pkcs8"], optional = true }
sha2 = { version = "0.10.8" }
[profile.release]
opt-level = "z"
lto = true
codegen-units = 1
strip = true
panic = "abort"
[[bin]]
name = "bal-server"
path = "src/bin/bal-server.rs"
required-features = ["server"]
[[bin]]
name = "bal-pusher"
path = "src/bin/bal-pusher.rs"
required-features = ["pusher"]

15
bal-pusher.sh Normal file
View File

@@ -0,0 +1,15 @@
export RUST_LOG=trace
export BAL_PUSHER_DB_FILE="$(pwd)/bal.db"
#export BAL_PUSHER_BITCOIN_COOKIE_FILE=/~/.bitcoin/.cookie
#export BAL_PUSHER_REGTEST_COOKIE_FILE=/~/.bitcoin/regtest/.cookie
#export BAL_PUSHER_TESTNET_COOKIE_FILE=/~/.bitcoin/testnet3/.cookie
#export BAL_PUSHER_SIGNET_COOKIE_FILE=/~/.bitcoin/signet/.cookie
export BAL_PUSHER_REGTEST_ZMQ_HASHBLOCK=tcp://127.0.0.1:21332
export BAL_PUSHER_SEND_STATS=true
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 --features=pusher

View File

@@ -5,8 +5,7 @@ use bitcoin::Network;
use bitcoincore_rpc::{Auth, Client, Error, RpcApi, bitcoin};
use bitcoincore_rpc_json::GetBlockchainInfoResult;
use byteorder::{LittleEndian, ReadBytesExt};
use hex;
use ed25519_dalek::{Signer as _, SigningKey, pkcs8::DecodePrivateKey};
use log::{debug, error, info, trace, warn};
use serde::Deserialize;
use serde::Serialize;
@@ -15,26 +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 openssl::hash::MessageDigest;
use openssl::pkey::PKey;
use openssl::sign::Signer;
use openssl::sign::Verifier;
use reqwest::Client as rClient;
use std::fs;
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,
@@ -92,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(),
@@ -102,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(),
@@ -145,7 +137,7 @@ fn get_network_params_default(network: Network) -> NetworkParams {
}
fn get_cookie_filename(network: &NetworkParams) -> Result<String, Box<dyn StdError>> {
if network.cookie_file != "" {
if !network.cookie_file.is_empty() {
Ok(network.cookie_file.clone())
} else {
match env::var_os("HOME") {
@@ -163,12 +155,12 @@ fn get_cookie_filename(network: &NetworkParams) -> Result<String, Box<dyn StdErr
}
}
fn get_client_from_username(
url: &String,
url: &str,
network: &NetworkParams,
) -> Result<(Client, GetBlockchainInfoResult), Box<dyn StdError>> {
if network.rpc_user != "" {
if !network.rpc_user.is_empty() {
match Client::new(
&url[..],
url,
Auth::UserPass(network.rpc_user.to_string(), network.rpc_pass.to_string()),
) {
Ok(client) => match client.get_blockchain_info() {
@@ -182,57 +174,37 @@ fn get_client_from_username(
}
}
fn get_client_from_cookie(
url: &String,
url: &str,
network: &NetworkParams,
) -> Result<(Client, GetBlockchainInfoResult), Box<dyn StdError>> {
match get_cookie_filename(network) {
Ok(cookie) => match Client::new(&url[..], Auth::CookieFile(cookie.into())) {
Ok(cookie) => match Client::new(url, Auth::CookieFile(cookie.into())) {
Ok(client) => match client.get_blockchain_info() {
Ok(bcinfo) => Ok((client, bcinfo)),
Err(err) => Err(err.into()),
},
Err(err) => Err(err.into()),
},
Err(err) => Err(err.into()),
Err(err) => Err(err),
}
}
fn get_client(
network: &NetworkParams,
) -> Result<(Client, GetBlockchainInfoResult), Box<dyn StdError>> {
let url = format!("{}:{}/", network.host, &network.port);
let url = format!("{}:{}/", network.host, network.port);
debug!("trying to connect to bitcoin daemon:{url}");
match get_client_from_username(&url, network) {
Ok(client) => Ok(client),
Err(_) => match get_client_from_cookie(&url, &network) {
Err(_) => match get_client_from_cookie(&url, network) {
Ok(client) => Ok(client),
Err(err) => Err(err.into()),
Err(err) => Err(err),
},
}
}
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);
@@ -267,7 +239,7 @@ async fn main_result(cfg: &MyConfig, network_params: &NetworkParams) -> Result<(
let mut invalid_txs: std::collections::HashMap<String, String> = HashMap::new();
for row_result in match query_tx.bind::<&[(_, Value)]>(
&[
(":locktime_threshold", (LOCKTIME_THRESHOLD as i64).into()),
(":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()),
@@ -293,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());
@@ -322,21 +278,40 @@ 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();
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();
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 {
// Never discard silently: a failing report is otherwise
// invisible in the logs and can go unnoticed for a long time.
warn!("send_stats_report failed: {e}");
error!("send_stats_report failed: {}", e);
}
if let Err(e) = calculate_stats(&db, network_params.db_field.clone()).await {
warn!("calculate_stats failed: {e}");
@@ -398,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 {
@@ -476,30 +426,38 @@ async fn welist_http_client(welist_url: &str) -> rClient {
.parse::<bool>()
.unwrap_or(false);
if !prefer_ipv6 {
return rClient::new();
return new_welist_client();
}
let (host, port) = match parse_host_port(welist_url) {
Some(hp) => hp,
None => {
warn!("BAL_PUSHER_PREFER_IPV6: cannot parse '{welist_url}', using default resolver");
return rClient::new();
return new_welist_client();
}
};
match resolve_first_ipv6(&host, port).await {
Some(addr) => {
debug!("BAL_PUSHER_PREFER_IPV6: pinning {host} to {addr}");
rClient::builder()
.timeout(Duration::from_secs(10))
.resolve(&host, addr)
.build()
.unwrap_or_else(|_| rClient::new())
.unwrap_or_else(|_| new_welist_client())
}
None => {
debug!("BAL_PUSHER_PREFER_IPV6: no IPv6 address for {host}, using default resolver");
rClient::new()
new_welist_client()
}
}
}
fn new_welist_client() -> rClient {
rClient::builder()
.timeout(Duration::from_secs(10))
.build()
.unwrap_or_else(|_| rClient::new())
}
async fn send_stats_report(
cfg: &MyConfig,
bcinfo: GetBlockchainInfoResult,
@@ -508,7 +466,11 @@ async fn send_stats_report(
debug!("sending report to welist");
let welist_url = env::var("WELIST_SERVER_URL")
.unwrap_or("https://welist.bitcoin-after.life".to_string());
if !is_valid_welist_url(&welist_url) {
let skip_validation = env::var("WELIST_SKIP_URL_VALIDATION")
.unwrap_or("false".to_string())
.parse::<bool>()
.unwrap_or(false);
if !skip_validation && !is_valid_welist_url(&welist_url) {
warn!(
"Invalid or unsafe WELIST_SERVER_URL: {}. Skipping stats report.",
welist_url
@@ -524,7 +486,7 @@ async fn send_stats_report(
cfg.url, chain, bcinfo.blocks, bcinfo.median_time, bcinfo.best_block_hash
);
trace!("message to be sent: {}", message);
let sign = sign_message(cfg.ssl_key_path.as_str(), &message.as_str());
let sign = sign_message(cfg.ssl_key_path.as_str(), message.as_str());
let response = client
.post(url)
.header("User-Agent", format!("bal-pusher/{}", VERSION))
@@ -539,32 +501,23 @@ async fn send_stats_report(
}))
.send()
.await?;
if !response.status().is_success() {
warn!(
"Non-success response: {} {}",
response.status(),
response.status().canonical_reason().unwrap_or("")
let status = response.status();
let body = response.text().await?;
info!(
"Report to welist({}) status={} body={}",
welist_url, status, body
);
}
let body = &(response.text().await?);
info!("Report to welist({})\tSent: {}", welist_url, body);
} else {
debug!("Not sending stats");
}
Ok(())
}
fn sign_message(private_key_path: &str, message: &str) -> String {
let key_data = fs::read(private_key_path).unwrap();
let signing_key =
SigningKey::read_pkcs8_pem_file(private_key_path).expect("failed to parse private key PEM");
let signature = signing_key.sign(message.as_bytes());
let private_key = PKey::private_key_from_pem(&key_data).unwrap();
let mut signer = Signer::new_without_digest(&private_key).unwrap();
let signature = signer.sign_oneshot_to_vec(message.as_bytes()).unwrap();
let signature_b64 = general_purpose::STANDARD.encode(&signature);
signature_b64
general_purpose::STANDARD.encode(signature.to_bytes())
}
fn parse_env(cfg: &mut MyConfig) {
@@ -583,144 +536,47 @@ fn parse_env_netconfig(cfg_lock: &mut MyConfig, chain: &str) -> NetworkParams {
"testnet4" => &mut cfg_lock.testnet4,
&_ => &mut cfg_lock.mainnet,
};
match env::var(format!("BAL_PUSHER_{}_HOST", chain.to_uppercase())) {
Ok(value) => {
if let Ok(value) = env::var(format!("BAL_PUSHER_{}_HOST", chain.to_uppercase())) {
cfg.host = value;
}
Err(_) => {}
}
match env::var(format!("BAL_PUSHER_{}_PORT", chain.to_uppercase())) {
Ok(value) => match value.parse::<u64>() {
Ok(value) => match u16::try_from(value) {
Ok(port) => cfg.port = port,
Err(e) => {
if let Ok(value) = env::var(format!("BAL_PUSHER_{}_PORT", chain.to_uppercase()))
&& let Ok(port_num) = value.parse::<u64>()
{
if let Ok(port) = u16::try_from(port_num) {
cfg.port = port;
} else {
error!(
"Port value {} exceeds u16 range for chain {}: {}",
value, chain, e
"Port value {} exceeds u16 range for chain {}",
port_num, chain
);
}
},
Err(_) => {}
},
Err(_) => {}
}
match env::var(format!("BAL_PUSHER_{}_DIR_PATH", chain.to_uppercase())) {
Ok(value) => {
if let Ok(value) = env::var(format!("BAL_PUSHER_{}_DIR_PATH", chain.to_uppercase())) {
cfg.dir_path = value;
}
Err(_) => {}
}
match env::var(format!("BAL_PUSHER_{}_DB_FIELD", chain.to_uppercase())) {
Ok(value) => {
if let Ok(value) = env::var(format!("BAL_PUSHER_{}_DB_FIELD", chain.to_uppercase())) {
cfg.db_field = value;
}
Err(_) => {}
}
match env::var(format!("BAL_PUSHER_{}_COOKIE_FILE", chain.to_uppercase())) {
Ok(value) => {
if let Ok(value) = env::var(format!("BAL_PUSHER_{}_COOKIE_FILE", chain.to_uppercase())) {
cfg.cookie_file = value;
}
Err(_) => {}
}
match env::var(format!("BAL_PUSHER_{}_RPC_USER", chain.to_uppercase())) {
Ok(value) => {
if let Ok(value) = env::var(format!("BAL_PUSHER_{}_RPC_USER", chain.to_uppercase())) {
cfg.rpc_user = value;
}
Err(_) => {}
}
match env::var(format!("BAL_PUSHER_{}_RPC_PASSWORD", chain.to_uppercase())) {
Ok(value) => {
if let Ok(value) = env::var(format!("BAL_PUSHER_{}_RPC_PASSWORD", chain.to_uppercase())) {
cfg.rpc_pass = value;
}
Err(_) => {}
}
println!(
"{}",
format!("BAL_PUSHER_{}_ZMQ_HASHBLOCK", chain.to_uppercase())
);
match env::var(format!("BAL_PUSHER_{}_ZMQ_HASHBLOCK", chain.to_uppercase())) {
Ok(value) => {
println!("value:{}", value);
if let Ok(value) = env::var(format!("BAL_PUSHER_{}_ZMQ_HASHBLOCK", chain.to_uppercase())) {
cfg.zmq_listener = value;
}
Err(_) => {}
}
cfg.clone()
}
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
struct ConnectionMonitor {
last_message_time: Instant,
timeout: Duration,
consecutive_timeouts: u32,
max_consecutive_timeouts: u32,
}
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();
}
}
enum ConnectionStatus {
Healthy,
Warning(Duration),
Lost(Duration),
}
#[tokio::main]
async fn main() -> std::io::Result<()> {
env_logger::init();
let mut cfg = MyConfig::default();
let dbfile = env::var("BAL_PUSHER_DB_FILE").unwrap();
parse_env(&mut cfg);
let mut args = std::env::args();
let _exe_name = args.next().unwrap();
@@ -755,31 +611,48 @@ async fn main() -> std::io::Result<()> {
}
match socket.set_subscribe(b"") {
Ok(_) => {}
Ok(_) => {
info!("ZMQ subscribed to all topics on {}", zmq_address);
}
Err(e) => {
error!("ZMQ subscribe failed: {}, exiting", e);
return Ok(());
}
}
let _ = main_result(&cfg, network_params).await;
if let Err(e) = main_result(&cfg, network_params).await {
error!("main_result failed on startup: {}", e);
}
info!("waiting new blocks..");
let mut last_seq: Vec<u8> = [0; 4].to_vec();
let mut counter = 0;
let max = 100;
socket.set_rcvtimeo(5000).unwrap(); // 5 seconds timeout
let mut consecutive_timeouts: u32 = 0;
loop {
let message = match socket.recv_multipart(0) {
Ok(m) => m,
Err(e) => {
warn!("ZMQ recv timeout or error: {}, retrying...", e);
consecutive_timeouts += 1;
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;
}
};
if consecutive_timeouts > 0 {
info!(
"ZMQ connection restored after {} consecutive timeouts",
consecutive_timeouts
);
}
consecutive_timeouts = 0;
let topic = message[0].clone();
let body = message[1].clone();
let seq = message[2].clone();
last_seq = seq;
debug!(
"ZMQ:GET TOPIC: {}",
String::from_utf8(topic.clone()).expect("invalid topic")
@@ -787,21 +660,13 @@ async fn main() -> std::io::Result<()> {
trace!("ZMQ:GET BODY: {}", hex::encode(&body));
if topic == b"hashblock" {
info!("NEW BLOCK: {}", hex::encode(&body));
let _ = main_result(&cfg, network_params).await;
if let Err(e) = main_result(&cfg, network_params).await {
error!("main_result failed on new block: {}", e);
}
}
thread::sleep(Duration::from_millis(100)); // Sleep for 100ms
}
}
fn seq_to_str(seq: &Vec<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 {

View File

@@ -7,7 +7,6 @@ use chrono::Utc;
use hex_conservative::FromHex;
use log::{debug, error, info, trace};
use serde::{Deserialize, Serialize};
use serde_json;
use sqlite::State;
use sqlite::{Connection, Value};
use std::collections::{HashMap, HashSet};
@@ -119,6 +118,7 @@ pub struct StatsResponse {
}
#[derive(Debug, Clone)]
#[expect(dead_code)]
struct ActixConfig {
max_body_size: usize,
timeout_secs: u64,
@@ -208,7 +208,7 @@ async fn echo_pub_key(data: web::Data<AppState>) -> impl Responder {
"Failed to read public key file {}: {}",
data.cfg.pub_key_path, e
);
HttpResponse::InternalServerError().body("Failed to read public key file")
HttpResponse::InternalServerError().body("error")
}
}
}
@@ -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>,
@@ -224,26 +256,15 @@ async fn echo_info(
) -> impl Responder {
let param = path.into_inner();
if !NETWORKS.contains(&param.as_str()) {
return HttpResponse::NotFound().body("Unknown network");
return HttpResponse::NotFound().body("error");
}
info!("echo info!!!{}", param);
let netconfig = data.cfg.get_net_config(&param);
if !netconfig.enabled {
debug!("network disabled {}", param);
return HttpResponse::BadRequest().body("network disabled");
return HttpResponse::BadRequest().body("error");
}
let remote_addr = req
.headers()
.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();
@@ -257,7 +278,7 @@ async fn echo_info(
Ok(g) => g,
Err(_p) => {
error!("DB mutex poisoned in echo_info (lookup phase)");
return HttpResponse::InternalServerError().body("DB mutex poisoned");
return HttpResponse::InternalServerError().body("error");
}
};
match get_last_used_address_by_ip(
@@ -275,10 +296,7 @@ async fn echo_info(
version: VERSION.to_string(),
});
}
None => {
let next = get_next_address_index(&db, &netconfig.name, &netconfig.address);
next
}
None => get_next_address_index(&db, &netconfig.name, &netconfig.address),
}
}; // lock released
@@ -288,8 +306,7 @@ async fn echo_info(
Ok(address) => address,
Err(e) => {
error!("Failed to derive address from xpub: {}", e);
return HttpResponse::BadRequest()
.body(format!("Failed to derive address: {}", e));
return HttpResponse::BadRequest().body("error");
}
};
@@ -299,7 +316,7 @@ async fn echo_info(
Ok(g) => g,
Err(_p) => {
error!("DB mutex poisoned in echo_info (save phase)");
return HttpResponse::InternalServerError().body("DB mutex poisoned");
return HttpResponse::InternalServerError().body("error");
}
};
save_new_address(&db, next_idx.0, &derived.0, &derived.1, &remote_addr);
@@ -322,30 +339,30 @@ async fn echo_info(
debug!("echo info reply: {}", json_data);
HttpResponse::Ok().json(info)
}
Err(err) => HttpResponse::InternalServerError().body(format!("error:{}", err)),
Err(_err) => HttpResponse::InternalServerError().body("error"),
}
}
async fn echo_stats(path: web::Path<String>, data: web::Data<AppState>) -> impl Responder {
let param = path.into_inner();
if !NETWORKS.contains(&param.as_str()) {
return HttpResponse::NotFound().body("Unknown network");
return HttpResponse::NotFound().body("error");
}
info!("echo stats!!! {}", data.cfg.expose_stats);
let netconfig = data.cfg.get_net_config(&param);
if !netconfig.enabled {
debug!("network disabled {}", param);
return HttpResponse::BadRequest().body("network disabled");
return HttpResponse::BadRequest().body("error");
}
if !data.cfg.expose_stats {
return HttpResponse::Forbidden().body("Stats not exposed");
return HttpResponse::Forbidden().body("error");
}
let mut stats: Vec<StatsResponse> = vec![];
let db = match data.db.lock() {
Ok(g) => g,
Err(_p) => {
error!("DB mutex poisoned in echo_stats");
return HttpResponse::InternalServerError().body("DB mutex poisoned");
return HttpResponse::InternalServerError().body("error");
}
};
let mut stmt = match db.prepare(
@@ -354,12 +371,12 @@ async fn echo_stats(path: web::Path<String>, data: web::Data<AppState>) -> impl
Ok(s) => s,
Err(e) => {
error!("Failed to prepare stats query: {}", e);
return HttpResponse::InternalServerError().body("Database error");
return HttpResponse::InternalServerError().body("error");
}
};
if let Err(e) = stmt.bind((1, Value::String(netconfig.name.clone()))) {
error!("Failed to bind chain in stats query: {}", e);
return HttpResponse::InternalServerError().body("Database error");
return HttpResponse::InternalServerError().body("error");
}
while let Ok(State::Row) = stmt.next() {
let report_date = stmt.read("report_date").unwrap_or("0".to_string());
@@ -417,47 +434,42 @@ 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);
debug!("echo stats reply for chain: {}", netconfig.name);
HttpResponse::Ok().json(stats)
}
Err(err) => HttpResponse::InternalServerError().body(format!("error:{}", err)),
}
}
async fn echo_search(body: Bytes, data: web::Data<AppState>) -> impl Responder {
info!("echo search!!!");
let strbody = match std::str::from_utf8(&body) {
Ok(s) => s,
Err(_) => {
return HttpResponse::BadRequest().body("Invalid UTF-8 body");
return HttpResponse::BadRequest().body("error");
}
};
info!("{}", strbody);
if strbody.is_empty() || strbody.len() != 64 || !strbody.chars().all(|c| c.is_ascii_hexdigit())
{
return HttpResponse::BadRequest().body("Invalid txid");
return HttpResponse::BadRequest().body("error");
}
let db = match data.db.lock() {
Ok(g) => g,
Err(_p) => {
error!("DB mutex poisoned in echo_search");
return HttpResponse::InternalServerError().body("DB mutex poisoned");
return HttpResponse::InternalServerError().body("error");
}
};
let mut statement = match db.prepare("SELECT * FROM tbl_tx WHERE txid = ? LIMIT 1") {
Ok(s) => s,
Err(e) => {
error!("Failed to prepare statement: {}", e);
return HttpResponse::InternalServerError().body("Database error");
return HttpResponse::InternalServerError().body("error");
}
};
if let Err(e) = statement.bind((1, strbody)) {
error!("Failed to bind parameter: {}", e);
return HttpResponse::InternalServerError().body("Database error");
return HttpResponse::InternalServerError().body("error");
}
if let Ok(State::Row) = statement.next() {
@@ -503,11 +515,14 @@ async fn echo_search(body: Bytes, data: web::Data<AppState>) -> impl Responder {
}
}
match serde_json::to_string(&response_data) {
Ok(json_data) => HttpResponse::Ok().json(json_data),
Err(_) => HttpResponse::BadRequest().body("Bad data received"),
Ok(json_data) => {
debug!("echo search reply: {}", json_data);
HttpResponse::Ok().json(&response_data)
}
Err(_) => HttpResponse::BadRequest().body("error"),
}
} else {
HttpResponse::BadRequest().body("Bad data received")
HttpResponse::BadRequest().body("error")
}
}
@@ -524,15 +539,14 @@ struct ParsedTx {
}
/// Parse all transactions from the request body **without** needing the DB lock.
/// Returns `Ok(parsed_txs)` if at least one tx was valid, or `Err(HttpResponse)` for early failure.
/// Skips transactions that don't have a valid willexecutor output.
fn parse_request_transactions(
strbody: &str,
_req_time: i64,
netconfig: &NetConfig,
known_addresses: &HashSet<String>,
) -> Result<Vec<(ParsedTx, String, u64)>, HttpResponse> {
) -> Vec<(ParsedTx, String, u64)> {
let mut result: Vec<(ParsedTx, String, u64)> = Vec::new();
let mut union_tx = true;
for line in strbody.split('\n') {
if line.is_empty() {
@@ -613,13 +627,8 @@ fn parse_request_transactions(
}
if !found {
error!("willexecutor output not found for tx {}", txid);
return Err(HttpResponse::BadRequest().body("Bad data received"));
}
if !union_tx {
// This is only used for SQL building later; we track it in the caller
} else {
union_tx = false;
trace!("willexecutor output not found for tx {}, skipping", txid);
continue;
}
result.push((
ParsedTx {
@@ -636,7 +645,7 @@ fn parse_request_transactions(
));
}
Ok(result)
result
}
async fn echo_push(
@@ -648,24 +657,24 @@ async fn echo_push(
let strbody = match std::str::from_utf8(&body) {
Ok(s) => s,
Err(_) => {
return HttpResponse::BadRequest().body("Invalid UTF-8 body");
return HttpResponse::BadRequest().body("error");
}
};
let param = path.into_inner();
if !NETWORKS.contains(&param.as_str()) {
return HttpResponse::NotFound().body("Unknown network");
return HttpResponse::NotFound().body("error");
}
let netconfig = data.cfg.get_net_config(&param);
if !netconfig.enabled {
trace!("network not enabled {}", &netconfig.name);
return HttpResponse::BadRequest().body("Network not enabled");
return HttpResponse::BadRequest().body("error");
}
let req_time = match Utc::now().timestamp_nanos_opt() {
Some(t) => t,
None => {
error!("Invalid timestamp");
return HttpResponse::BadRequest().body("Invalid timestamp");
return HttpResponse::BadRequest().body("error");
}
};
@@ -675,7 +684,7 @@ async fn echo_push(
Ok(g) => g,
Err(_p) => {
error!("DB mutex poisoned acquiring addresses in echo_push");
return HttpResponse::InternalServerError().body("DB mutex poisoned");
return HttpResponse::InternalServerError().body("error");
}
};
if netconfig.xpub {
@@ -683,7 +692,7 @@ async fn echo_push(
Ok(addrs) => addrs,
Err(e) => {
error!("Failed to load addresses from xpub: {}", e);
return HttpResponse::InternalServerError().body("Database error");
return HttpResponse::InternalServerError().body("error");
}
}
} else {
@@ -692,12 +701,9 @@ async fn echo_push(
}; // lock released here
// Parse all transactions (CPU-bound, no DB needed)
let parsed = match parse_request_transactions(strbody, req_time, netconfig, &known_addresses) {
Ok(v) => v,
Err(resp) => return resp,
};
let parsed = parse_request_transactions(strbody, req_time, netconfig, &known_addresses);
if parsed.is_empty() {
return HttpResponse::Ok().body("thx");
return HttpResponse::BadRequest().body("error");
}
let all_txids: Vec<String> = parsed.iter().map(|(p, _, _)| p.txid.clone()).collect();
@@ -708,14 +714,14 @@ async fn echo_push(
Ok(g) => g,
Err(_p) => {
error!("DB mutex poisoned in echo_push duplicate check");
return HttpResponse::InternalServerError().body("DB mutex poisoned");
return HttpResponse::InternalServerError().body("error");
}
};
match check_duplicate_txids(&db, &all_txids) {
Ok(dups) => dups,
Err(e) => {
error!("Duplicate check failed: {}", e);
return HttpResponse::InternalServerError().body("Database error");
return HttpResponse::InternalServerError().body("error");
}
}
}; // lock released here
@@ -731,7 +737,7 @@ async fn echo_push(
Ok(g) => g,
Err(_p) => {
error!("DB mutex poisoned in echo_push insert phase");
return HttpResponse::InternalServerError().body("DB mutex poisoned");
return HttpResponse::InternalServerError().body("error");
}
};
@@ -818,7 +824,7 @@ async fn echo_push(
if let Err(err) = execute_insert(&db, sqltxs, ptx, sqlinps, pinps, sqlouts, pouts) {
error!("execute_insert failed: {}", err);
return HttpResponse::BadRequest().body("Bad data received");
return HttpResponse::BadRequest().body("error");
}
} // lock released
@@ -900,7 +906,7 @@ async fn main() -> std::io::Result<()> {
let db = match open_db(&cfg.db_file) {
Ok(c) => c,
Err(e) => {
return Err(std::io::Error::new(std::io::ErrorKind::Other, e));
return Err(std::io::Error::other(e));
}
};
@@ -915,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;

View File

@@ -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>{
@@ -178,10 +178,11 @@ pub fn create_database(db: &Connection) {
}
}
*/
pub fn insert_xpub(db: &Connection, network: &String, xpub: &String) {
if xpub != "" {
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(?, ?);") {
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);

View File

@@ -11,6 +11,7 @@ use sha2::{Digest, Sha256};
use std::str::FromStr;
// Mainnet (BIP44/BIP49/BIP84)
#[allow(dead_code)]
enum BS58Prefix {
Xpub,
Ypub,
@@ -102,44 +103,29 @@ fn calc_checksum(desc: &str) -> Result<String, String> {
Ok(checksum)
}
pub fn get_bitcoincore_descriptor(xpub: &String) -> String {
pub fn get_bitcoincore_descriptor(xpub: &str) -> String {
let fingerprint = match calculate_fingerprint(xpub) {
Ok(f) => f,
Err(_) => return String::new(), // Invalid xpub, return empty descriptor
};
let mut bip = 84;
let cpub = xpub.to_string();
match &xpub[0..4] {
"vpub" => {
bip = 84;
}
"zpub" => {
bip = 84;
}
&_ => {
bip = 84;
}
};
let xpub_converted = match convert_xpub(xpub) {
Ok(c) => c,
Err(_) => return String::new(), // Invalid xpub, return empty descriptor
};
let descriptor = format!("wpkh([{}/84h/0h/0h]{}/0/*)", fingerprint, xpub_converted);
let descriptor = match calc_checksum(&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()
String::new()
}
};
descriptor
//format!("{}#{}",descriptor,checksum)
}
fn convert_xpub(xpub: &String) -> Result<String, String> {
}
fn convert_xpub(xpub: &str) -> Result<String, String> {
if xpub.len() >= 4 && (&xpub[0..4] == "xpub" || &xpub[0..4] == "ypub" || &xpub[0..4] == "zpub")
{
convert_to(xpub, BS58Prefix::Xpub)
@@ -165,7 +151,7 @@ fn base58check_decode(s: &str) -> Result<Vec<u8>, String> {
return Err("Data troppo corta".to_string());
}
let (payload, checksum) = data.split_at(data.len() - 4);
let hash = Sha256::digest(&Sha256::digest(payload));
let hash = Sha256::digest(Sha256::digest(payload));
if hash[0..4] != checksum[..] {
return Err("Checksum invalido".to_string());
}
@@ -173,7 +159,7 @@ fn base58check_decode(s: &str) -> Result<Vec<u8>, String> {
}
fn base58check_encode(data: &[u8]) -> String {
let checksum = &Sha256::digest(&Sha256::digest(data))[0..4];
let checksum = &Sha256::digest(Sha256::digest(data))[0..4];
let full = [data, checksum].concat();
bs58::encode(full).into_string()
}
@@ -205,7 +191,7 @@ pub fn new_address_from_xpub(
) -> 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 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;

View File

@@ -1,7 +1,6 @@
use bal_server::db::open_db;
use sqlite::State;
use std::fs;
use std::path::Path;
#[test]
fn test_open_db_blocks_traversal() {
@@ -74,7 +73,7 @@ fn test_open_db_rejects_symlink() {
let _ = fs::remove_file(real);
let _ = fs::remove_file(link);
fs::File::create(real).unwrap();
fs::soft_link(real, link).unwrap();
std::os::unix::fs::symlink(real, link).unwrap();
let res = open_db(link);
assert!(res.is_err(), "Symlink DB path should be rejected");

View File

@@ -36,11 +36,8 @@ fn test_db_null_unwrap_or() {
let mut found_value = None;
let _ = db.iterate("SELECT * FROM test_stats;", |pairs| {
let row: HashMap<_, _> = pairs
.into_iter()
.map(|(k, v)| (k.to_string(), v.map(|s| s)))
.collect();
let totals = row["totals"].clone().unwrap_or("0").to_string();
let row: HashMap<_, _> = pairs.iter().map(|(k, v)| (k.to_string(), *v)).collect();
let totals = row["totals"].unwrap_or("0").to_string();
found_value = Some(totals);
true
});

View File

@@ -1,5 +1,4 @@
use std::fs;
use std::path::Path;
#[test]
fn test_gitignore_protection_env() {
@@ -23,30 +22,20 @@ fn test_gitignore_protection_env() {
];
for pattern in required_patterns {
let has_exact = gitignore.contains(&pattern);
let has_wildcard = gitignore.contains(&format!("*.env.local"))
|| gitignore.contains(&format!(".env.local"));
let has_wildcard = gitignore.contains("*.env.local") || gitignore.contains(".env.local");
let has_env = gitignore.contains("*.env") || gitignore.contains(".env");
// For .env.local, either .env.local or *.env.local is acceptable
let is_env_local = pattern == "*.env.local" || pattern == ".env.local";
if is_env_local {
assert!(
has_wildcard,
".gitignore must contain pattern '*.env.local' or '.env.local' to protect secrets",
);
} else if pattern == ".env.production" || pattern == ".env.secret" {
assert!(
gitignore.contains(pattern),
".gitignore must contain pattern '{}' to protect secrets",
pattern
);
} else if pattern == "*.env" {
assert!(
has_env,
".gitignore must contain pattern '*.env' or '.env' to protect secrets",
);
} else if pattern == ".env" {
} else if pattern == ".env.production"
|| pattern == ".env.secret"
|| pattern == "*.env"
|| pattern == ".env"
{
assert!(
has_env,
".gitignore must contain pattern '*.env' or '.env' to protect secrets",
@@ -65,7 +54,6 @@ fn test_gitignore_protection_env() {
#[test]
fn test_no_private_key_in_git() {
// Check that .gitignore includes private_key.pem
let gitignore = match fs::read_to_string(".gitignore") {
Ok(c) => c,
Err(e) => {
@@ -88,20 +76,17 @@ fn test_no_private_key_in_git() {
".gitignore must block chiave_privata.key"
);
// Check that no private key files are tracked by git
let output = std::process::Command::new("git")
.args(&["ls-files", "*.pem", "*.key"])
.args(["ls-files", "*.pem", "*.key"])
.output()
.expect("Failed to run git ls-files");
let tracked_keys = String::from_utf8(output.stdout).unwrap();
let tracked_keys: Vec<&str> = tracked_keys.lines().collect();
// Only non-empty entries and only public_key.pem should be tracked
for tracked in tracked_keys.iter().filter(|s| !s.is_empty()) {
if !tracked.contains("public_key.pem") {
assert!(
false,
panic!(
"Private key file is tracked by git: {}. Remove it with git rm --cached",
tracked
);
@@ -113,38 +98,34 @@ fn test_no_private_key_in_git() {
#[test]
fn test_no_token_in_source_files() {
// Scan source files for hardcoded tokens
let mut found_issues = Vec::new();
// Scan .sh files for hardcoded 40-char hex strings
for entry in fs::read_dir(".").unwrap().filter_map(|e| e.ok()) {
let path = entry.path();
if !path.is_file() {
continue;
}
if let Some(ext) = path.extension() {
if ext == "sh" {
if let Some(ext) = path.extension()
&& ext == "sh"
{
let content = fs::read_to_string(&path).unwrap();
for (line_num, line) in content.lines().enumerate() {
// Skip comments and example/template files
if line.trim().starts_with("#")
if line.trim().starts_with('#')
|| line.to_lowercase().contains("example")
|| line.to_lowercase().contains("template")
{
continue;
}
// Check for 40-64 hex chars that could be API tokens (not in .env.example comments)
if line.trim().len() >= 40 {
let hex_chars = line
let hex_chars: Vec<_> = line
.trim()
.chars()
.filter(|c| c.is_ascii_hexdigit())
.collect::<Vec<_>>();
if hex_chars.len() >= 40 && hex_chars.len() <= 64 {
// Check if it looks like it's part of a TOKEN assignment
if line.to_lowercase().contains("token")
.collect();
if (40..=64).contains(&hex_chars.len())
&& (line.to_lowercase().contains("token")
|| line.to_lowercase().contains("api")
|| line.to_lowercase().contains("secret")
|| line.to_lowercase().contains("secret"))
{
found_issues.push(format!(
"Potential hardcoded token in {}: line {}: {}",
@@ -157,16 +138,13 @@ fn test_no_token_in_source_files() {
}
}
}
}
}
if !found_issues.is_empty() {
println!("FAIL: Found potential hardcoded tokens:");
for issue in &found_issues {
println!(" {}", issue);
}
assert!(
false,
panic!(
"Found potential hardcoded tokens in shell scripts: {:?}",
found_issues
);