use evidence_chain::{EvidenceCategory, EvidenceChain, EvidenceLink};
use crate::heuristics::{Heuristic, HeuristicStatus};
use crate::match_result::HeuristicMatch;
use crate::models::TxFeatures;
pub struct UtxoConsolidationHeuristic;
impl Heuristic for UtxoConsolidationHeuristic {
fn id(&self) -> &'static str {
"institutional-consolidation-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.is_consolidation {
return None;
}
let btc_price = f.block_price_usd?;
let total_btc = f.total_input_value as f64 / 100_000_000.0;
let total_usd_estimated = total_btc * btc_price;
if total_usd_estimated < 100_000.0 {
return None;
}
if f.input_count < 10 {
return None;
}
if !f.is_input_script_homogeneous {
return None;
}
let max_output = f.output_value_max as f64;
let total_output = f.total_output_value as f64;
if total_output > 0.0 && (max_output / total_output) < 0.7 {
return None;
}
if f.has_equal_outputs && f.output_count >= 2 {
return None;
}
let summary = format!(
"Institutional consolidation: {} inputs → {} outputs (${:.0} USD) at block {}",
f.input_count, f.output_count, total_usd_estimated, f.block_height
);
Some(HeuristicMatch::new(
self.id(),
self.version(),
"institutional_consolidation",
"structural_consolidation",
self.trigger_scope(),
summary,
serde_json::json!({
"input_count": f.input_count,
"output_count": f.output_count,
"total_input_value": f.total_input_value,
"fee": f.fee,
"fee_rate_sat_vb": f.fee_rate_sat_vb,
"is_input_script_homogeneous": f.is_input_script_homogeneous,
"has_equal_outputs": f.has_equal_outputs,
"dominant_input_script": format!("{:?}", f.dominant_input_script),
}),
))
}
fn build_evidence(&self, f: &TxFeatures) -> Option<EvidenceChain> {
if f.is_coinbase || !f.is_consolidation {
return None;
}
let btc_price = f.block_price_usd?;
let total_btc = f.total_input_value as f64 / 100_000_000.0;
let total_usd_estimated = total_btc * btc_price;
if total_usd_estimated < 100_000.0 {
return None;
}
if f.input_count < 10 || !f.is_input_script_homogeneous {
return None;
}
let total_output = f.total_output_value as f64;
if total_output > 0.0 && (f.output_value_max as f64 / total_output) < 0.7 {
return None;
}
if f.has_equal_outputs && f.output_count >= 2 {
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());
chain.add_link(
EvidenceLink::new(
EvidenceCategory::Structural,
"Large institutional consolidation (over $100,000 USD)".to_string(),
txid_hex.clone(),
)
.with_metric(total_usd_estimated, "USD")
.with_threshold(100_000.0, total_usd_estimated >= 100_000.0),
);
chain.add_link(
EvidenceLink::new(
EvidenceCategory::Structural,
format!("{} inputs consumed in single transaction", f.input_count),
txid_hex.clone(),
)
.with_metric(f.input_count as f64, "inputs")
.with_threshold(10.0, f.input_count >= 10),
);
chain.add_link(
EvidenceLink::new(
EvidenceCategory::Structural,
format!(
"{} outputs produced (consolidation limit ≤3)",
f.output_count
),
txid_hex,
)
.with_metric(f.output_count as f64, "outputs")
.with_threshold(3.0, f.output_count <= 3),
);
chain.finalize();
Some(chain)
}
}
#[cfg(test)]
mod tests {
use super::*;
use chrono::Utc;
fn make_features(input_count: i32, output_count: i32, is_coinbase: bool) -> TxFeatures {
TxFeatures {
txid: vec![0x01u8; 32],
block_height: 840_000,
block_timestamp: Utc::now(),
input_count,
output_count,
is_coinbase,
total_input_value: 200_000_000,
total_output_value: 199_995_000,
fee: 5_000,
output_value_min: 100_000,
output_value_max: 199_895_000,
output_value_median: 497_500.0,
output_value_variance: 1e14,
fee_rate_sat_vb: Some(10.0),
input_p2pkh_count: input_count,
output_p2pkh_count: output_count,
is_consolidation: (input_count >= 3 && output_count <= 2)
|| (input_count >= 20 && output_count <= 3),
is_distribution: input_count <= 3 && output_count >= 10,
is_simple_send: input_count == 1 && output_count == 2,
is_sweep: input_count >= 2 && output_count == 1,
tx_vsize_vbytes: 500,
tx_version: 1,
input_utxo_refs: (0..input_count as u32)
.map(|i| (vec![0xffu8; 32], i))
.collect(),
output_values: vec![100_000; output_count as usize],
is_input_script_homogeneous: true,
block_price_usd: Some(60_000.0),
..Default::default()
}
}
#[test]
fn test_consolidation_detected() {
let h = UtxoConsolidationHeuristic;
let f = make_features(15, 2, false);
let result = h.evaluate(&f);
assert!(
result.is_some(),
"15 inputs → 2 outputs with USD price should fire"
);
let r = result.unwrap();
assert_eq!(r.event_type, "institutional_consolidation");
}
#[test]
fn test_consolidation_no_price_skips() {
let h = UtxoConsolidationHeuristic;
let mut f = make_features(15, 2, false);
f.block_price_usd = None;
assert!(h.evaluate(&f).is_none(), "no price should skip");
}
#[test]
fn test_consolidation_low_usd_value_skips() {
let h = UtxoConsolidationHeuristic;
let mut f = make_features(15, 2, false);
f.total_input_value = 10_000_000;
assert!(h.evaluate(&f).is_none(), "low USD value should skip");
}
#[test]
fn test_consolidation_exact_threshold() {
let h = UtxoConsolidationHeuristic;
let f = make_features(10, 2, false);
assert!(
h.evaluate(&f).is_some(),
"10 inputs + sufficient USD value should fire"
);
}
#[test]
fn test_not_consolidation_too_few_inputs() {
let h = UtxoConsolidationHeuristic;
let f = make_features(5, 2, false);
assert!(
h.evaluate(&f).is_none(),
"5 inputs is no longer high volume consolidation"
);
}
#[test]
fn test_not_consolidation_too_many_outputs() {
let h = UtxoConsolidationHeuristic;
let f = make_features(35, 5, false);
assert!(
h.evaluate(&f).is_none(),
"35 inputs + 5 outputs is not consolidation"
);
}
#[test]
fn test_consolidation_skips_coinbase() {
let h = UtxoConsolidationHeuristic;
let f = make_features(30, 1, true);
assert!(h.evaluate(&f).is_none(), "coinbase is never consolidation");
}
}