14 Commits

Author SHA1 Message Date
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
7fe5fd3139 fix(pusher): log send_stats_report errors instead of discarding
main_result called send_stats_report/calculate_stats with 'let _ = ...',
silently dropping any failure. A broken welist route (e.g. unreachable
IPv4 path) is then invisible in the logs and can go unnoticed for a long
time. Log failures with warn! so connectivity problems are diagnosable.
2026-07-19 14:09:41 +02:00
b46f85f436 feat(pusher): optional IPv6 preference for welist reports (BAL_PUSHER_PREFER_IPV6)
The welist host publishes both A and AAAA records. On networks where the
IPv4 route is broken (connection stalls after the TCP handshake) while
IPv6 works, the default connector may pick the broken family and the
report request hangs.

When BAL_PUSHER_PREFER_IPV6 is truthy, the pusher now resolves the welist
host itself and pins the reqwest client to its first IPv6 address; the
original hostname is still used for the Host header and TLS SNI. When the
variable is unset (default) or no AAAA record exists, behavior is
completely unchanged.

Includes unit tests for the URL host/port parsing and documentation in
docs/07_deployment_and_ops.md.
2026-07-19 14:09:01 +02: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
12 changed files with 510 additions and 597 deletions

4
.gitignore vendored
View File

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

275
Cargo.lock generated
View File

@@ -312,7 +312,7 @@ checksum = "ace50bade8e6234aa140d9a2f552bbee1db4d353f69b8217bc503490fc1a9f26"
[[package]] [[package]]
name = "bal_server" name = "bal_server"
version = "0.3.0" version = "0.3.1"
dependencies = [ dependencies = [
"actix-governor", "actix-governor",
"actix-rt", "actix-rt",
@@ -325,12 +325,11 @@ dependencies = [
"byteorder", "byteorder",
"bytes", "bytes",
"chrono", "chrono",
"confy", "ed25519-dalek",
"env_logger", "env_logger",
"hex", "hex",
"hex-conservative 0.1.1", "hex-conservative 0.1.1",
"log", "log",
"openssl",
"regex", "regex",
"reqwest", "reqwest",
"serde", "serde",
@@ -364,6 +363,12 @@ version = "0.22.1"
source = "registry+https://github.com/rust-lang/crates.io-index" source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "72b3254f16251a8381aa12e40e3c4d2f0199f8c6508fbecb9d91f575e0fbb8c6" checksum = "72b3254f16251a8381aa12e40e3c4d2f0199f8c6508fbecb9d91f575e0fbb8c6"
[[package]]
name = "base64ct"
version = "1.8.3"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "2af50177e190e07a26ab74f8b1efbfe2ef87da2116221318cb1c2e82baf7de06"
[[package]] [[package]]
name = "bech32" name = "bech32"
version = "0.11.0" version = "0.11.0"
@@ -592,16 +597,10 @@ source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "d3fd119d74b830634cea2a0f58bbd0d54540518a14397557951e79340abc28c0" checksum = "d3fd119d74b830634cea2a0f58bbd0d54540518a14397557951e79340abc28c0"
[[package]] [[package]]
name = "confy" name = "const-oid"
version = "0.6.1" version = "0.9.6"
source = "registry+https://github.com/rust-lang/crates.io-index" source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "45b1f4c00870f07dc34adcac82bb6a72cc5aabca8536ba1797e01df51d2ce9a0" checksum = "c2459377285ad874054d797f3ccebf984978aa39129f6eafde5cdc8315b612f8"
dependencies = [
"directories",
"serde",
"thiserror",
"toml",
]
[[package]] [[package]]
name = "const-oid" name = "const-oid"
@@ -747,6 +746,33 @@ dependencies = [
"hybrid-array", "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]] [[package]]
name = "dashmap" name = "dashmap"
version = "5.5.3" version = "5.5.3"
@@ -760,6 +786,17 @@ dependencies = [
"parking_lot_core", "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]] [[package]]
name = "deranged" name = "deranged"
version = "0.5.8" version = "0.5.8"
@@ -809,7 +846,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "f1dd6dbb5841937940781866fa1281a1ff7bd3bf827091440879f9994983d5c2" checksum = "f1dd6dbb5841937940781866fa1281a1ff7bd3bf827091440879f9994983d5c2"
dependencies = [ dependencies = [
"block-buffer 0.12.1", "block-buffer 0.12.1",
"const-oid", "const-oid 0.10.2",
"crypto-common 0.2.2", "crypto-common 0.2.2",
] ]
@@ -824,27 +861,6 @@ dependencies = [
"walkdir", "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]] [[package]]
name = "displaydoc" name = "displaydoc"
version = "0.2.5" version = "0.2.5"
@@ -856,6 +872,30 @@ dependencies = [
"syn", "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]] [[package]]
name = "either" name = "either"
version = "1.15.0" version = "1.15.0"
@@ -916,6 +956,12 @@ version = "2.3.0"
source = "registry+https://github.com/rust-lang/crates.io-index" source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "37909eebbb50d72f9059c3b6d82c0463f2ff062c9e95845c43a6c9c0355411be" checksum = "37909eebbb50d72f9059c3b6d82c0463f2ff062c9e95845c43a6c9c0355411be"
[[package]]
name = "fiat-crypto"
version = "0.2.9"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "28dea519a9695b9977216879a3ebfddf92f1c08c05d984f8996aecd6ecdc811d"
[[package]] [[package]]
name = "find-msvc-tools" name = "find-msvc-tools"
version = "0.1.4" version = "0.1.4"
@@ -1578,16 +1624,6 @@ version = "0.2.186"
source = "registry+https://github.com/rust-lang/crates.io-index" source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "68ab91017fe16c622486840e4c83c9a37afeff978bd239b5293d61ece587de66" 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]] [[package]]
name = "linux-raw-sys" name = "linux-raw-sys"
version = "0.11.0" version = "0.11.0"
@@ -1759,15 +1795,6 @@ version = "0.1.6"
source = "registry+https://github.com/rust-lang/crates.io-index" source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "d05e27ee213611ffe7d6348b942e8f942b37114c00cc03cec254295a4a17852e" 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]] [[package]]
name = "openssl-sys" name = "openssl-sys"
version = "0.9.110" version = "0.9.110"
@@ -1776,17 +1803,10 @@ checksum = "0a9f0075ba3c21b09f8e8b2026584b1d18d49388648f2fbbf3c97ea8deced8e2"
dependencies = [ dependencies = [
"cc", "cc",
"libc", "libc",
"openssl-src",
"pkg-config", "pkg-config",
"vcpkg", "vcpkg",
] ]
[[package]]
name = "option-ext"
version = "0.2.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "04744f49eae99ab78e0d5c0b603ab218f515ea8cfe5a456d7629ad883a3b6e7d"
[[package]] [[package]]
name = "parking_lot" name = "parking_lot"
version = "0.12.5" version = "0.12.5"
@@ -1810,6 +1830,15 @@ dependencies = [
"windows-link 0.2.1", "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]] [[package]]
name = "percent-encoding" name = "percent-encoding"
version = "2.3.2" version = "2.3.2"
@@ -1828,6 +1857,16 @@ version = "0.1.0"
source = "registry+https://github.com/rust-lang/crates.io-index" source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "8b870d8c151b6f2fb93e84a13146138f05d02ed11c7e7c54f8826aaaf7c9f184" 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]] [[package]]
name = "pkg-config" name = "pkg-config"
version = "0.3.30" version = "0.3.30"
@@ -1994,17 +2033,6 @@ dependencies = [
"bitflags 2.9.0", "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]] [[package]]
name = "regex" name = "regex"
version = "1.10.4" version = "1.10.4"
@@ -2326,6 +2354,15 @@ dependencies = [
"libc", "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]] [[package]]
name = "simd-adler32" name = "simd-adler32"
version = "0.3.10" version = "0.3.10"
@@ -2373,6 +2410,16 @@ dependencies = [
"lock_api", "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]] [[package]]
name = "sqlite" name = "sqlite"
version = "0.34.0" version = "0.34.0"
@@ -2499,26 +2546,6 @@ dependencies = [
"windows-sys 0.61.2", "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]] [[package]]
name = "time" name = "time"
version = "0.3.44" version = "0.3.44"
@@ -3032,15 +3059,6 @@ dependencies = [
"windows-link 0.1.1", "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]] [[package]]
name = "windows-sys" name = "windows-sys"
version = "0.52.0" version = "0.52.0"
@@ -3068,21 +3086,6 @@ dependencies = [
"windows-link 0.2.1", "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]] [[package]]
name = "windows-targets" name = "windows-targets"
version = "0.52.5" version = "0.52.5"
@@ -3116,12 +3119,6 @@ dependencies = [
"windows_x86_64_msvc 0.53.1", "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]] [[package]]
name = "windows_aarch64_gnullvm" name = "windows_aarch64_gnullvm"
version = "0.52.5" version = "0.52.5"
@@ -3134,12 +3131,6 @@ version = "0.53.1"
source = "registry+https://github.com/rust-lang/crates.io-index" source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "a9d8416fa8b42f5c947f8482c43e7d89e73a173cead56d044f6a56104a6d1b53" checksum = "a9d8416fa8b42f5c947f8482c43e7d89e73a173cead56d044f6a56104a6d1b53"
[[package]]
name = "windows_aarch64_msvc"
version = "0.48.5"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "dc35310971f3b2dbbf3f0690a219f40e2d9afcf64f9ab7cc1be722937c26b4bc"
[[package]] [[package]]
name = "windows_aarch64_msvc" name = "windows_aarch64_msvc"
version = "0.52.5" version = "0.52.5"
@@ -3152,12 +3143,6 @@ version = "0.53.1"
source = "registry+https://github.com/rust-lang/crates.io-index" source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "b9d782e804c2f632e395708e99a94275910eb9100b2114651e04744e9b125006" checksum = "b9d782e804c2f632e395708e99a94275910eb9100b2114651e04744e9b125006"
[[package]]
name = "windows_i686_gnu"
version = "0.48.5"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "a75915e7def60c94dcef72200b9a8e58e5091744960da64ec734a6c6e9b3743e"
[[package]] [[package]]
name = "windows_i686_gnu" name = "windows_i686_gnu"
version = "0.52.5" version = "0.52.5"
@@ -3182,12 +3167,6 @@ version = "0.53.1"
source = "registry+https://github.com/rust-lang/crates.io-index" source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "fa7359d10048f68ab8b09fa71c3daccfb0e9b559aed648a8f95469c27057180c" checksum = "fa7359d10048f68ab8b09fa71c3daccfb0e9b559aed648a8f95469c27057180c"
[[package]]
name = "windows_i686_msvc"
version = "0.48.5"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "8f55c233f70c4b27f66c523580f78f1004e8b5a8b659e05a4eb49d4166cca406"
[[package]] [[package]]
name = "windows_i686_msvc" name = "windows_i686_msvc"
version = "0.52.5" version = "0.52.5"
@@ -3200,12 +3179,6 @@ version = "0.53.1"
source = "registry+https://github.com/rust-lang/crates.io-index" source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "1e7ac75179f18232fe9c285163565a57ef8d3c89254a30685b57d83a38d326c2" 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]] [[package]]
name = "windows_x86_64_gnu" name = "windows_x86_64_gnu"
version = "0.52.5" version = "0.52.5"
@@ -3218,12 +3191,6 @@ version = "0.53.1"
source = "registry+https://github.com/rust-lang/crates.io-index" source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "9c3842cdd74a865a8066ab39c8a7a473c0778a3f29370b5fd6b4b9aa7df4a499" 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]] [[package]]
name = "windows_x86_64_gnullvm" name = "windows_x86_64_gnullvm"
version = "0.52.5" version = "0.52.5"
@@ -3236,12 +3203,6 @@ version = "0.53.1"
source = "registry+https://github.com/rust-lang/crates.io-index" source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "0ffa179e2d07eee8ad8f57493436566c7cc30ac536a3379fdf008f47f6bb7ae1" 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]] [[package]]
name = "windows_x86_64_msvc" name = "windows_x86_64_msvc"
version = "0.52.5" version = "0.52.5"

View File

@@ -1,42 +1,58 @@
[package] [package]
name = "bal_server" name = "bal_server"
version = "0.3.0" version = "0.3.1"
edition = "2024" edition = "2024"
# See more keys and their definitions at https://doc.rust-lang.org/cargo/reference/manifest.html # 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] [dependencies]
base64 = { version = "0.22.1" }
bs58 = { version = "0.4.0" } bs58 = { version = "0.4.0" }
bytes = { version = "1.2" } bytes = { version = "1.2" }
bitcoin = { version = "0.32.5" } bitcoin = { version = "0.32.5" }
bitcoincore-rpc = { version = "0.19.0" } bitcoincore-rpc = { version = "0.19.0" }
bitcoincore-rpc-json = { 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" } env_logger = { version = "0.11.5" }
hex = { version = "0.4.3" } 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" } 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 = { version = "1.0.152", features = ["derive"] }
serde_json = { version = "1.0.116" } serde_json = { version = "1.0.116" }
sqlite = { version = "0.34.0" } sqlite = { version = "0.34.0" }
regex = { version = "1.10.4" } 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"] } tokio = { version = "1", features = ["rt", "net","macros","rt-multi-thread"] }
url = { version = "2" } 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]] [[bin]]
name = "bal-server" name = "bal-server"
path = "src/bin/bal-server.rs" path = "src/bin/bal-server.rs"
required-features = ["server"]
[[bin]] [[bin]]
name = "bal-pusher" name = "bal-pusher"
path = "src/bin/bal-pusher.rs" 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

@@ -45,6 +45,7 @@ WELIST_URL=https://welist.example.com/api/stats
- `BAL_SSL_KEY_PATH`: The path to the Ed25519 private key (`private_key.pem`) used to sign the statistics payload before sending it to the `welist` server. This is a critical secret. - `BAL_SSL_KEY_PATH`: The path to the Ed25519 private key (`private_key.pem`) used to sign the statistics payload before sending it to the `welist` server. This is a critical secret.
- `SEND_STATS`: A boolean flag to enable the reporting of statistics to the remote `welist` server. - `SEND_STATS`: A boolean flag to enable the reporting of statistics to the remote `welist` server.
- `WELIST_URL`: The URL to which the statistics are sent. If `SEND_STATS` is `true`, this URL must be reachable. If the server is unreachable, the pusher will log an error but might not crash (see `08_security_audit.md` for DoS analysis). - `WELIST_URL`: The URL to which the statistics are sent. If `SEND_STATS` is `true`, this URL must be reachable. If the server is unreachable, the pusher will log an error but might not crash (see `08_security_audit.md` for DoS analysis).
- `BAL_PUSHER_PREFER_IPV6`: Optional boolean flag (default `false`). When set to `true`, the pusher resolves the `welist` host itself and pins the HTTP connection to its first IPv6 (AAAA) address, still using the hostname for the `Host` header and TLS SNI. This works around networks where the IPv4 route to the `welist` host is broken while IPv6 works — the default connector may otherwise pick the unreachable family and the request would stall. Leave unset unless you hit this specific connectivity problem.
--- ---

View File

@@ -5,8 +5,7 @@ use bitcoin::Network;
use bitcoincore_rpc::{Auth, Client, Error, RpcApi, bitcoin}; use bitcoincore_rpc::{Auth, Client, Error, RpcApi, bitcoin};
use bitcoincore_rpc_json::GetBlockchainInfoResult; use bitcoincore_rpc_json::GetBlockchainInfoResult;
use byteorder::{LittleEndian, ReadBytesExt}; use ed25519_dalek::{Signer as _, SigningKey, pkcs8::DecodePrivateKey};
use hex;
use log::{debug, error, info, trace, warn}; use log::{debug, error, info, trace, warn};
use serde::Deserialize; use serde::Deserialize;
use serde::Serialize; use serde::Serialize;
@@ -15,24 +14,19 @@ use sqlite::{Connection, Value};
use std::collections::HashMap; use std::collections::HashMap;
use std::env; use std::env;
use std::error::Error as StdError; use std::error::Error as StdError;
use std::io::Cursor;
use std::str; use std::str;
use std::{thread, time::Duration}; use std::{thread, time::Duration};
use zmq::{Context, DEALER, DONTWAIT, Socket}; use zmq::{Context, Socket};
use bal_server::db::open_db; use bal_server::db::open_db;
use bal_server::validation::is_valid_welist_url; use bal_server::validation::is_valid_welist_url;
use base64::{Engine as _, engine::general_purpose}; 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 reqwest::Client as rClient;
use std::fs; use std::net::SocketAddr;
use std::time::Instant; use url::Url;
const LOCKTIME_THRESHOLD: i64 = 5000000; const LOCKTIME_THRESHOLD: i64 = 5000000;
const VERSION: &str = "0.0.2"; const VERSION: &str = env!("CARGO_PKG_VERSION");
#[derive(Debug, Clone, Serialize, Deserialize)] #[derive(Debug, Clone, Serialize, Deserialize)]
struct MyConfig { struct MyConfig {
db_file: String, db_file: String,
@@ -90,7 +84,7 @@ fn get_network_params(cfg: &MyConfig, network: Network) -> &NetworkParams {
fn get_network_params_default(network: Network) -> NetworkParams { fn get_network_params_default(network: Network) -> NetworkParams {
match network { match network {
Network::Testnet => NetworkParams { Network::Testnet => NetworkParams {
host: "http://i27.0.0.1".to_string(), host: "http://127.0.0.1".to_string(),
port: 18332, port: 18332,
dir_path: "testnet3/".to_string(), dir_path: "testnet3/".to_string(),
db_field: "testnet".to_string(), db_field: "testnet".to_string(),
@@ -100,7 +94,7 @@ fn get_network_params_default(network: Network) -> NetworkParams {
zmq_listener: "tcp://127.0.0.1:23332".to_string(), zmq_listener: "tcp://127.0.0.1:23332".to_string(),
}, },
Network::Testnet4 => NetworkParams { Network::Testnet4 => NetworkParams {
host: "http://i27.0.0.1".to_string(), host: "http://127.0.0.1".to_string(),
port: 48332, port: 48332,
dir_path: "testnet4/".to_string(), dir_path: "testnet4/".to_string(),
db_field: "testnet4".to_string(), db_field: "testnet4".to_string(),
@@ -143,7 +137,7 @@ fn get_network_params_default(network: Network) -> NetworkParams {
} }
fn get_cookie_filename(network: &NetworkParams) -> Result<String, Box<dyn StdError>> { 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()) Ok(network.cookie_file.clone())
} else { } else {
match env::var_os("HOME") { match env::var_os("HOME") {
@@ -161,12 +155,12 @@ fn get_cookie_filename(network: &NetworkParams) -> Result<String, Box<dyn StdErr
} }
} }
fn get_client_from_username( fn get_client_from_username(
url: &String, url: &str,
network: &NetworkParams, network: &NetworkParams,
) -> Result<(Client, GetBlockchainInfoResult), Box<dyn StdError>> { ) -> Result<(Client, GetBlockchainInfoResult), Box<dyn StdError>> {
if network.rpc_user != "" { if !network.rpc_user.is_empty() {
match Client::new( match Client::new(
&url[..], url,
Auth::UserPass(network.rpc_user.to_string(), network.rpc_pass.to_string()), Auth::UserPass(network.rpc_user.to_string(), network.rpc_pass.to_string()),
) { ) {
Ok(client) => match client.get_blockchain_info() { Ok(client) => match client.get_blockchain_info() {
@@ -180,57 +174,37 @@ fn get_client_from_username(
} }
} }
fn get_client_from_cookie( fn get_client_from_cookie(
url: &String, url: &str,
network: &NetworkParams, network: &NetworkParams,
) -> Result<(Client, GetBlockchainInfoResult), Box<dyn StdError>> { ) -> Result<(Client, GetBlockchainInfoResult), Box<dyn StdError>> {
match get_cookie_filename(network) { 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(client) => match client.get_blockchain_info() {
Ok(bcinfo) => Ok((client, bcinfo)), Ok(bcinfo) => Ok((client, bcinfo)),
Err(err) => Err(err.into()), Err(err) => Err(err.into()),
}, },
Err(err) => Err(err.into()), Err(err) => Err(err.into()),
}, },
Err(err) => Err(err.into()), Err(err) => Err(err),
} }
} }
fn get_client( fn get_client(
network: &NetworkParams, network: &NetworkParams,
) -> Result<(Client, GetBlockchainInfoResult), Box<dyn StdError>> { ) -> 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}"); debug!("trying to connect to bitcoin daemon:{url}");
match get_client_from_username(&url, network) { match get_client_from_username(&url, network) {
Ok(client) => Ok(client), Ok(client) => Ok(client),
Err(_) => match get_client_from_cookie(&url, &network) { Err(_) => match get_client_from_cookie(&url, network) {
Ok(client) => Ok(client), Ok(client) => Ok(client),
Err(err) => Err(err.into()), Err(err) => Err(err),
}, },
} }
} }
async fn main_result(cfg: &MyConfig, network_params: &NetworkParams) -> Result<(), Error> { 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) { match get_client(network_params) {
Ok((rpc, bcinfo)) => { Ok((rpc, bcinfo)) => {
info!("connected"); 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!("median time: {}", bcinfo.median_time);
//info!("height time: {}",bcinfo.median_time); //info!("height time: {}",bcinfo.median_time);
info!("blocks: {}", bcinfo.blocks); info!("blocks: {}", bcinfo.blocks);
@@ -265,7 +239,7 @@ async fn main_result(cfg: &MyConfig, network_params: &NetworkParams) -> Result<(
let mut invalid_txs: std::collections::HashMap<String, String> = HashMap::new(); let mut invalid_txs: std::collections::HashMap<String, String> = HashMap::new();
for row_result in match query_tx.bind::<&[(_, Value)]>( 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_time", (average_time as i64).into()),
(":bestblock_height", (bcinfo.blocks as i64).into()), (":bestblock_height", (bcinfo.blocks as i64).into()),
(":network", network_params.db_field.clone().into()), (":network", network_params.db_field.clone().into()),
@@ -291,26 +265,10 @@ async fn main_result(cfg: &MyConfig, network_params: &NetworkParams) -> Result<(
info!("to be pushed: {}: {}", txid, locktime); info!("to be pushed: {}: {}", txid, locktime);
match rpc.send_raw_transaction(tx) { match rpc.send_raw_transaction(tx) {
Ok(o) => { 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); info!("tx: {} pusshata PUSHED\n{}", txid, o);
pushed_txs.push(txid.to_string()); pushed_txs.push(txid.to_string());
} }
Err(err) => { 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); warn!("Error: {}\n{}", err, txid);
//store err in invalid_txs //store err in invalid_txs
invalid_txs.insert(txid.to_string(), err.to_string()); invalid_txs.insert(txid.to_string(), err.to_string());
@@ -320,19 +278,44 @@ async fn main_result(cfg: &MyConfig, network_params: &NetworkParams) -> Result<(
for txid in &pushed_txs { for txid in &pushed_txs {
let sql = "UPDATE tbl_tx SET status = 1 WHERE txid = ?"; let sql = "UPDATE tbl_tx SET status = 1 WHERE txid = ?";
let mut stmt = db.prepare(sql).unwrap(); match db.prepare(sql) {
stmt.bind((1, Value::String(txid.clone()))).unwrap(); 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(); let _ = stmt.next();
} }
Err(e) => {
error!("Failed to prepare status update: {}", e);
}
}
}
for (txid, txerr) in &invalid_txs { for (txid, txerr) in &invalid_txs {
let sql = "UPDATE tbl_tx SET status = 2, push_err = ? WHERE txid = ?"; let sql = "UPDATE tbl_tx SET status = 2, push_err = ? WHERE txid = ?";
let mut stmt = db.prepare(sql).unwrap(); match db.prepare(sql) {
stmt.bind((1, Value::String(txerr.clone()))).unwrap(); Ok(mut stmt) => {
stmt.bind((2, Value::String(txid.clone()))).unwrap(); 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(); let _ = stmt.next();
} }
let _ = send_stats_report(cfg, bcinfo).await; Err(e) => {
let _ = calculate_stats(&db, network_params.db_field.clone()).await; error!("Failed to prepare error update: {}", e);
}
}
}
if let Err(e) = send_stats_report(cfg, bcinfo).await {
error!("send_stats_report failed: {}", e);
}
if let Err(e) = calculate_stats(&db, network_params.db_field.clone()).await {
warn!("calculate_stats failed: {e}");
}
} }
Err(erx) => { Err(erx) => {
error!("impossible to get client: {}, retrying on next block", erx); error!("impossible to get client: {}, retrying on next block", erx);
@@ -390,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) { if let Err(err) = db.execute(&sql) {
error!("error inserting creating stats table {err}"); error!("error inserting creating stats table {err}");
} else { } else {
@@ -422,6 +380,84 @@ ON CONFLICT(chain) DO UPDATE SET
} }
Ok(()) Ok(())
} }
/// Parse the `(host, port)` pair from a base URL like `https://host[:port]`.
///
/// Falls back to the scheme's well-known default port (443 for `https`,
/// 80 for plain `http`), or to 443 when the scheme is unknown.
fn parse_host_port(base_url: &str) -> Option<(String, u16)> {
let url = Url::parse(base_url).ok()?;
let host = url
.host_str()?
.trim_start_matches('[')
.trim_end_matches(']')
.to_string();
let port = url.port_or_known_default().unwrap_or(443);
Some((host, port))
}
/// Resolve `host:port` and return the first IPv6 (AAAA) address, if any.
///
/// Returns `None` when the host has no IPv6 address.
async fn resolve_first_ipv6(host: &str, port: u16) -> Option<SocketAddr> {
use std::net::ToSocketAddrs;
let host = host.to_string();
tokio::task::spawn_blocking(move || {
format!("{}:{}", host, port)
.to_socket_addrs()
.ok()
.and_then(|mut addrs| addrs.find(|a| a.is_ipv6()))
})
.await
.ok()
.flatten()
}
/// Build the HTTP client used for welist reports.
///
/// When `BAL_PUSHER_PREFER_IPV6` is truthy, the welist host is resolved and
/// the client is pinned to its first IPv6 address (the original hostname is
/// still used for the `Host` header and TLS SNI). This works around networks
/// where the IPv4 route to the welist host is broken while IPv6 works: the
/// default connector may otherwise pick the broken family and the request
/// stalls. When the variable is unset (the default), behavior is unchanged.
async fn welist_http_client(welist_url: &str) -> rClient {
let prefer_ipv6 = env::var("BAL_PUSHER_PREFER_IPV6")
.unwrap_or("false".to_string())
.parse::<bool>()
.unwrap_or(false);
if !prefer_ipv6 {
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 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(|_| new_welist_client())
}
None => {
debug!("BAL_PUSHER_PREFER_IPV6: no IPv6 address for {host}, using default resolver");
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( async fn send_stats_report(
cfg: &MyConfig, cfg: &MyConfig,
bcinfo: GetBlockchainInfoResult, bcinfo: GetBlockchainInfoResult,
@@ -430,14 +466,18 @@ async fn send_stats_report(
debug!("sending report to welist"); debug!("sending report to welist");
let welist_url = env::var("WELIST_SERVER_URL") let welist_url = env::var("WELIST_SERVER_URL")
.unwrap_or("https://welist.bitcoin-after.life".to_string()); .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!( warn!(
"Invalid or unsafe WELIST_SERVER_URL: {}. Skipping stats report.", "Invalid or unsafe WELIST_SERVER_URL: {}. Skipping stats report.",
welist_url welist_url
); );
return Ok(()); return Ok(());
} }
let client = rClient::new(); let client = welist_http_client(&welist_url).await;
let url = format!("{}/ping", welist_url); let url = format!("{}/ping", welist_url);
debug!("welist url: {}", url); debug!("welist url: {}", url);
let chain = bcinfo.chain.to_string().to_lowercase(); let chain = bcinfo.chain.to_string().to_lowercase();
@@ -446,7 +486,7 @@ async fn send_stats_report(
cfg.url, chain, bcinfo.blocks, bcinfo.median_time, bcinfo.best_block_hash cfg.url, chain, bcinfo.blocks, bcinfo.median_time, bcinfo.best_block_hash
); );
trace!("message to be sent: {}", message); 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 let response = client
.post(url) .post(url)
.header("User-Agent", format!("bal-pusher/{}", VERSION)) .header("User-Agent", format!("bal-pusher/{}", VERSION))
@@ -461,32 +501,23 @@ async fn send_stats_report(
})) }))
.send() .send()
.await?; .await?;
if !response.status().is_success() { let status = response.status();
warn!( let body = response.text().await?;
"Non-success response: {} {}", info!(
response.status(), "Report to welist({}) status={} body={}",
response.status().canonical_reason().unwrap_or("") welist_url, status, body
); );
}
let body = &(response.text().await?);
info!("Report to welist({})\tSent: {}", welist_url, body);
} else { } else {
debug!("Not sending stats"); debug!("Not sending stats");
} }
Ok(()) Ok(())
} }
fn sign_message(private_key_path: &str, message: &str) -> String { 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(); general_purpose::STANDARD.encode(signature.to_bytes())
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
} }
fn parse_env(cfg: &mut MyConfig) { fn parse_env(cfg: &mut MyConfig) {
@@ -505,144 +536,47 @@ fn parse_env_netconfig(cfg_lock: &mut MyConfig, chain: &str) -> NetworkParams {
"testnet4" => &mut cfg_lock.testnet4, "testnet4" => &mut cfg_lock.testnet4,
&_ => &mut cfg_lock.mainnet, &_ => &mut cfg_lock.mainnet,
}; };
match env::var(format!("BAL_PUSHER_{}_HOST", chain.to_uppercase())) { if let Ok(value) = env::var(format!("BAL_PUSHER_{}_HOST", chain.to_uppercase())) {
Ok(value) => {
cfg.host = value; cfg.host = value;
} }
Err(_) => {} if let Ok(value) = env::var(format!("BAL_PUSHER_{}_PORT", chain.to_uppercase()))
} && let Ok(port_num) = value.parse::<u64>()
match env::var(format!("BAL_PUSHER_{}_PORT", chain.to_uppercase())) { {
Ok(value) => match value.parse::<u64>() { if let Ok(port) = u16::try_from(port_num) {
Ok(value) => match u16::try_from(value) { cfg.port = port;
Ok(port) => cfg.port = port, } else {
Err(e) => {
error!( error!(
"Port value {} exceeds u16 range for chain {}: {}", "Port value {} exceeds u16 range for chain {}",
value, chain, e port_num, chain
); );
} }
},
Err(_) => {}
},
Err(_) => {}
} }
match env::var(format!("BAL_PUSHER_{}_DIR_PATH", chain.to_uppercase())) { if let Ok(value) = env::var(format!("BAL_PUSHER_{}_DIR_PATH", chain.to_uppercase())) {
Ok(value) => {
cfg.dir_path = value; cfg.dir_path = value;
} }
Err(_) => {} if let Ok(value) = env::var(format!("BAL_PUSHER_{}_DB_FIELD", chain.to_uppercase())) {
}
match env::var(format!("BAL_PUSHER_{}_DB_FIELD", chain.to_uppercase())) {
Ok(value) => {
cfg.db_field = value; cfg.db_field = value;
} }
Err(_) => {} if let Ok(value) = env::var(format!("BAL_PUSHER_{}_COOKIE_FILE", chain.to_uppercase())) {
}
match env::var(format!("BAL_PUSHER_{}_COOKIE_FILE", chain.to_uppercase())) {
Ok(value) => {
cfg.cookie_file = value; cfg.cookie_file = value;
} }
Err(_) => {} if let Ok(value) = env::var(format!("BAL_PUSHER_{}_RPC_USER", chain.to_uppercase())) {
}
match env::var(format!("BAL_PUSHER_{}_RPC_USER", chain.to_uppercase())) {
Ok(value) => {
cfg.rpc_user = value; cfg.rpc_user = value;
} }
Err(_) => {} if let Ok(value) = env::var(format!("BAL_PUSHER_{}_RPC_PASSWORD", chain.to_uppercase())) {
}
match env::var(format!("BAL_PUSHER_{}_RPC_PASSWORD", chain.to_uppercase())) {
Ok(value) => {
cfg.rpc_pass = value; cfg.rpc_pass = value;
} }
Err(_) => {} if let Ok(value) = env::var(format!("BAL_PUSHER_{}_ZMQ_HASHBLOCK", chain.to_uppercase())) {
}
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);
cfg.zmq_listener = value; cfg.zmq_listener = value;
} }
Err(_) => {}
}
cfg.clone() 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] #[tokio::main]
async fn main() -> std::io::Result<()> { async fn main() -> std::io::Result<()> {
env_logger::init(); env_logger::init();
let mut cfg = MyConfig::default(); let mut cfg = MyConfig::default();
let dbfile = env::var("BAL_PUSHER_DB_FILE").unwrap();
parse_env(&mut cfg); parse_env(&mut cfg);
let mut args = std::env::args(); let mut args = std::env::args();
let _exe_name = args.next().unwrap(); let _exe_name = args.next().unwrap();
@@ -677,31 +611,48 @@ async fn main() -> std::io::Result<()> {
} }
match socket.set_subscribe(b"") { match socket.set_subscribe(b"") {
Ok(_) => {} Ok(_) => {
info!("ZMQ subscribed to all topics on {}", zmq_address);
}
Err(e) => { Err(e) => {
error!("ZMQ subscribe failed: {}, exiting", e); error!("ZMQ subscribe failed: {}, exiting", e);
return Ok(()); 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.."); 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 socket.set_rcvtimeo(5000).unwrap(); // 5 seconds timeout
let mut consecutive_timeouts: u32 = 0;
loop { loop {
let message = match socket.recv_multipart(0) { let message = match socket.recv_multipart(0) {
Ok(m) => m, Ok(m) => m,
Err(e) => { 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; continue;
} }
}; };
if consecutive_timeouts > 0 {
info!(
"ZMQ connection restored after {} consecutive timeouts",
consecutive_timeouts
);
}
consecutive_timeouts = 0;
let topic = message[0].clone(); let topic = message[0].clone();
let body = message[1].clone(); let body = message[1].clone();
let seq = message[2].clone();
last_seq = seq;
debug!( debug!(
"ZMQ:GET TOPIC: {}", "ZMQ:GET TOPIC: {}",
String::from_utf8(topic.clone()).expect("invalid topic") String::from_utf8(topic.clone()).expect("invalid topic")
@@ -709,18 +660,52 @@ async fn main() -> std::io::Result<()> {
trace!("ZMQ:GET BODY: {}", hex::encode(&body)); trace!("ZMQ:GET BODY: {}", hex::encode(&body));
if topic == b"hashblock" { if topic == b"hashblock" {
info!("NEW BLOCK: {}", hex::encode(&body)); 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 thread::sleep(Duration::from_millis(100)); // Sleep for 100ms
} }
} }
fn seq_to_str(seq: &Vec<u8>) -> String {
if seq.len() == 4 { #[cfg(test)]
let mut rdr = Cursor::new(seq); mod tests {
let sequence = rdr use super::*;
.read_u32::<LittleEndian>()
.expect("Failed to read integer"); #[test]
return sequence.to_string(); fn parse_host_port_https_default_port() {
assert_eq!(
parse_host_port("https://welist.bitcoin-after.life"),
Some(("welist.bitcoin-after.life".to_string(), 443))
);
}
#[test]
fn parse_host_port_explicit_port_and_path() {
assert_eq!(
parse_host_port("https://example.com:8443/ping"),
Some(("example.com".to_string(), 8443))
);
}
#[test]
fn parse_host_port_http_default_port() {
assert_eq!(
parse_host_port("http://example.com"),
Some(("example.com".to_string(), 80))
);
}
#[test]
fn parse_host_port_ipv6_literal_brackets_stripped() {
assert_eq!(
parse_host_port("https://[2a13:2c0::1]:443"),
Some(("2a13:2c0::1".to_string(), 443))
);
}
#[test]
fn parse_host_port_invalid_url() {
assert_eq!(parse_host_port("not a url"), None);
} }
"Unknown".to_string()
} }

View File

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

View File

@@ -11,6 +11,7 @@ use sha2::{Digest, Sha256};
use std::str::FromStr; use std::str::FromStr;
// Mainnet (BIP44/BIP49/BIP84) // Mainnet (BIP44/BIP49/BIP84)
#[allow(dead_code)]
enum BS58Prefix { enum BS58Prefix {
Xpub, Xpub,
Ypub, Ypub,
@@ -102,44 +103,29 @@ fn calc_checksum(desc: &str) -> Result<String, String> {
Ok(checksum) Ok(checksum)
} }
pub fn get_bitcoincore_descriptor(xpub: &String) -> String { pub fn get_bitcoincore_descriptor(xpub: &str) -> String {
let fingerprint = match calculate_fingerprint(xpub) { let fingerprint = match calculate_fingerprint(xpub) {
Ok(f) => f, Ok(f) => f,
Err(_) => return String::new(), // Invalid xpub, return empty descriptor 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) { let xpub_converted = match convert_xpub(xpub) {
Ok(c) => c, Ok(c) => c,
Err(_) => return String::new(), // Invalid xpub, return empty descriptor Err(_) => return String::new(), // Invalid xpub, return empty descriptor
}; };
let descriptor = format!("wpkh([{}/84h/0h/0h]{}/0/*)", fingerprint, xpub_converted); let descriptor = format!("wpkh([{}/84h/0h/0h]{}/0/*)", fingerprint, xpub_converted);
let descriptor = match calc_checksum(&descriptor) { match calc_checksum(&descriptor) {
Ok(checksum) => { Ok(checksum) => {
let clean_descriptor = descriptor.split('#').next().unwrap_or(&descriptor); let clean_descriptor = descriptor.split('#').next().unwrap_or(&descriptor);
format!("{}#{}", clean_descriptor, checksum) format!("{}#{}", clean_descriptor, checksum)
} }
Err(err) => { Err(err) => {
eprintln!("Error: {}", 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") if xpub.len() >= 4 && (&xpub[0..4] == "xpub" || &xpub[0..4] == "ypub" || &xpub[0..4] == "zpub")
{ {
convert_to(xpub, BS58Prefix::Xpub) 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()); return Err("Data troppo corta".to_string());
} }
let (payload, checksum) = data.split_at(data.len() - 4); 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[..] { if hash[0..4] != checksum[..] {
return Err("Checksum invalido".to_string()); 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 { 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(); let full = [data, checksum].concat();
bs58::encode(full).into_string() bs58::encode(full).into_string()
} }
@@ -205,7 +191,7 @@ pub fn new_address_from_xpub(
) -> Result<(String, String), Box<dyn std::error::Error>> { ) -> Result<(String, String), Box<dyn std::error::Error>> {
let xpub = Xpub::from_str(&convert_to(zpub, BS58Prefix::Xpub)?)?; let xpub = Xpub::from_str(&convert_to(zpub, BS58Prefix::Xpub)?)?;
let path = format!("m/0/{}", index); 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 secp = Secp256k1::new();
let derived_xpub = xpub.derive_pub(&secp, &derivation_path)?; let derived_xpub = xpub.derive_pub(&secp, &derivation_path)?;
let public_key = derived_xpub.public_key; let public_key = derived_xpub.public_key;

View File

@@ -1,7 +1,6 @@
use bal_server::db::open_db; use bal_server::db::open_db;
use sqlite::State; use sqlite::State;
use std::fs; use std::fs;
use std::path::Path;
#[test] #[test]
fn test_open_db_blocks_traversal() { 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(real);
let _ = fs::remove_file(link); let _ = fs::remove_file(link);
fs::File::create(real).unwrap(); fs::File::create(real).unwrap();
fs::soft_link(real, link).unwrap(); std::os::unix::fs::symlink(real, link).unwrap();
let res = open_db(link); let res = open_db(link);
assert!(res.is_err(), "Symlink DB path should be rejected"); 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 mut found_value = None;
let _ = db.iterate("SELECT * FROM test_stats;", |pairs| { let _ = db.iterate("SELECT * FROM test_stats;", |pairs| {
let row: HashMap<_, _> = pairs let row: HashMap<_, _> = pairs.iter().map(|(k, v)| (k.to_string(), *v)).collect();
.into_iter() let totals = row["totals"].unwrap_or("0").to_string();
.map(|(k, v)| (k.to_string(), v.map(|s| s)))
.collect();
let totals = row["totals"].clone().unwrap_or("0").to_string();
found_value = Some(totals); found_value = Some(totals);
true true
}); });

View File

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