use evidence_chain::{EvidenceCategory, EvidenceChain, EvidenceLink};
use crate::heuristics::{Heuristic, HeuristicStatus};
use crate::match_result::HeuristicMatch;
use crate::models::TxFeatures;
pub struct DistributionHeuristic;
impl Heuristic for DistributionHeuristic {
fn id(&self) -> &'static str {
"institutional-distribution-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_distribution || f.output_count < 20 {
return None;
}
let btc_price = f.block_price_usd?;
let total_btc = f.total_output_value as f64 / 100_000_000.0;
let total_usd = total_btc * btc_price;
if total_usd < 50_000.0 {
return None;
}
let total_outs = f.output_count.max(1);
let max_script_count = [
f.output_p2pkh_count,
f.output_p2sh_count,
f.output_p2wpkh_count,
f.output_p2wsh_count,
f.output_p2tr_count,
]
.iter()
.copied()
.max()
.unwrap_or(0);
let output_homogeneity = max_script_count as f64 / total_outs as f64;
if output_homogeneity < 0.70 {
return None;
}
let summary = format!(
"Institutional distribution: {} inputs → {} outputs (${:.0} USD) at block {}",
f.input_count, f.output_count, total_usd, f.block_height
);
Some(HeuristicMatch::new(
self.id(),
self.version(),
"institutional_distribution",
"distribution",
self.trigger_scope(),
summary,
serde_json::json!({
"input_count": f.input_count,
"output_count": f.output_count,
"total_output_value": f.total_output_value,
"output_value_min": f.output_value_min,
"output_value_max": f.output_value_max,
"output_value_variance": f.output_value_variance,
"output_homogeneity": output_homogeneity,
}),
))
}
fn build_evidence(&self, f: &TxFeatures) -> Option<EvidenceChain> {
if f.is_coinbase || !f.is_distribution || f.output_count < 20 {
return None;
}
let btc_price = f.block_price_usd?;
let total_btc = f.total_output_value as f64 / 100_000_000.0;
let total_usd = total_btc * btc_price;
if total_usd < 50_000.0 {
return None;
}
let total_outs = f.output_count.max(1);
let max_script_count = [
f.output_p2pkh_count,
f.output_p2sh_count,
f.output_p2wpkh_count,
f.output_p2wsh_count,
f.output_p2tr_count,
]
.iter()
.copied()
.max()
.unwrap_or(0);
let output_homogeneity = max_script_count as f64 / total_outs as f64;
if output_homogeneity < 0.70 {
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 distribution (over $50,000 USD)".to_string(),
txid_hex.clone(),
)
.with_metric(total_usd, "USD")
.with_threshold(50_000.0, total_usd >= 50_000.0),
);
chain.add_link(
EvidenceLink::new(
EvidenceCategory::Structural,
format!(
"{} outputs (high fan-out, distribution threshold ≥20, homogeneity={:.2})",
f.output_count, output_homogeneity
),
txid_hex.clone(),
)
.with_metric(f.output_count as f64, "outputs")
.with_threshold(20.0, f.output_count >= 20),
);
let value_range = f.output_value_max - f.output_value_min;
chain.add_link(
EvidenceLink::new(
EvidenceCategory::Value,
format!(
"Output value range: {} sat (min={}, max={})",
value_range, f.output_value_min, f.output_value_max
),
txid_hex,
)
.with_metric(f.output_value_variance, "sat²"),
);
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![0x03u8; 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_990_000,
fee: 10_000,
output_value_min: 50_000,
output_value_max: 1_000_000,
output_value_median: 200_000.0,
output_value_variance: 1e12,
fee_rate_sat_vb: Some(5.0),
input_p2sh_count: input_count,
output_p2wpkh_count: output_count,
is_consolidation: input_count >= 10 && 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: 2_000,
tx_version: 1,
input_utxo_refs: (0..input_count as u32)
.map(|i| (vec![0xddu8; 32], i))
.collect(),
output_values: vec![200_000; output_count as usize],
block_price_usd: Some(60_000.0),
..Default::default()
}
}
#[test]
fn test_batching_detected() {
let h = DistributionHeuristic;
let f = make_features(1, 50, false);
let result = h.evaluate(&f);
assert!(
result.is_some(),
"1 input + 50 outputs should fire distribution"
);
let r = result.unwrap();
assert_eq!(r.event_type, "institutional_distribution");
assert_eq!(r.pattern, "distribution");
}
#[test]
fn test_batching_exact_threshold() {
let h = DistributionHeuristic;
let f = make_features(3, 20, false);
assert!(
h.evaluate(&f).is_some(),
"3 inputs + 20 outputs = exact threshold"
);
}
#[test]
fn test_not_batching_too_many_inputs() {
let h = DistributionHeuristic;
let f = make_features(5, 50, false);
assert!(h.evaluate(&f).is_none(), "5 inputs is not distribution");
}
#[test]
fn test_not_batching_too_few_outputs() {
let h = DistributionHeuristic;
let f = make_features(1, 19, false);
assert!(h.evaluate(&f).is_none(), "19 outputs is not distribution");
}
#[test]
fn test_batching_skips_coinbase() {
let h = DistributionHeuristic;
let f = make_features(1, 50, true);
assert!(h.evaluate(&f).is_none(), "coinbase is never distribution");
}
#[test]
fn test_distribution_no_price_skips() {
let h = DistributionHeuristic;
let mut f = make_features(1, 50, false);
f.block_price_usd = None;
assert!(h.evaluate(&f).is_none(), "no price should skip");
}
#[test]
fn test_distribution_low_usd_skips() {
let h = DistributionHeuristic;
let mut f = make_features(1, 50, false);
f.total_output_value = 10_000_000;
assert!(h.evaluate(&f).is_none(), "low USD value should skip");
}
}