use chrono::{DateTime, TimeZone, Utc};
use evidence_chain::{EvidenceCategory, EvidenceChain, EvidenceLink};
use crate::heuristics::{Heuristic, HeuristicStatus};
use crate::match_result::HeuristicMatch;
use crate::models::TxFeatures;
#[derive(Debug, Clone)]
pub struct TxOutInfo {
pub confirmations: u64,
#[allow(dead_code)]
pub block_time: Option<DateTime<Utc>>,
}
pub trait BitcoinRpc: Send + Sync {
fn get_tx_out(&self, txid_hex: &str, vout: u32) -> anyhow::Result<Option<TxOutInfo>>;
}
pub struct JsonRpcClient {
url: String,
user: Option<String>,
pass: Option<String>,
agent: ureq::Agent,
max_retries: u8,
cache: std::sync::Mutex<std::collections::HashMap<(String, u32), Option<u64>>>,
}
impl JsonRpcClient {
pub fn new(url: &str, user: &str, pass: &str) -> Self {
let agent = ureq::AgentBuilder::new()
.timeout_connect(std::time::Duration::from_secs(5))
.timeout_read(std::time::Duration::from_secs(10))
.timeout_write(std::time::Duration::from_secs(10))
.max_idle_connections(200)
.max_idle_connections_per_host(200)
.build();
Self {
url: url.trim().to_string(),
user: if user.is_empty() {
None
} else {
Some(user.to_string())
},
pass: if pass.is_empty() {
None
} else {
Some(pass.to_string())
},
agent,
max_retries: 2,
cache: std::sync::Mutex::new(std::collections::HashMap::new()),
}
}
fn call(&self, method: &str, params: serde_json::Value) -> anyhow::Result<serde_json::Value> {
let body = serde_json::json!({
"jsonrpc": "2.0",
"id": "bitcoin-heuristics",
"method": method,
"params": params,
});
let mut attempt: u8 = 0;
loop {
let mut req = self
.agent
.post(&self.url)
.set("Content-Type", "application/json");
if let (Some(user), Some(pass)) = (&self.user, &self.pass) {
let credentials = base64_encode(&format!("{user}:{pass}"));
req = req.set("Authorization", &format!("Basic {credentials}"));
}
let resp = match req.send_string(&body.to_string()) {
Ok(r) => r,
Err(e) => {
if attempt < self.max_retries {
attempt += 1;
let backoff_ms = 200u64.saturating_mul(1u64 << (attempt - 1));
tracing::warn!(error = %e, attempt, "RPC request failed, retrying");
std::thread::sleep(std::time::Duration::from_millis(backoff_ms));
continue;
} else {
return Err(anyhow::anyhow!("RPC request failed after retries: {e}"));
}
}
};
let resp_json: serde_json::Value = resp
.into_json()
.map_err(|e| anyhow::anyhow!("Failed to parse RPC response: {e}"))?;
if let Some(error) = resp_json.get("error").filter(|e| !e.is_null()) {
return Err(anyhow::anyhow!("RPC error: {error}"));
}
return Ok(resp_json["result"].clone());
}
}
}
fn base64_encode(input: &str) -> String {
use std::fmt::Write;
const TABLE: &[u8; 64] = b"ABCDEFGHIJKLMNOPQRSTUVWXYZabcdefghijklmnopqrstuvwxyz0123456789+/";
let bytes = input.as_bytes();
let mut out = String::new();
let mut i = 0;
while i < bytes.len() {
let b0 = bytes[i];
let b1 = if i + 1 < bytes.len() { bytes[i + 1] } else { 0 };
let b2 = if i + 2 < bytes.len() { bytes[i + 2] } else { 0 };
let _ = write!(out, "{}", TABLE[(b0 >> 2) as usize] as char);
let _ = write!(out, "{}", TABLE[((b0 & 3) << 4 | b1 >> 4) as usize] as char);
let _ = write!(
out,
"{}",
if i + 1 < bytes.len() {
TABLE[((b1 & 0xF) << 2 | b2 >> 6) as usize] as char
} else {
'='
}
);
let _ = write!(
out,
"{}",
if i + 2 < bytes.len() {
TABLE[(b2 & 0x3F) as usize] as char
} else {
'='
}
);
i += 3;
}
out
}
impl BitcoinRpc for JsonRpcClient {
fn get_tx_out(&self, txid_hex: &str, vout: u32) -> anyhow::Result<Option<TxOutInfo>> {
let cache_key = (txid_hex.to_string(), vout);
if let Ok(cache) = self.cache.lock() {
if let Some(cached) = cache.get(&cache_key) {
return Ok(cached.map(|confs| TxOutInfo {
confirmations: confs,
block_time: None,
}));
}
}
let result = self.call("gettxout", serde_json::json!([txid_hex, vout, false]))?;
if result.is_null() {
if let Ok(mut cache) = self.cache.lock() {
cache.insert(cache_key, None);
}
return Ok(None);
}
let confirmations = result["confirmations"].as_u64().unwrap_or(0);
let block_time = result["blockTime"]
.as_i64()
.map(|ts| Utc.timestamp_opt(ts, 0).single().unwrap_or_else(Utc::now));
if let Ok(mut cache) = self.cache.lock() {
cache.insert(cache_key, Some(confirmations));
}
Ok(Some(TxOutInfo {
confirmations,
block_time,
}))
}
}
pub struct LongTermSupplyActivationHeuristic<R: BitcoinRpc> {
pub rpc: R,
pub min_activation_confirmations: u64,
}
impl<R: BitcoinRpc> LongTermSupplyActivationHeuristic<R> {
pub fn new(rpc: R, min_activation_confirmations: u64) -> Self {
Self {
rpc,
min_activation_confirmations,
}
}
}
impl<R: BitcoinRpc> Heuristic for LongTermSupplyActivationHeuristic<R> {
fn id(&self) -> &'static str {
"long-term-supply-activation-v1"
}
fn version(&self) -> &'static str {
"1.0.0"
}
fn status(&self) -> HeuristicStatus {
HeuristicStatus::Active
}
fn evaluate(&self, f: &TxFeatures) -> Option<HeuristicMatch> {
if f.is_coinbase || f.input_utxo_refs.is_empty() {
return None;
}
let total_usd =
(f.total_input_value as f64 * f.block_price_usd.unwrap_or(0.0)) / 100_000_000.0;
if total_usd < 50_000.0 {
return None;
}
let mut activation_refs: Vec<String> = Vec::new();
let mut max_confirmations: u64 = 0;
let mut sum_confirmations: u64 = 0;
for (prev_txid, prev_vout) in &f.input_utxo_refs {
let txid_hex: String = prev_txid.iter().rev().map(|b| format!("{b:02x}")).collect();
let age_blocks: Option<u64> =
if let Some(&age) = f.utxo_ages_blocks.get(&(prev_txid.clone(), *prev_vout)) {
if age >= 0 {
Some(age as u64)
} else {
None
}
} else {
match self.rpc.get_tx_out(&txid_hex, *prev_vout) {
Ok(Some(info)) => Some(info.confirmations),
Ok(None) => None,
Err(e) => {
tracing::warn!(
error = %e,
txid = %txid_hex,
vout = prev_vout,
"Failed to query UTXO for activation check"
);
None
}
}
};
if let Some(age) = age_blocks {
if age >= self.min_activation_confirmations {
activation_refs.push(format!("{}:{}", txid_hex, prev_vout));
max_confirmations = max_confirmations.max(age);
sum_confirmations += age;
}
}
}
if activation_refs.is_empty() {
return None;
}
let activation_count = activation_refs.len();
let total_input_count = f.input_utxo_refs.len().max(1);
let activation_ratio = activation_count as f64 / total_input_count as f64;
let avg_confirmations = sum_confirmations / activation_count as u64;
if activation_ratio < 0.5 && avg_confirmations < 144_000 {
return None;
}
let approx_years = max_confirmations / 52_560;
let summary = format!(
"Long-term supply activation: {activation_count} supply UTXO(s) \
(≥{} blocks, ~{approx_years} year(s)) activated at block {}",
self.min_activation_confirmations, f.block_height
);
Some(HeuristicMatch::new(
self.id(),
self.version(),
"long_term_supply_activation",
"long_term_supply_activation",
self.trigger_scope(),
summary,
serde_json::json!({
"activation_utxo_count": activation_count,
"max_confirmations": max_confirmations,
"avg_activation_confirmations": avg_confirmations,
"activation_ratio": activation_ratio,
"min_activation_confirmations": self.min_activation_confirmations,
"input_count": f.input_count,
"script_type_evolution": f.script_type_evolution,
"total_usd_value": total_usd,
}),
))
}
fn build_evidence(&self, f: &TxFeatures) -> Option<EvidenceChain> {
if f.is_coinbase || f.input_utxo_refs.is_empty() {
return None;
}
let total_usd =
(f.total_input_value as f64 * f.block_price_usd.unwrap_or(0.0)) / 100_000_000.0;
if total_usd < 50_000.0 {
return None;
}
let txid_hex: String = f.txid.iter().rev().map(|b| format!("{b:02x}")).collect();
let mut chain = EvidenceChain::new(self.id(), self.version());
let (first_prev_txid, first_prev_vout) = &f.input_utxo_refs[0];
let utxo_ref = format!(
"{}:{}",
first_prev_txid
.iter()
.rev()
.map(|b| format!("{b:02x}"))
.collect::<String>(),
first_prev_vout
);
chain.add_link(
EvidenceLink::new(
EvidenceCategory::Temporal,
format!(
"Input UTXO dormant for ≥{} blocks (~{} years) before activation",
self.min_activation_confirmations,
self.min_activation_confirmations / 52_560
),
utxo_ref,
)
.with_metric(self.min_activation_confirmations as f64, "blocks")
.with_threshold(self.min_activation_confirmations as f64, true),
);
chain.add_link(
EvidenceLink::new(
EvidenceCategory::Value,
format!(
"Total input value {} sat (~${:.2}) reactivated in tx {}",
f.total_input_value,
total_usd,
&txid_hex[..8]
),
txid_hex,
)
.with_metric(total_usd, "usd"),
);
chain.finalize();
Some(chain)
}
}
#[cfg(test)]
mod tests {
use super::*;
use chrono::Utc;
struct MockRpc {
dormant_confirmations: u64,
fail: bool,
}
impl BitcoinRpc for MockRpc {
fn get_tx_out(&self, _txid_hex: &str, _vout: u32) -> anyhow::Result<Option<TxOutInfo>> {
if self.fail {
return Err(anyhow::anyhow!("mock RPC failure"));
}
Ok(Some(TxOutInfo {
confirmations: self.dormant_confirmations,
block_time: None,
}))
}
}
struct MockRpcNeverDormant;
impl BitcoinRpc for MockRpcNeverDormant {
fn get_tx_out(&self, _: &str, _: u32) -> anyhow::Result<Option<TxOutInfo>> {
Ok(Some(TxOutInfo {
confirmations: 100,
block_time: None,
}))
}
}
struct MockRpcSpent;
impl BitcoinRpc for MockRpcSpent {
fn get_tx_out(&self, _: &str, _: u32) -> anyhow::Result<Option<TxOutInfo>> {
Ok(None)
}
}
fn make_features(utxo_refs: Vec<(Vec<u8>, u32)>) -> TxFeatures {
TxFeatures {
txid: vec![0x06u8; 32],
block_height: 840_000,
block_timestamp: Utc::now(),
input_count: utxo_refs.len() as i32,
output_count: 2,
is_coinbase: false,
total_input_value: 10_000_000_000,
total_output_value: 9_995_000_000,
fee: 5_000_000,
output_value_max: 9_994_000_000,
block_price_usd: Some(100_000.0),
fee_rate_sat_vb: Some(3.0),
input_p2pkh_count: utxo_refs.len() as i32,
output_p2pkh_count: 2,
tx_vsize_vbytes: 300,
tx_version: 1,
input_utxo_refs: utxo_refs,
output_values: vec![1_000_000, 3_995_000],
..Default::default()
}
}
#[test]
fn test_activation_detected() {
let rpc = MockRpc {
dormant_confirmations: 200_000,
fail: false,
};
let h = LongTermSupplyActivationHeuristic::new(rpc, 157_680);
let f = make_features(vec![(vec![0xaau8; 32], 0)]);
let result = h.evaluate(&f);
assert!(result.is_some(), "UTXO with 200k confirmations should fire");
assert_eq!(result.unwrap().event_type, "long_term_supply_activation");
}
#[test]
fn test_not_long_term_too_recent() {
let rpc = MockRpcNeverDormant;
let h = LongTermSupplyActivationHeuristic::new(rpc, 157_680);
let f = make_features(vec![(vec![0xbbu8; 32], 0)]);
assert!(h.evaluate(&f).is_none(), "Recent UTXO does not fire");
}
#[test]
fn test_spent_utxo_skipped() {
let rpc = MockRpcSpent;
let h = LongTermSupplyActivationHeuristic::new(rpc, 157_680);
let f = make_features(vec![(vec![0xccu8; 32], 0)]);
assert!(h.evaluate(&f).is_none(), "Spent UTXO does not fire");
}
#[test]
fn test_rpc_failure_skipped() {
let rpc = MockRpc {
dormant_confirmations: 0,
fail: true,
};
let h = LongTermSupplyActivationHeuristic::new(rpc, 157_680);
let f = make_features(vec![(vec![0xddu8; 32], 0)]);
assert!(h.evaluate(&f).is_none(), "RPC failure should not fire");
}
#[test]
fn test_coinbase_skipped() {
let rpc = MockRpc {
dormant_confirmations: 200_000,
fail: false,
};
let h = LongTermSupplyActivationHeuristic::new(rpc, 157_680);
let mut f = make_features(vec![(vec![0xeeu8; 32], 0)]);
f.is_coinbase = true;
assert!(h.evaluate(&f).is_none(), "coinbase ignored");
}
#[test]
fn test_no_utxo_refs_skipped() {
let rpc = MockRpc {
dormant_confirmations: 200_000,
fail: false,
};
let h = LongTermSupplyActivationHeuristic::new(rpc, 157_680);
let f = make_features(vec![]);
assert!(h.evaluate(&f).is_none(), "without UTXO refs does not fire");
}
#[test]
fn test_activation_uses_utxo_ages_map() {
struct PanicRpc;
impl BitcoinRpc for PanicRpc {
fn get_tx_out(&self, _: &str, _: u32) -> anyhow::Result<Option<TxOutInfo>> {
panic!("RPC must not be called when utxo_ages_blocks map is populated")
}
}
let h = LongTermSupplyActivationHeuristic::new(PanicRpc, 157_680);
let mut f = make_features(vec![(vec![0xaau8; 32], 0)]);
f.utxo_ages_blocks.insert((vec![0xaau8; 32], 0), 200_000);
let result = h.evaluate(&f);
assert!(result.is_some(), "should detect with map age without RPC");
assert_eq!(result.unwrap().event_type, "long_term_supply_activation");
}
#[test]
fn test_long_term_not_activated_in_map() {
struct PanicRpc;
impl BitcoinRpc for PanicRpc {
fn get_tx_out(&self, _: &str, _: u32) -> anyhow::Result<Option<TxOutInfo>> {
panic!("RPC must not be called when map is populated")
}
}
let h = LongTermSupplyActivationHeuristic::new(PanicRpc, 157_680);
let mut f = make_features(vec![(vec![0xbbu8; 32], 0)]);
f.utxo_ages_blocks.insert((vec![0xbbu8; 32], 0), 1_000);
assert!(h.evaluate(&f).is_none(), "age < threshold should not fire");
}
#[test]
fn test_activation_fallback_to_rpc() {
let rpc = MockRpc {
dormant_confirmations: 200_000,
fail: false,
};
let h = LongTermSupplyActivationHeuristic::new(rpc, 157_680);
let f = make_features(vec![(vec![0xffu8; 32], 0)]);
let result = h.evaluate(&f);
assert!(result.is_some(), "fallback to RPC should detect activation");
}
#[test]
fn test_base64_encode() {
assert_eq!(base64_encode(""), "");
assert_eq!(base64_encode("a"), "YQ==");
assert_eq!(base64_encode("abc"), "YWJj");
assert_eq!(base64_encode("user:pass"), "dXNlcjpwYXNz");
}
}