dApp Docs/链上安全监控与威胁检测指南
Development reference. Not independently verified for production.

链上安全监控与威胁检测指南

数据来源:MSG Chain 代码库核实

主网状态: No-Go — 当前 MSGChain 主网裁决为 No-Go,以下内容反映代码实际状态,不代表生产可用。


目录

  1. 链上威胁建模
  2. 实时交易监控架构
  3. 合约行为基线
  4. 告警规则引擎
  5. MEV 攻击检测
  6. 链上蜜罐 / Rug Pull 检测
  7. MSG Chain 安全监控栈
  8. 自动化响应
  9. 威胁情报
  10. 安全仪表盘设计与指标
  11. 合规监控
  12. 案例:监控 + 自动响应实战演练

1. 链上威胁建模

1.1 MSG Chain 攻击面总览

MSG Chain 基于 Cosmos SDK + DAR 共识 + CosmWasm 虚拟机,攻击面涵盖共识层、虚拟机层、应用层和合约层。

攻击面层次
┌─────────────────────────────────────────────────────┐
│  共识层攻击 (DAR 操纵、Stake 集中化、长程攻击)       │
├─────────────────────────────────────────────────────┤
│  虚拟机层攻击 (Gas 攻击、存储碰撞、沙箱逃逸)         │
├─────────────────────────────────────────────────────┤
│  应用层攻击 (IBC 伪造包、跨链重放、治理操纵)         │
├─────────────────────────────────────────────────────┤
│  合约层攻击 (重入、权限滥用、预言机操纵、闪电贷)      │
└─────────────────────────────────────────────────────┘

1.2 智能合约攻击分类

1.2.1 重入攻击 (Reentrancy)

原理: CosmWasm 通过 SubMsg + Reply 机制支持异步回调。攻击者利用 reply handler 在合约状态更新前重新进入关键函数。

MSG Chain 上重入攻击特征:

检测指标:

指标 正常范围 异常信号
单区块内同一合约调用次数 1-3 次 > 10 次
同一交易内状态读取/写入比 1:1 - 2:1 > 5:1
Reply 调用栈深度 0-1 层 > 3 层
同一 sender 连续调用间隔 > 10 区块 同区块
pub fn detect_reentrancy(
    deps: DepsMut,
    env: Env,
    calls: Vec<ContractCall>,
) -> StdResult<DetectionResult> {
    let mut call_counts: HashMap<String, u32> = HashMap::new();
    let mut reentrant_calls = Vec::new();

    for call in &calls {
        let key = format!("{}:{}", call.sender, call.contract);
        *call_counts.entry(key.clone()).or_insert(0) += 1;

        if *call_counts.get(&key).unwrap_or(&0) > 3 {
            reentrant_calls.push(ReentrantCall {
                contract: call.contract.clone(),
                sender: call.sender.clone(),
                function: call.function.clone(),
                tx_hash: call.tx_hash.clone(),
                block_height: env.block.height,
            });
        }
    }

    Ok(DetectionResult {
        detected: !reentrant_calls.is_empty(),
        severity: if reentrant_calls.len() > 5 {
            Severity::Critical
        } else if !reentrant_calls.is_empty() {
            Severity::High
        } else {
            Severity::None
        },
        evidence: reentrant_calls,
    })
}

1.2.2 闪电贷攻击 (Flash Loan Attack)

原理: CosmWasm 虽不原生支持闪电贷,但通过复合 SubMsg 链可以实现类似效果。攻击者在单笔交易中借入大量资金操纵价格后归还。

MSG Chain 上闪电贷攻击特征:

检测模型:

检测逻辑:
1. 解析交易中的所有 SubMsg 序列
2. 构建资金流向图 (start → intermediate → end)
3. 检测闭环: 同一合约的资金流出和流入发生在同一交易
4. 检查价格变化: 交易前后预言机价格差异 > 阈值
5. 检查流动性池余额变化: 临时耗尽某侧流动性
pub struct FlashLoanDetector {
    pub min_loan_amount: Uint128,
    pub max_calls_in_tx: u32,
    pub price_deviation_threshold: Decimal,
}

impl FlashLoanDetector {
    pub fn analyze_tx(
        &self,
        deps: DepsMut,
        tx: AnalyzedTransaction,
    ) -> DetectionResult {
        let mut alerts = Vec::new();
        let mut flow_graph = FlowGraph::new();

        for msg in &tx.messages {
            if let Some(transfer) = self.extract_transfer(msg) {
                flow_graph.add_edge(
                    &transfer.from, &transfer.to, transfer.amount,
                );
            }
        }

        if flow_graph.has_cycle() && tx.total_volume > self.min_loan_amount {
            alerts.push(Alert {
                alert_type: AlertType::FlashLoan,
                severity: Severity::Critical,
                message: format!(
                    "Flash loan detected: {} amsg cycled in tx {}",
                    tx.total_volume, tx.tx_hash,
                ),
            });
        }

        if let Some(price_change) = self.detect_price_manipulation(deps, &tx) {
            if price_change > self.price_deviation_threshold {
                alerts.push(Alert {
                    alert_type: AlertType::PriceManipulation,
                    severity: Severity::Critical,
                    message: format!("Price swing {:.2}% in tx", price_change * 100),
                });
            }
        }

        DetectionResult {
            detected: !alerts.is_empty(),
            severity: alerts.iter()
                .map(|a| a.severity)
                .max().unwrap_or(Severity::None),
            alerts,
        }
    }
}

1.2.3 预言机操纵 (Oracle Manipulation)

原理: 攻击者通过操纵链上数据源(AMM 池价、TWAP 计算窗口等)影响依赖预言机的合约决策。

MSG Chain 上预言机操纵攻击特征:

检测策略:

检测项 正常 异常
价格源数量 >= 3 1
价格偏差 (与链下参考价) < 2% > 10%
价格更新频率 每区块 长时间未更新后突变
交易额 vs 池流动性 < 1% > 20% 单笔
pub fn detect_oracle_manipulation(
    deps: DepsMut,
    env: Env,
    oracle_prices: Vec<PriceFeed>,
    reference_prices: HashMap<String, Decimal>,
) -> StdResult<Vec<OracleAlert>> {
    let mut alerts = Vec::new();

    for price in &oracle_prices {
        if let Some(ref_price) = reference_prices.get(&price.asset) {
            let deviation = (price.price - *ref_price) / *ref_price;
            if deviation.abs() > Decimal::percent(5) {
                alerts.push(OracleAlert {
                    asset: price.asset.clone(),
                    oracle_price: price.price,
                    reference_price: *ref_price,
                    deviation: deviation.abs(),
                    source: price.source.clone(),
                    severity: if deviation.abs() > Decimal::percent(20) {
                        Severity::Critical
                    } else {
                        Severity::High
                    },
                    detected_at: env.block.time.seconds(),
                });
            }
        }

        let history = PRICE_HISTORY
            .may_load(deps.storage, &price.asset)?
            .unwrap_or_default();

        if let Some(last_price) = history.last() {
            let change = (price.price - last_price.price) / last_price.price;
            if change.abs() > Decimal::percent(10) {
                alerts.push(OracleAlert {
                    asset: price.asset.clone(),
                    oracle_price: price.price,
                    reference_price: last_price.price,
                    deviation: change.abs(),
                    source: format!("{}_sudden_change", price.source),
                    severity: Severity::High,
                    detected_at: env.block.time.seconds(),
                });
            }
        }
    }

    Ok(alerts)
}

1.2.4 权限滥用 (Privilege Abuse)

原理: 合约中权限控制不完善或存在未预期的特权路径,攻击者利用提升的权限执行越权操作。

MSG Chain 上权限滥用攻击特征:

检测规则:

权限操作 调用者 风险等级
migrate 非合约 owner Critical
sudo 非验证者 Critical
execute > admin 函数 非授权地址 High
Constitution 更新 非 Agent owner High
白名单修改 非管理员 Medium
pub fn detect_privilege_abuse(
    deps: DepsMut,
    tx: AnalyzedTransaction,
    expected_admin: &Addr,
    expected_operators: &[Addr],
) -> DetectionResult {
    let mut alerts = Vec::new();

    for msg in &tx.messages {
        match msg {
            TransactionMsg::Migrate { sender, .. } => {
                if sender != expected_admin {
                    alerts.push(Alert {
                        alert_type: AlertType::UnauthorizedMigrate,
                        severity: Severity::Critical,
                        message: format!("Non-admin migrate attempt by {}", sender),
                    });
                }
            }
            TransactionMsg::Execute { sender, function, .. }
                if function == "update_owner" || function == "transfer_ownership" =>
            {
                if sender != expected_admin {
                    alerts.push(Alert {
                        alert_type: AlertType::OwnershipChange,
                        severity: Severity::Critical,
                        message: format!("Unauthorized ownership change by {}", sender),
                    });
                }
            }
            _ => {}
        }
    }

    DetectionResult {
        detected: !alerts.is_empty(),
        severity: alerts.iter().map(|a| a.severity).max().unwrap_or(Severity::None),
        alerts,
    }
}

1.2.5 逻辑漏洞 (Logic Bugs)

常见逻辑漏洞分类:

漏洞类型 描述 典型案例
四舍五入错误 精度损失导致资金分配不均 奖励分配向某一方倾斜
边界条件 0 值处理、最大值溢出 空数组导致 panic
状态机错误 允许非法状态转换 已关闭的池仍可交易
签名验证遗漏 关键操作缺少签名验证 任意地址可提现
原子性违反 多步操作中间状态不一致 部分失败不回滚
pub fn check_state_machine(
    deps: DepsMut,
    contract_state: &ContractState,
    proposed_action: &Action,
) -> StdResult<StateCheckResult> {
    let allowed_transitions: HashMap<&str, Vec<&str>> = [
        ("initialized", vec!["active", "paused"]),
        ("active", vec!["paused", "closing"]),
        ("paused", vec!["active", "closing"]),
        ("closing", vec!["closed"]),
        ("closed", vec![]),
    ].iter().cloned().collect();

    let current_state = &contract_state.phase;
    let target_state = proposed_action.target_state();

    let valid = allowed_transitions
        .get(current_state.as_str())
        .map(|transitions| transitions.contains(&target_state.as_str()))
        .unwrap_or(false);

    if !valid {
        return Ok(StateCheckResult {
            valid: false,
            reason: format!("Invalid state transition: {} -> {}", current_state, target_state),
            severity: Severity::High,
        });
    }

    Ok(StateCheckResult {
        valid: true,
        reason: String::new(),
        severity: Severity::None,
    })
}

1.3 威胁评分模型

威胁评分 = W1 x 严重程度 + W2 x 资产价值 + W3 x 可操作性 + W4 x 紧急程度

严重程度: 0-10 (基于攻击影响)
资产价值: 0-10 (受影响的 TVL)
可操作性: 0-10 (攻击可行性和成熟度)
紧急程度: 0-10 (攻击是否正在进行)

W1 = 0.40, W2 = 0.25, W3 = 0.20, W4 = 0.15

评分分级:
  Critical:  >= 8.0
  High:      6.0 - 7.9
  Medium:    4.0 - 5.9
  Low:       2.0 - 3.9
  Info:      < 2.0
pub struct ThreatScore {
    pub severity: f64,
    pub tvl_at_risk: Uint128,
    pub exploitability: f64,
    pub urgency: f64,
}

impl ThreatScore {
    pub fn calculate(&self) -> f64 {
        let normalized_tvl = (self.tvl_at_risk.u128() as f64).log10().min(10.0) / 10.0;
        0.40 * self.severity
            + 0.25 * (normalized_tvl * 10.0)
            + 0.20 * self.exploitability
            + 0.15 * self.urgency
    }

    pub fn severity_level(&self) -> &'static str {
        let score = self.calculate();
        if score >= 8.0 { "CRITICAL" }
        else if score >= 6.0 { "HIGH" }
        else if score >= 4.0 { "MEDIUM" }
        else if score >= 2.0 { "LOW" }
        else { "INFO" }
    }
}

2. 实时交易监控架构

2.1 整体架构

MSG Chain 的实时交易监控架构分为四层:数据采集层、预处理层、分析引擎层和响应层。

┌──────────────────────────────────────────────────────────┐
│                      响应层                              │
│  ┌──────────┐  ┌──────────┐  ┌──────────┐  ┌──────────┐ │
│  │ 自动暂停  │  │ 多签触发 │  │ 紧急提案 │  │ 告警通知 │ │
│  └──────────┘  └──────────┘  └──────────┘  └──────────┘ │
├──────────────────────────────────────────────────────────┤
│                     分析引擎层                           │
│  ┌──────────┐  ┌──────────┐  ┌──────────┐  ┌──────────┐ │
│  │ 行为基线  │  │ 规则引擎  │  │ ML 检测  │  │ MEV 检测 │ │
│  └──────────┘  └──────────┘  └──────────┘  └──────────┘ │
├──────────────────────────────────────────────────────────┤
│                    预处理层                              │
│  ┌──────────┐  ┌──────────┐  ┌──────────┐  ┌──────────┐ │
│  │ Preflight│  │ 交易解码  │  │ 特征提取  │  │ 协议分类 │ │
│  └──────────┘  └──────────┘  └──────────┘  └──────────┘ │
├──────────────────────────────────────────────────────────┤
│                    数据采集层                            │
│  ┌──────────┐  ┌──────────┐  ┌──────────┐  ┌──────────┐ │
│  │ RPC 节点  │  │ Mempool  │  │ Indexer  │  │ 日志采集 │ │
│  └──────────┘  └──────────┘  └──────────┘  └──────────┘ │
└──────────────────────────────────────────────────────────┘

2.2 Mempool 监控

2.2.1 Mempool 数据采集

Mempool 是攻击检测的第一道防线。在交易进入区块之前捕获异常行为,允许抢先阻断。

import { CosmWasmClient } from '@cosmjs/cosmwasm-stargate';

interface MempoolTx {
  hash: string;
  from: string;
  to: string;
  value: string;
  gas: string;
  gasPrice: string;
  memo: string;
  msgs: TransactionMsg[];
  timestamp: number;
  rawHex: string;
}

interface MempoolMonitorConfig {
  rpcEndpoint: string;
  wsEndpoint: string;
  chainId: string;
  bech32Prefix: string;
  preflightEndpoint: string;
}

class MempoolMonitor {
  private ws: WebSocket | null = null;
  private client: CosmWasmClient | null = null;
  private config: MempoolMonitorConfig;
  private subscribers: Array<(tx: MempoolTx) => void> = [];
  private txQueue: MempoolTx[] = [];

  constructor(config: MempoolMonitorConfig) {
    this.config = config;
  }

  async connect(): Promise<void> {
    this.client = await CosmWasmClient.connect(this.config.rpcEndpoint);
    this.ws = new WebSocket(this.config.wsEndpoint);

    this.ws.on('open', () => {
      this.ws!.send(JSON.stringify({
        jsonrpc: '2.0',
        method: 'subscribe',
        params: ['tm.event="Tx"'],
        id: 1,
      }));
    });

    this.ws.on('message', (data: WebSocket.Data) => {
      const parsed = JSON.parse(data.toString());
      if (parsed.result?.events?.tx) {
        const tx = this.decodeMempoolTx(parsed.result);
        if (tx) {
          this.txQueue.push(tx);
          this.notifySubscribers(tx);
        }
      }
    });
  }

  private decodeMempoolTx(raw: any): MempoolTx | null {
    try {
      return {
        hash: raw.events['tx.hash']?.[0] || '',
        from: raw.events['tx.from']?.[0] || '',
        to: raw.events['tx.to']?.[0] || '',
        value: raw.events['tx.value']?.[0] || '0',
        gas: raw.events['tx.gas']?.[0] || '0',
        gasPrice: raw.events['tx.gas_price']?.[0] || '0',
        memo: raw.events['tx.memo']?.[0] || '',
        msgs: this.decodeMsgs(raw),
        timestamp: Date.now(),
        rawHex: raw.events['tx']?.[0] || '',
      };
    } catch (err) {
      console.error('Failed to decode mempool tx:', err);
      return null;
    }
  }

  private decodeMsgs(raw: any): TransactionMsg[] {
    const msgs: TransactionMsg[] = [];
    const msgTypes = raw.events['message.action'] || [];
    const contracts = raw.events['wasm._contract_address'] || [];
    const senders = raw.events['message.sender'] || [];

    for (let i = 0; i < msgTypes.length; i++) {
      msgs.push({
        type: msgTypes[i] || 'unknown',
        contract: contracts[i] || '',
        sender: senders[i] || '',
      });
    }

    return msgs;
  }

  subscribe(callback: (tx: MempoolTx) => void): () => void {
    this.subscribers.push(callback);
    return () => {
      this.subscribers = this.subscribers.filter(cb => cb !== callback);
    };
  }

  private notifySubscribers(tx: MempoolTx): void {
    for (const cb of this.subscribers) {
      try { cb(tx); } catch (err) {
        console.error('Subscriber error:', err);
      }
    }
  }

  getPendingTxs(): MempoolTx[] {
    const now = Date.now();
    this.txQueue = this.txQueue.filter(tx => now - tx.timestamp < 30000);
    return this.txQueue;
  }
}

2.2.2 Pre-flight Check (交易模拟)

在交易进入区块之前,通过模拟执行验证其安全性。Pre-flight check 使用与节点相同的 wasmd 运行时进行沙箱执行。

pub struct PreflightConfig {
    pub max_gas: u64,
    pub max_messages: u32,
    pub max_data_size: usize,
    pub blocked_functions: Vec<String>,
    pub suspicious_patterns: Vec<String>,
    pub allowed_contracts: Option<Vec<String>>,
}

pub struct PreflightResult {
    pub passed: bool,
    pub gas_estimate: u64,
    pub warnings: Vec<String>,
    pub errors: Vec<String>,
    pub state_diffs: Vec<StateDiff>,
    pub affected_contracts: Vec<String>,
    pub risk_score: f64,
}

pub struct PreflightChecker {
    config: PreflightConfig,
}

impl PreflightChecker {
    pub fn new(config: PreflightConfig) -> Self {
        Self { config }
    }

    pub fn check_tx(&self, tx: &AnalyzedTransaction) -> PreflightResult {
        let mut warnings = Vec::new();
        let mut errors = Vec::new();

        if tx.messages.len() > self.config.max_messages as usize {
            errors.push(format!("Transaction has {} messages, limit is {}",
                tx.messages.len(), self.config.max_messages));
        }

        for msg in &tx.messages {
            if let Some(func) = &msg.function_name {
                if self.config.blocked_functions.contains(func) {
                    errors.push(format!("Blocked function call: {}", func));
                }
            }
        }

        for pattern in &self.config.suspicious_patterns {
            if tx.memo.contains(pattern) || tx.raw_hex.contains(pattern) {
                warnings.push(format!("Suspicious pattern: {}", pattern));
            }
        }

        let gas_estimate = self.estimate_gas(tx);
        let affected_contracts: Vec<String> = tx.messages.iter()
            .filter_map(|m| m.contract_address.clone())
            .collect();

        let risk_score = self.calculate_risk(tx, &warnings, &errors, gas_estimate);

        PreflightResult {
            passed: errors.is_empty() && risk_score < 0.8,
            gas_estimate,
            warnings,
            errors,
            state_diffs: Vec::new(),
            affected_contracts,
            risk_score,
        }
    }

    fn estimate_gas(&self, tx: &AnalyzedTransaction) -> u64 {
        let base_gas: u64 = 100_000;
        let msg_gas: u64 = tx.messages.len() as u64 * 50_000;
        let data_gas: u64 = (tx.raw_hex.len() / 1024) as u64 * 10_000;
        base_gas + msg_gas + data_gas
    }

    fn calculate_risk(&self, tx: &AnalyzedTransaction, warnings: &[String],
        errors: &[String], gas_estimate: u64) -> f64 {
        let mut score = 0.0;
        if !errors.is_empty() { score += 0.5; }
        score += (warnings.len() as f64).min(5.0) * 0.05;
        if gas_estimate > self.config.max_gas { score += 0.2; }
        score += (tx.messages.len() as f64).min(10.0) * 0.03;
        let total_value: u128 = tx.messages.iter()
            .filter_map(|m| m.amount)
            .map(|c| c.amount.u128())
            .sum();
        if total_value > 1_000_000_000_000u128 { score += 0.2; }
        score.min(1.0)
    }
}

2.3 异常交易检测

2.3.1 特征工程

pub struct TxFeatureVector {
    pub msg_count: u32,
    pub contract_call_count: u32,
    pub unique_contracts: u32,
    pub total_gas: u64,
    pub gas_price: u64,
    pub memo_length: usize,
    pub delegate_call_count: u32,
    pub reply_count: u32,
    pub submsg_depth: u32,
    pub cross_contract_calls: u32,
    pub total_inflow: Uint128,
    pub total_outflow: Uint128,
    pub max_single_transfer: Uint128,
    pub unique_denoms: u32,
    pub sender_age_blocks: u64,
    pub sender_tx_count_24h: u32,
    pub time_since_last_tx: u64,
    pub has_selfdestruct: bool,
    pub has_sudo_call: bool,
    pub has_migrate_call: bool,
    pub has_update_owner: bool,
    pub has_proxy_upgrade: bool,
    pub unusual_function_names: u32,
    pub function_calls: HashMap<String, u32>,
}

pub struct FeatureExtractor;

impl FeatureExtractor {
    pub fn extract(tx: &AnalyzedTransaction) -> TxFeatureVector {
        let msg_count = tx.messages.len() as u32;
        let contract_call_count = tx.messages.iter()
            .filter(|m| m.is_contract_call()).count() as u32;
        let unique_contracts = tx.messages.iter()
            .filter_map(|m| m.contract_address.as_ref())
            .collect::<std::collections::HashSet<_>>()
            .len() as u32;
        let submsg_depth = tx.messages.iter()
            .map(|m| m.submsg_depth).max().unwrap_or(0);
        let cross_contract_calls = tx.messages.iter()
            .filter(|m| m.is_cross_contract_call()).count() as u32;

        let total_inflow = tx.messages.iter()
            .filter(|m| m.direction == TransferDirection::In)
            .map(|m| m.amount.unwrap_or_default())
            .fold(Uint128::zero(), |a, b| a + b);

        let total_outflow = tx.messages.iter()
            .filter(|m| m.direction == TransferDirection::Out)
            .map(|m| m.amount.unwrap_or_default())
            .fold(Uint128::zero(), |a, b| a + b);

        let mut func_calls = HashMap::new();
        for msg in &tx.messages {
            if let Some(ref func) = msg.function_name {
                *func_calls.entry(func.clone()).or_insert(0) += 1;
            }
        }

        TxFeatureVector {
            msg_count,
            contract_call_count,
            unique_contracts,
            total_gas: tx.gas_used,
            gas_price: tx.gas_price,
            memo_length: tx.memo.len(),
            delegate_call_count: 0,
            reply_count: tx.reply_count,
            submsg_depth,
            cross_contract_calls,
            total_inflow,
            total_outflow,
            max_single_transfer: tx.messages.iter()
                .filter_map(|m| m.amount)
                .max().unwrap_or(Uint128::zero()),
            unique_denoms: tx.messages.iter()
                .filter_map(|m| m.denom.as_ref())
                .collect::<std::collections::HashSet<_>>()
                .len() as u32,
            sender_age_blocks: 0,
            sender_tx_count_24h: 0,
            time_since_last_tx: 0,
            has_selfdestruct: tx.messages.iter()
                .any(|m| m.function_name.as_deref() == Some("self_destruct")),
            has_sudo_call: tx.messages.iter()
                .any(|m| m.function_name.as_deref() == Some("sudo")),
            has_migrate_call: tx.messages.iter()
                .any(|m| m.function_name.as_deref() == Some("migrate")),
            has_update_owner: tx.messages.iter()
                .any(|m| m.function_name.as_deref() == Some("update_owner")),
            has_proxy_upgrade: tx.messages.iter()
                .any(|m| m.function_name.as_deref() == Some("upgrade")),
            unusual_function_names: tx.messages.iter()
                .filter(|m| m.function_name.as_deref()
                    .map(|f| !is_standard_function(f)).unwrap_or(false))
                .count() as u32,
            function_calls: func_calls,
        }
    }
}

fn is_standard_function(name: &str) -> bool {
    const STANDARD: &[&str] = &[
        "transfer", "transfer_from", "approve", "burn", "mint",
        "deposit", "withdraw", "swap", "provide_liquidity",
        "withdraw_liquidity", "stake", "unstake", "claim",
        "execute", "instantiate", "send",
    ];
    STANDARD.contains(&name)
}

2.3.2 异常检测算法 — Isolation Forest

pub struct IsolationForest {
    pub trees: Vec<IsolationTreeNode>,
    pub sample_size: usize,
    pub num_trees: usize,
}

impl IsolationForest {
    pub fn new(num_trees: usize, sample_size: usize) -> Self {
        Self { trees: Vec::new(), sample_size, num_trees }
    }

    pub fn train(&mut self, data: &[TxFeatureVector]) {
        let feature_matrix: Vec<Vec<f64>> = data.iter()
            .map(|f| self.vectorize(f)).collect();

        for _ in 0..self.num_trees {
            let sample = self.sample(&feature_matrix);
            let tree = self.build_tree(&sample, 0, 100);
            self.trees.push(tree);
        }
    }

    fn vectorize(&self, fv: &TxFeatureVector) -> Vec<f64> {
        vec![
            fv.msg_count as f64,
            fv.contract_call_count as f64,
            fv.unique_contracts as f64,
            fv.total_gas as f64,
            fv.gas_price as f64,
            fv.memo_length as f64,
            fv.submsg_depth as f64,
            fv.cross_contract_calls as f64,
            (fv.total_inflow.u128() as f64).log10().max(0.0),
            (fv.total_outflow.u128() as f64).log10().max(0.0),
            fv.unique_denoms as f64,
            fv.sender_tx_count_24h as f64,
            fv.has_selfdestruct as u32 as f64,
            fv.has_sudo_call as u32 as f64,
            fv.has_migrate_call as u32 as f64,
            fv.has_update_owner as u32 as f64,
            fv.has_proxy_upgrade as u32 as f64,
        ]
    }

    fn sample(&self, data: &[Vec<f64>]) -> Vec<Vec<f64>> {
        use rand::seq::SliceRandom;
        let mut rng = rand::thread_rng();
        let n = self.sample_size.min(data.len());
        data.choose_multiple(&mut rng, n).cloned().collect()
    }

    fn build_tree(&self, data: &[Vec<f64>], depth: usize, max_depth: usize) -> IsolationTreeNode {
        if depth >= max_depth || data.len() <= 1 {
            return IsolationTreeNode {
                split_feature: 0, split_value: 0.0,
                left: None, right: None,
                size: data.len(), is_leaf: true,
            };
        }

        let num_features = data[0].len();
        let feature = rand::random::<usize>() % num_features;
        let min_val = data.iter().map(|d| d[feature]).fold(f64::INFINITY, f64::min);
        let max_val = data.iter().map(|d| d[feature]).fold(f64::NEG_INFINITY, f64::max);

        if (max_val - min_val).abs() < f64::EPSILON {
            return IsolationTreeNode {
                split_feature: feature, split_value: min_val,
                left: None, right: None,
                size: data.len(), is_leaf: true,
            };
        }

        let split_value = min_val + rand::random::<f64>() * (max_val - min_val);
        let left_data: Vec<Vec<f64>> = data.iter()
            .filter(|d| d[feature] < split_value).cloned().collect();
        let right_data: Vec<Vec<f64>> = data.iter()
            .filter(|d| d[feature] >= split_value).cloned().collect();

        IsolationTreeNode {
            split_feature: feature, split_value,
            left: Some(Box::new(self.build_tree(&left_data, depth + 1, max_depth))),
            right: Some(Box::new(self.build_tree(&right_data, depth + 1, max_depth))),
            size: data.len(), is_leaf: false,
        }
    }

    pub fn anomaly_score(&self, fv: &TxFeatureVector) -> f64 {
        let features = self.vectorize(fv);
        let avg_path_length: f64 = self.trees.iter()
            .map(|tree| self.path_length(tree, &features, 0))
            .sum::<f64>() / self.trees.len() as f64;
        let c = self.c_factor(self.sample_size);
        2.0f64.powf(-avg_path_length / c)
    }

    fn path_length(&self, node: &IsolationTreeNode, features: &[f64], depth: usize) -> f64 {
        if node.is_leaf { return depth as f64 + self.c_factor(node.size); }
        if features[node.split_feature] < node.split_value {
            if let Some(ref left) = node.left { return self.path_length(left, features, depth + 1); }
        } else {
            if let Some(ref right) = node.right { return self.path_length(right, features, depth + 1); }
        }
        depth as f64 + self.c_factor(node.size)
    }

    fn c_factor(&self, n: usize) -> f64 {
        if n <= 1 { return 0.0; }
        if n == 2 { return 1.0; }
        let h = (n as f64).ln() + std::f64::consts::E - 1.0;
        2.0 * h - (2.0 * (n - 1) as f64 / n as f64)
    }

    pub fn is_anomaly(&self, fv: &TxFeatureVector, threshold: f64) -> bool {
        self.anomaly_score(fv) > threshold
    }
}

pub struct IsolationTreeNode {
    pub split_feature: usize,
    pub split_value: f64,
    pub left: Option<Box<IsolationTreeNode>>,
    pub right: Option<Box<IsolationTreeNode>>,
    pub size: usize,
    pub is_leaf: bool,
}

2.3.3 滑动窗口统计检测

pub struct SlidingWindowDetector {
    pub window_size: usize,
    pub history: VecDeque<TxFeatureVector>,
    pub mean: Option<StatsSummary>,
    pub std: Option<StatsSummary>,
}

pub struct StatsSummary {
    pub mean_msg_count: f64,
    pub mean_gas: f64,
    pub mean_inflow: f64,
    pub mean_outflow: f64,
    pub std_msg_count: f64,
    pub std_gas: f64,
    pub std_inflow: f64,
    pub std_outflow: f64,
}

impl SlidingWindowDetector {
    pub fn new(window_size: usize) -> Self {
        Self {
            window_size,
            history: VecDeque::with_capacity(window_size),
            mean: None,
            std: None,
        }
    }

    pub fn add_sample(&mut self, fv: TxFeatureVector) {
        if self.history.len() >= self.window_size {
            self.history.pop_front();
        }
        self.history.push_back(fv);
        self.update_statistics();
    }

    fn update_statistics(&mut self) {
        if self.history.is_empty() { return; }
        let n = self.history.len() as f64;

        let mean_msg = self.history.iter().map(|f| f.msg_count as f64).sum::<f64>() / n;
        let mean_gas = self.history.iter().map(|f| f.total_gas as f64).sum::<f64>() / n;
        let mean_inflow = self.history.iter()
            .map(|f| (f.total_inflow.u128() as f64).log10().max(0.0)).sum::<f64>() / n;
        let mean_outflow = self.history.iter()
            .map(|f| (f.total_outflow.u128() as f64).log10().max(0.0)).sum::<f64>() / n;

        let var_msg = self.history.iter()
            .map(|f| (f.msg_count as f64 - mean_msg).powi(2)).sum::<f64>() / n;
        let var_gas = self.history.iter()
            .map(|f| (f.total_gas as f64 - mean_gas).powi(2)).sum::<f64>() / n;
        let var_inflow = self.history.iter()
            .map(|f| ((f.total_inflow.u128() as f64).log10().max(0.0) - mean_inflow).powi(2))
            .sum::<f64>() / n;
        let var_outflow = self.history.iter()
            .map(|f| ((f.total_outflow.u128() as f64).log10().max(0.0) - mean_outflow).powi(2))
            .sum::<f64>() / n;

        self.mean = Some(StatsSummary {
            mean_msg_count: mean_msg, mean_gas, mean_inflow, mean_outflow,
            std_msg_count: 0.0, std_gas: 0.0, std_inflow: 0.0, std_outflow: 0.0,
        });
        self.std = Some(StatsSummary {
            mean_msg_count: var_msg.sqrt(), mean_gas: var_gas.sqrt(),
            mean_inflow: var_inflow.sqrt(), mean_outflow: var_outflow.sqrt(),
            std_msg_count: 0.0, std_gas: 0.0, std_inflow: 0.0, std_outflow: 0.0,
        });
    }

    pub fn detect_anomaly(&self, fv: &TxFeatureVector, sigma: f64) -> AnomalyResult {
        let mut signals = Vec::new();
        if let (Some(ref mean), Some(ref std)) = (self.mean.as_ref(), self.std.as_ref()) {
            let msg_z = (fv.msg_count as f64 - mean.mean_msg_count)
                / (std.mean_msg_count).max(1.0);
            if msg_z.abs() > sigma {
                signals.push(format!("msg_count z-score: {:.2}", msg_z));
            }
            let gas_z = (fv.total_gas as f64 - mean.mean_gas)
                / (std.mean_gas).max(1.0);
            if gas_z.abs() > sigma {
                signals.push(format!("gas z-score: {:.2}", gas_z));
            }
        }
        AnomalyResult {
            is_anomaly: !signals.is_empty(),
            signals,
            severity: if signals.len() >= 3 { Severity::High }
            else if !signals.is_empty() { Severity::Medium }
            else { Severity::None },
        }
    }
}

2.4 Indexer 事件流架构

interface IndexerConfig {
  rpcEndpoint: string;
  chainId: string;
  startHeight: number;
  pollIntervalMs: number;
  batchSize: number;
}

interface BlockEvent {
  height: number;
  hash: string;
  time: string;
  txs: TxEvent[];
  proposer: string;
  numTxs: number;
  gasUsed: string;
  gasLimit: string;
}

interface TxEvent {
  hash: string;
  height: number;
  index: number;
  sender: string;
  gasUsed: string;
  gasWanted: string;
  messages: WasmMsg[];
  events: ChainEvent[];
  success: boolean;
  error?: string;
}

interface WasmMsg {
  type: string;
  contract: string;
  sender: string;
  function: string;
  args: Record<string, unknown>;
  funds: { denom: string; amount: string }[];
}

interface ChainEvent {
  type: string;
  attributes: { key: string; value: string }[];
}

class SecurityIndexer {
  private client: CosmWasmClient | null = null;
  private config: IndexerConfig;
  private currentHeight: number;
  private polling: boolean = false;
  private processedBlocks: Set<number> = new Set();
  private listeners: Map<string, Function[]> = new Map();

  constructor(config: IndexerConfig) {
    this.config = config;
    this.currentHeight = config.startHeight;
  }

  async start(): Promise<void> {
    this.client = await CosmWasmClient.connect(this.config.rpcEndpoint);
    this.polling = true;
    this.poll();
  }

  stop(): void { this.polling = false; }

  private async poll(): Promise<void> {
    while (this.polling) {
      try {
        const latestHeight = await this.client!.getHeight();
        while (this.currentHeight < latestHeight) {
          const endHeight = Math.min(
            this.currentHeight + this.config.batchSize, latestHeight);
          const blocks = await this.fetchBlocks(this.currentHeight, endHeight);
          for (const block of blocks) {
            await this.processBlock(block);
            this.processedBlocks.add(block.height);
          }
          this.currentHeight = endHeight + 1;
        }
      } catch (err) {
        console.error('Indexer poll error:', err);
      }
      await new Promise(resolve => setTimeout(resolve, this.config.pollIntervalMs));
    }
  }

  private async fetchBlocks(from: number, to: number): Promise<BlockEvent[]> {
    const blocks: BlockEvent[] = [];
    for (let h = from; h <= to; h++) {
      try {
        const block = await this.client!.getBlock(h);
        const txs = await this.getBlockTxs(h);
        blocks.push({
          height: h, hash: block.id, time: block.header.time,
          txs, proposer: block.header.proposerAddress,
          numTxs: txs.length, gasUsed: '0', gasLimit: '0',
        });
      } catch (err) { console.error(`Failed to fetch block ${h}:`, err); }
    }
    return blocks;
  }

  private async getBlockTxs(height: number): Promise<TxEvent[]> {
    const txs: TxEvent[] = [];
    const result = await this.client!.searchTx([
      { key: 'tx.height', value: height.toString() },
    ]);
    for (const tx of result) {
      const wasmMsgs = this.extractWasmMsgs(tx);
      txs.push({
        hash: tx.hash, height: tx.height, index: tx.index,
        sender: this.extractSender(tx),
        gasUsed: tx.gasUsed?.toString() || '0',
        gasWanted: tx.gasWanted?.toString() || '0',
        messages: wasmMsgs,
        events: tx.events.map(e => ({
          type: e.type,
          attributes: e.attributes.map(a => ({ key: a.key, value: a.value })),
        })),
        success: tx.code === 0,
        error: tx.rawLog,
      });
    }
    return txs;
  }

  private extractWasmMsgs(tx: any): WasmMsg[] {
    const msgs: WasmMsg[] = [];
    const events = tx.events || [];
    for (const event of events) {
      if (event.type.startsWith('wasm-') || event.type.startsWith('wasm.')) {
        const attrs = this.eventToMap(event.attributes);
        msgs.push({
          type: attrs['action'] || event.type,
          contract: attrs['_contract_address'] || '',
          sender: attrs['sender'] || '',
          function: attrs['action'] || '',
          args: attrs,
          funds: [],
        });
      }
    }
    return msgs;
  }

  private extractSender(tx: any): string {
    const events = tx.events || [];
    for (const event of events) {
      if (event.type === 'message') {
        const senderAttr = event.attributes.find((a: any) => a.key === 'sender');
        if (senderAttr) return senderAttr.value;
      }
    }
    return '';
  }

  private eventToMap(attrs: { key: string; value: string }[]): Record<string, string> {
    const map: Record<string, string> = {};
    for (const attr of attrs) { map[attr.key] = attr.value; }
    return map;
  }

  private async processBlock(block: BlockEvent): Promise<void> {
    this.emit('block', block);
    for (const tx of block.txs) {
      this.emit('tx', tx);
      for (const msg of tx.messages) {
        this.emit('wasm_msg', { ...msg, txHash: tx.hash, blockHeight: block.height });
      }
    }
  }

  on(event: string, callback: Function): void {
    if (!this.listeners.has(event)) this.listeners.set(event, []);
    this.listeners.get(event)!.push(callback);
  }

  private emit(event: string, data: any): void {
    const cbs = this.listeners.get(event) || [];
    for (const cb of cbs) { try { cb(data); } catch {} }
  }

  getProcessedHeight(): number { return this.currentHeight - 1; }
  getBlockCount(): number { return this.processedBlocks.size; }
}

3. 合约行为基线

3.1 基线建立方法

合约行为基线是通过对历史链上数据的统计分析,建立合约正常行为模式的量化描述。

3.1.1 基线数据采集

interface ContractBaseline {
  contractAddress: string;
  codeId: number;
  label: string;
  instantiateHeight: number;
  dailyCallCount: StatisticalSummary;
  hourlyCallDistribution: number[];
  gasPerCall: StatisticalSummary;
  gasByFunction: Record<string, StatisticalSummary>;
  dailyVolume: StatisticalSummary;
  avgTransferSize: StatisticalSummary;
  transferFrequency: StatisticalSummary;
  uniqueCallersDaily: StatisticalSummary;
  topCallers: { address: string; count: number; percentage: number }[];
  functionDistribution: Record<string, number>;
  unusualFunctions: string[];
  peakHours: number[];
  quietHours: number[];
  weekendMultiplier: number;
  crossContractCount: StatisticalSummary;
  trustedContracts: string[];
  avgStateWritesPerCall: StatisticalSummary;
  avgStateReadsPerCall: StatisticalSummary;
}

interface StatisticalSummary {
  mean: number;
  median: number;
  stddev: number;
  p95: number;
  p99: number;
  min: number;
  max: number;
  sampleCount: number;
}

class BaselineCollector {
  private baselines: Map<string, ContractBaseline> = new Map();

  async collectBaseline(
    indexer: SecurityIndexer,
    contractAddress: string,
    lookbackDays: number,
  ): Promise<ContractBaseline> {
    const endHeight = await this.getCurrentHeight();
    const blocksPerDay = 14400;
    const startHeight = endHeight - lookbackDays * blocksPerDay;
    const txs = await this.fetchContractTxs(indexer, contractAddress, startHeight, endHeight);
    const baseline = this.computeBaseline(contractAddress, txs);
    this.baselines.set(contractAddress, baseline);
    return baseline;
  }

  private computeBaseline(address: string, txs: TxEvent[]): ContractBaseline {
    const callCounts: number[] = [];
    const gasValues: number[] = [];
    const transferValues: number[] = [];
    const dailyCallers: Map<number, Set<string>> = new Map();
    const functionCalls: Record<string, number> = {};

    for (const tx of txs) {
      const day = Math.floor(tx.height / 14400);
      callCounts.push(tx.messages.length);
      gasValues.push(parseInt(tx.gasUsed));
      if (!dailyCallers.has(day)) dailyCallers.set(day, new Set());
      dailyCallers.get(day)!.add(tx.sender);
      for (const msg of tx.messages) {
        if (msg.function) {
          functionCalls[msg.function] = (functionCalls[msg.function] || 0) + 1;
        }
      }
    }

    return {
      contractAddress: address, codeId: 0, label: '', instantiateHeight: 0,
      dailyCallCount: this.summarize(callCounts),
      hourlyCallDistribution: new Array(24).fill(0),
      gasPerCall: this.summarize(gasValues), gasByFunction: {},
      dailyVolume: this.summarize([]),
      avgTransferSize: this.summarize(transferValues),
      transferFrequency: this.summarize([]),
      uniqueCallersDaily: this.summarize(
        Array.from(dailyCallers.values()).map(s => s.size)),
      topCallers: [], functionDistribution: functionCalls,
      unusualFunctions: [], peakHours: [], quietHours: [],
      weekendMultiplier: 1.0, crossContractCount: this.summarize([]),
      trustedContracts: [], avgStateWritesPerCall: this.summarize([]),
      avgStateReadsPerCall: this.summarize([]),
    };
  }

  private summarize(values: number[]): StatisticalSummary {
    if (values.length === 0) {
      return { mean: 0, median: 0, stddev: 0, p95: 0, p99: 0, min: 0, max: 0, sampleCount: 0 };
    }
    const sorted = [...values].sort((a, b) => a - b);
    const n = sorted.length;
    const mean = sorted.reduce((a, b) => a + b, 0) / n;
    const variance = sorted.reduce((sum, v) => sum + (v - mean) ** 2, 0) / n;
    return {
      mean, median: sorted[Math.floor(n / 2)], stddev: Math.sqrt(variance),
      p95: sorted[Math.floor(n * 0.95)], p99: sorted[Math.floor(n * 0.99)],
      min: sorted[0], max: sorted[n - 1], sampleCount: n,
    };
  }

  private async fetchContractTxs(
    indexer: SecurityIndexer, contractAddress: string,
    startHeight: number, endHeight: number,
  ): Promise<TxEvent[]> { return []; }

  private async getCurrentHeight(): Promise<number> { return 0; }

  getBaseline(address: string): ContractBaseline | undefined {
    return this.baselines.get(address);
  }
}

3.2 异常检测算法

3.2.1 统计检测

pub enum AnomalySignal { None, Medium(f64), High(f64) }

pub struct StatisticalAnomalyDetector;

impl StatisticalAnomalyDetector {
    pub fn zscore_detection(value: f64, mean: f64, stddev: f64, threshold: f64) -> AnomalySignal {
        if stddev < f64::EPSILON { return AnomalySignal::None; }
        let z = (value - mean).abs() / stddev;
        if z > threshold { AnomalySignal::High(z) }
        else if z > threshold * 0.7 { AnomalySignal::Medium(z) }
        else { AnomalySignal::None }
    }

    pub fn moving_average_deviation(current: f64, ma_short: f64, ma_long: f64, threshold: f64) -> AnomalySignal {
        if ma_long < f64::EPSILON { return AnomalySignal::None; }
        let deviation = (ma_short - ma_long).abs() / ma_long;
        if deviation > threshold { AnomalySignal::High(deviation) }
        else if deviation > threshold * 0.5 { AnomalySignal::Medium(deviation) }
        else { AnomalySignal::None }
    }
}

3.2.2 时序模式异常检测

pub struct TimeSeriesDetector {
    pub period_blocks: u64,
    pub history: VecDeque<Vec<f64>>,
    pub current_period: Vec<f64>,
    pub period_start: u64,
}

impl TimeSeriesDetector {
    pub fn new(period_blocks: u64, max_history: usize) -> Self {
        Self {
            period_blocks,
            history: VecDeque::with_capacity(max_history),
            current_period: Vec::new(),
            period_start: 0,
        }
    }

    pub fn add_observation(&mut self, block_height: u64, value: f64) -> Option<AnomalyResult> {
        if self.period_start == 0 { self.period_start = block_height; }
        self.current_period.push(value);
        if block_height - self.period_start >= self.period_blocks {
            let completed_period = std::mem::take(&mut self.current_period);
            self.period_start = block_height;
            let period_stat = self.compute_period_stats(&completed_period);
            self.history.push_back(period_stat.clone());
            if self.history.len() >= 3 { return self.detect_anomaly(&period_stat); }
        }
        None
    }

    fn compute_period_stats(&self, data: &[f64]) -> Vec<f64> {
        if data.is_empty() { return vec![0.0; 5]; }
        let n = data.len() as f64;
        let sum: f64 = data.iter().sum();
        let mean = sum / n;
        let variance = data.iter().map(|v| (v - mean).powi(2)).sum::<f64>() / n;
        let max = data.iter().cloned().fold(f64::NEG_INFINITY, f64::max);
        let min = data.iter().cloned().fold(f64::INFINITY, f64::min);
        vec![mean, variance.sqrt(), max, min, n]
    }

    fn detect_anomaly(&self, current: &[f64]) -> AnomalyResult {
        let mut signals = Vec::new();
        let hist_mean: Vec<f64> = (0..5)
            .map(|i| self.history.iter().map(|h| h[i]).sum::<f64>() / self.history.len() as f64)
            .collect();
        let hist_std: Vec<f64> = (0..5).map(|i| {
            let m = hist_mean[i];
            let var = self.history.iter().map(|h| (h[i] - m).powi(2)).sum::<f64>() / self.history.len() as f64;
            var.sqrt()
        }).collect();

        for (i, (c, (m, s))) in current.iter().zip(hist_mean.iter().zip(hist_std.iter())).enumerate() {
            if *s > f64::EPSILON {
                let z = (c - m).abs() / s;
                if z > 3.0 {
                    let labels = ["mean", "stddev", "max", "min", "count"];
                    signals.push(format!("Period {} z-score: {:.2}", labels[i], z));
                }
            }
        }

        AnomalyResult {
            is_anomaly: !signals.is_empty(),
            signals,
            severity: if signals.len() >= 2 { Severity::High }
            else if !signals.is_empty() { Severity::Medium }
            else { Severity::None },
        }
    }
}

pub struct AnomalyResult {
    pub is_anomaly: bool,
    pub signals: Vec<String>,
    pub severity: Severity,
}

3.3 基线偏差告警

pub struct BaselineDeviationAlert {
    pub contract: String, pub metric: String,
    pub expected: f64, pub actual: f64,
    pub deviation: f64, pub zscore: f64,
    pub severity: Severity, pub timestamp: u64, pub block_height: u64,
}

pub struct BaselineMonitor {
    pub baselines: HashMap<String, ContractBaseline>,
    pub deviation_threshold: f64,
}

impl BaselineMonitor {
    pub fn new(deviation_threshold: f64) -> Self {
        Self { baselines: HashMap::new(), deviation_threshold }
    }

    pub fn check_baseline(&self, contract: &str, current_stats: &TxFeatureVector,
        block_height: u64) -> Vec<BaselineDeviationAlert> {
        let mut alerts = Vec::new();
        let baseline = match self.baselines.get(contract) { Some(b) => b, None => return alerts };

        if baseline.gas_per_call.sample_count > 10 {
            let z = (current_stats.total_gas as f64 - baseline.gas_per_call.mean)
                / baseline.gas_per_call.stddev.max(1.0);
            if z.abs() > 3.0 {
                alerts.push(BaselineDeviationAlert {
                    contract: contract.to_string(), metric: "gas_per_call".to_string(),
                    expected: baseline.gas_per_call.mean,
                    actual: current_stats.total_gas as f64,
                    deviation: z.abs(), zscore: z,
                    severity: if z.abs() > 5.0 { Severity::Critical } else { Severity::High },
                    timestamp: std::time::SystemTime::now().duration_since(std::time::UNIX_EPOCH).unwrap().as_secs(),
                    block_height,
                });
            }
        }

        for (func, count) in &current_stats.function_calls {
            let expected = baseline.function_distribution.get(func).copied().unwrap_or(0);
            if expected == 0 && *count > 0 {
                alerts.push(BaselineDeviationAlert {
                    contract: contract.to_string(),
                    metric: format!("function_call:{}", func),
                    expected: 0.0, actual: *count as f64,
                    deviation: f64::INFINITY, zscore: f64::INFINITY,
                    severity: Severity::High, timestamp: 0, block_height,
                });
            }
        }

        alerts
    }
}

4. 告警规则引擎

4.1 规则定义

rules:
  - id: large_transfer_critical
    name: 大额转账
    description: 单笔转账超过 100 万 amsg
    severity: CRITICAL
    condition:
      type: threshold
      metric: transfer_amount
      operator: ">"
      value: 1000000000000
      window: immediate
    actions:
      - type: slack
        channel: "#security-critical"
      - type: pause_contract
      - type: notify_validator

  - id: frequent_instantiate
    name: 频繁合约部署
    description: 同一地址 1 小时内部署超过 10 个合约
    severity: HIGH
    condition:
      type: rate
      metric: instantiate_count
      operator: ">"
      value: 10
      window_seconds: 3600
      group_by: sender
    actions:
      - type: block_address
        duration_seconds: 86400
      - type: notify_security_team

  - id: ownership_change
    name: 合约所有权变更
    description: 合约 owner 地址被修改
    severity: CRITICAL
    condition:
      type: event
      event_type: wasm
      attributes:
        action: ["update_owner", "transfer_ownership", "set_admin"]
    actions:
      - type: slack
        channel: "#security-critical"
      - type: multisig_approval_required
      - type: snapshot_state

  - id: proxy_upgrade
    name: 代理合约升级
    description: 检测代理合约 upgrade / migrate 操作
    severity: CRITICAL
    condition:
      type: event
      event_type: wasm
      attributes:
        action: ["upgrade", "migrate", "update_code"]
    actions:
      - type: slack
        channel: "#security-critical"
      - type: hold_transactions
      - type: require_multi_sig

  - id: gas_anomaly
    name: Gas 消耗异常
    description: 合约 Gas 消耗偏离基线超过 5 个标准差
    severity: HIGH
    condition:
      type: statistical
      metric: gas_used
      deviation: 5.0
      window: 100
    actions:
      - type: slack
        channel: "#security-anomaly"
      - type: alert_developer

  - id: call_frequency_critical
    name: 超高频率调用
    description: 同一合约 1 分钟内被调用超过 1000 次
    severity: HIGH
    condition:
      type: rate
      metric: contract_calls
      operator: ">"
      value: 1000
      window_seconds: 60
      group_by: contract
    actions:
      - type: slack
        channel: "#security-critical"
      - type: pause_contract

  - id: sandwich_attack_pattern
    name: 三明治攻击模式
    description: 检测 Mempool 中三明治攻击交易模式
    severity: HIGH
    condition:
      type: pattern
      pattern: sandwich_attack
      window_seconds: 30
    actions:
      - type: slack
        channel: "#security-mev"
      - type: log_tx_details

  - id: agent_constitution_violation
    name: Agent 宪法违规
    description: Agent 尝试执行宪法禁止的操作
    severity: CRITICAL
    condition:
      type: event
      event_type: wasm
      attributes:
        error_code: ["12", "CONSTITUTION_VIOLATION"]
    actions:
      - type: slack
        channel: "#security-agent"
      - type: pause_agent
      - type: notify_agent_owner

  - id: contract_selfdestruct
    name: 合约自毁
    description: 合约执行 self-destruct 操作
    severity: CRITICAL
    condition:
      type: event
      event_type: wasm
      attributes:
        action: ["self_destruct", "destroy"]
    actions:
      - type: slack
        channel: "#security-critical"
      - type: snapshot_final_state

  - id: contract_paused
    name: 合约暂停
    description: 合约进入暂停状态
    severity: HIGH
    condition:
      type: event
      event_type: wasm
      attributes:
        action: ["emergency_pause", "pause"]
    actions:
      - type: slack
        channel: "#security-operations"
      - type: investigate

4.2 规则引擎实现

type AlertSeverity = 'CRITICAL' | 'HIGH' | 'MEDIUM' | 'LOW' | 'INFO';

interface AlertRule {
  id: string;
  name: string;
  description: string;
  severity: AlertSeverity;
  condition: RuleCondition;
  actions: RuleAction[];
  enabled: boolean;
  cooldownSeconds: number;
  lastTriggered?: number;
}

type RuleCondition =
  | ThresholdCondition
  | RateCondition
  | EventCondition
  | StatisticalCondition
  | PatternCondition;

interface ThresholdCondition {
  type: 'threshold';
  metric: string;
  operator: '>' | '<' | '>=' | '<=' | '==' | '!=';
  value: number;
  window: 'immediate' | '1m' | '5m' | '1h';
}

interface RateCondition {
  type: 'rate';
  metric: string;
  operator: '>' | '<';
  value: number;
  window_seconds: number;
  group_by?: string;
}

interface EventCondition {
  type: 'event';
  event_type: string;
  attributes: Record<string, string | string[]>;
}

interface StatisticalCondition {
  type: 'statistical';
  metric: string;
  deviation: number;
  window: 'immediate' | number;
}

interface PatternCondition {
  type: 'pattern';
  pattern: string;
  window_seconds: number;
}

interface RuleAction {
  type: string;
  [key: string]: unknown;
}

interface AlertEvent {
  ruleId: string;
  ruleName: string;
  severity: AlertSeverity;
  message: string;
  timestamp: number;
  context: Record<string, unknown>;
  txHash?: string;
  blockHeight?: number;
  contractAddress?: string;
}

class RuleEngine {
  private rules: Map<string, AlertRule> = new Map();
  private metricStore: Map<string, number[]> = new Map();
  private triggeredRules: Map<string, number> = new Map();
  private listeners: Map<string, Function[]> = new Map();

  loadRules(rules: AlertRule[]): void {
    for (const rule of rules) this.rules.set(rule.id, rule);
  }

  addRule(rule: AlertRule): void { this.rules.set(rule.id, rule); }
  removeRule(ruleId: string): boolean { return this.rules.delete(ruleId); }

  evaluate(metric: string, value: number, context: Record<string, unknown> = {}): AlertEvent[] {
    const alerts: AlertEvent[] = [];
    const now = Date.now();

    if (!this.metricStore.has(metric)) this.metricStore.set(metric, []);
    const values = this.metricStore.get(metric)!;
    values.push(value);
    if (values.length > 10000) values.shift();

    for (const rule of this.rules.values()) {
      if (!rule.enabled) continue;
      const lastTriggered = this.triggeredRules.get(rule.id) || 0;
      if (now - lastTriggered < rule.cooldownSeconds * 1000) continue;

      const matched = this.evaluateCondition(rule.condition, metric, value, values);
      if (matched) {
        const alert: AlertEvent = {
          ruleId: rule.id, ruleName: rule.name, severity: rule.severity,
          message: rule.description, timestamp: now,
          context: { ...context, metric, value },
        };
        alerts.push(alert);
        this.triggeredRules.set(rule.id, now);
        this.emit('alert', alert);
      }
    }
    return alerts;
  }

  private evaluateCondition(condition: RuleCondition, metric: string, value: number, values: number[]): boolean {
    switch (condition.type) {
      case 'threshold': return this.evaluateThreshold(condition as ThresholdCondition, value);
      case 'rate': return this.evaluateRate(condition as RateCondition, values);
      case 'event': return true;
      case 'statistical': return this.evaluateStatistical(condition as StatisticalCondition, value, values);
      case 'pattern': return false;
    }
  }

  private evaluateThreshold(condition: ThresholdCondition, value: number): boolean {
    switch (condition.operator) {
      case '>': return value > condition.value;
      case '<': return value < condition.value;
      case '>=': return value >= condition.value;
      case '<=': return value <= condition.value;
      case '==': return value === condition.value;
      case '!=': return value !== condition.value;
      default: return false;
    }
  }

  private evaluateRate(condition: RateCondition, values: number[]): boolean {
    return condition.operator === '>' ? values.length > condition.value : values.length < condition.value;
  }

  private evaluateStatistical(condition: StatisticalCondition, value: number, values: number[]): boolean {
    if (values.length < 10) return false;
    const n = values.length;
    const mean = values.reduce((a, b) => a + b, 0) / n;
    const variance = values.reduce((sum, v) => sum + (v - mean) ** 2, 0) / n;
    const stddev = Math.sqrt(variance);
    if (stddev < 0.001) return false;
    return Math.abs(value - mean) / stddev > condition.deviation;
  }

  processTransaction(tx: TxEvent): AlertEvent[] {
    const alerts: AlertEvent[] = [];
    alerts.push(...this.evaluate('gas_used', parseInt(tx.gasUsed), { txHash: tx.hash }));
    alerts.push(...this.evaluate('msg_count', tx.messages.length, { txHash: tx.hash }));

    for (const event of tx.events) {
      if (!event.type.startsWith('wasm')) continue;
      for (const rule of this.rules.values()) {
        if (rule.condition.type !== 'event') continue;
        const evCondition = rule.condition as EventCondition;
        if (evCondition.event_type !== 'wasm' && evCondition.event_type !== event.type) continue;

        const attrs: Record<string, string> = {};
        for (const attr of event.attributes) attrs[attr.key] = attr.value;

        let allMatch = true;
        for (const [key, expected] of Object.entries(evCondition.attributes)) {
          const actual = attrs[key];
          if (!actual) { allMatch = false; break; }
          if (Array.isArray(expected)) { if (!expected.includes(actual)) { allMatch = false; break; } }
          else { if (actual !== expected) { allMatch = false; break; } }
        }

        if (allMatch) {
          const now = Date.now();
          const lastTriggered = this.triggeredRules.get(rule.id) || 0;
          if (now - lastTriggered >= rule.cooldownSeconds * 1000) {
            alerts.push({
              ruleId: rule.id, ruleName: rule.name, severity: rule.severity,
              message: rule.description, timestamp: now,
              context: { event, txHash: tx.hash },
              txHash: tx.hash, blockHeight: tx.height,
              contractAddress: attrs['_contract_address'],
            });
            this.triggeredRules.set(rule.id, now);
          }
        }
      }
    }
    return alerts;
  }

  on(event: string, callback: Function): void {
    if (!this.listeners.has(event)) this.listeners.set(event, []);
    this.listeners.get(event)!.push(callback);
  }

  private emit(event: string, data: any): void {
    const cbs = this.listeners.get(event) || [];
    for (const cb of cbs) try { cb(data); } catch {}
  }
}

4.3 告警去重与抑制

interface AlertDedupConfig {
  windowMs: number;
  maxAlertsPerWindow: number;
  groupBy: string[];
  suppressionRules: SuppressionRule[];
}

interface SuppressionRule {
  pattern: string;
  durationMs: number;
  reason: string;
}

class AlertDeduplicator {
  private config: AlertDedupConfig;
  private alertHistory: Map<string, AlertEvent[]> = new Map();
  private suppressions: Map<string, number> = new Map();

  constructor(config: AlertDedupConfig) { this.config = config; }

  shouldProcess(alert: AlertEvent): boolean {
    if (this.isSuppressed(alert)) return false;

    const dedupKey = this.getDedupKey(alert);
    const now = Date.now();

    if (!this.alertHistory.has(dedupKey)) this.alertHistory.set(dedupKey, []);
    const history = this.alertHistory.get(dedupKey)!;
    while (history.length > 0 && now - history[0].timestamp > this.config.windowMs) history.shift();

    if (history.length >= this.config.maxAlertsPerWindow) return false;

    history.push(alert);
    return true;
  }

  private getDedupKey(alert: AlertEvent): string {
    const parts = [alert.ruleId];
    for (const key of this.config.groupBy) {
      const value = alert.context[key] || (alert as any)[key];
      if (value !== undefined) parts.push(String(value));
    }
    return parts.join(':');
  }

  private isSuppressed(alert: AlertEvent): boolean {
    const now = Date.now();
    for (const rule of this.config.suppressionRules) {
      const expiresAt = this.suppressions.get(`suppress:${rule.pattern}`) || 0;
      if (now < expiresAt && (alert.ruleId.includes(rule.pattern) || alert.message.includes(rule.pattern))) {
        return true;
      }
    }
    return false;
  }

  addSuppression(pattern: string, durationMs: number, reason: string): void {
    this.suppressions.set(`suppress:${pattern}`, Date.now() + durationMs);
  }
}

5. MEV 攻击检测

5.1 三明治攻击检测

三明治攻击交易模式:
┌────────────────────────────────────────────┐
│  Block N                                     │
│  ┌──────────────────────────────────────┐   │
│  │ Tx A: 攻击者买入 (Front-run)          │   │
│  │ Tx B: 受害者交易 (目标交易)           │   │
│  │ Tx C: 攻击者卖出 (Back-run)           │   │
│  └──────────────────────────────────────┘   │
│                                              │
│  三要素:                                     │
│  1. 同一地址执行 A 和 C                      │
│  2. A 和 C 使用同一交易对                    │
│  3. B 在 A 和 C 之间                         │
└────────────────────────────────────────────┘
interface SwapEvent {
  txHash: string;
  blockHeight: number;
  txIndex: number;
  sender: string;
  pool: string;
  tokenIn: string;
  tokenOut: string;
  amountIn: string;
  amountOut: string;
  type: 'buy' | 'sell';
}

interface SandwichAttack {
  frontRunTx: string;
  victimTx: string;
  backRunTx: string;
  attacker: string;
  victim: string;
  pool: string;
  profit: string;
  blockHeight: number;
  timestamp: number;
}

class SandwichDetector {
  detect(swaps: SwapEvent[]): SandwichAttack[] {
    const attacks: SandwichAttack[] = [];
    const groupedByBlock = this.groupByBlock(swaps);

    for (const [blockHeight, blockSwaps] of groupedByBlock) {
      blockSwaps.sort((a, b) => a.txIndex - b.txIndex);

      for (let i = 0; i < blockSwaps.length - 2; i++) {
        for (let j = i + 2; j < blockSwaps.length; j++) {
          const first = blockSwaps[i];
          const last = blockSwaps[j];

          if (first.sender !== last.sender) continue;
          if (first.type !== 'buy' || last.type !== 'sell') continue;
          if (first.pool !== last.pool) continue;

          for (let k = i + 1; k < j; k++) {
            const victim = blockSwaps[k];
            if (victim.sender === first.sender) continue;
            if (victim.pool !== first.pool) continue;

            attacks.push({
              frontRunTx: first.txHash, victimTx: victim.txHash, backRunTx: last.txHash,
              attacker: first.sender, victim: victim.sender, pool: first.pool,
              profit: this.calculateProfit(first, last),
              blockHeight: parseInt(blockHeight), timestamp: Date.now(),
            });
            break;
          }
        }
      }
    }
    return attacks;
  }

  private groupByBlock(swaps: SwapEvent[]): Map<string, SwapEvent[]> {
    const grouped = new Map<string, SwapEvent[]>();
    for (const swap of swaps) {
      const key = String(swap.blockHeight);
      if (!grouped.has(key)) grouped.set(key, []);
      grouped.get(key)!.push(swap);
    }
    return grouped;
  }

  private calculateProfit(front: SwapEvent, back: SwapEvent): string {
    const profit = BigInt(back.amountOut) - BigInt(front.amountIn);
    return profit > 0n ? profit.toString() : '0';
  }
}

5.2 抢跑检测 (Front-running Detection)

interface PendingTx {
  hash: string; from: string; to: string;
  gasPrice: string; nonce: number;
  memo: string; firstSeen: number; data: string;
}

interface FrontrunAttempt {
  originalTx: string; frontrunTx: string;
  attacker: string; victim: string;
  targetContract: string; timeDelta: number;
  gasPriceDelta: string; method: string;
}

class FrontrunDetector {
  private pendingTxs: Map<string, PendingTx> = new Map();

  addPendingTx(tx: PendingTx): void {
    this.pendingTxs.set(`${tx.from}:${tx.nonce}`, tx);
    const now = Date.now();
    for (const [k, v] of this.pendingTxs) {
      if (now - v.firstSeen > 120000) this.pendingTxs.delete(k);
    }
  }

  detect(newTx: PendingTx, knownTxs: PendingTx[]): FrontrunAttempt[] {
    const attempts: FrontrunAttempt[] = [];
    const now = Date.now();

    for (const known of knownTxs) {
      if (known.to !== newTx.to) continue;
      const knownGas = BigInt(known.gasPrice || '0');
      const newGas = BigInt(newTx.gasPrice || '0');
      if (newGas <= knownGas * 110n / 100n) continue;
      if (now - known.firstSeen > 30000) continue;
      if (!this.sameFunctionCall(known.data, newTx.data)) continue;
      if (known.from === newTx.from) continue;

      attempts.push({
        originalTx: known.hash, frontrunTx: newTx.hash,
        attacker: newTx.from, victim: known.from,
        targetContract: newTx.to, timeDelta: newTx.firstSeen - known.firstSeen,
        gasPriceDelta: (newGas - knownGas).toString(),
        method: this.extractMethod(newTx.data),
      });
    }
    return attempts;
  }

  private sameFunctionCall(data1: string, data2: string): boolean {
    if (!data1 || !data2) return false;
    return data1.length >= 8 && data2.length >= 8
      && data1.substring(0, 8) === data2.substring(0, 8);
  }

  private extractMethod(data: string): string {
    if (!data || data.length < 8) return 'unknown';
    return data.substring(0, 8);
  }
}

5.3 时间强盗检测 (Time-bandit Attack Detection)

pub struct TimeBanditDetector {
    pub expected_block_producer: HashMap<u64, String>,
    pub block_time_history: VecDeque<BlockTimeRecord>,
    pub tx_reordering_threshold: f64,
}

pub struct BlockTimeRecord {
    pub height: u64, pub proposer: String,
    pub timestamp: u64, pub num_txs: u32,
    pub gas_used: u64, pub mev_extracted: Uint128,
}

impl TimeBanditDetector {
    pub fn new(tx_reordering_threshold: f64) -> Self {
        Self {
            expected_block_producer: HashMap::new(),
            block_time_history: VecDeque::with_capacity(1000),
            tx_reordering_threshold,
        }
    }

    pub fn detect_block_time_anomaly(&self, record: &BlockTimeRecord) -> DetectionResult {
        let mut alerts = Vec::new();
        if self.block_time_history.len() < 10 {
            return DetectionResult { detected: false, severity: Severity::None, alerts: Vec::new() };
        }

        let mean_interval = self.block_time_history.iter()
            .zip(self.block_time_history.iter().skip(1))
            .map(|(a, b)| b.timestamp - a.timestamp)
            .sum::<u64>() as f64 / (self.block_time_history.len() - 1) as f64;

        if record.num_txs == 0 && mean_interval > 0.0 {
            alerts.push(Alert {
                alert_type: AlertType::TimeBandit, severity: Severity::High,
                message: format!("Empty block at height {} proposer {}", record.height, record.proposer),
            });
        }

        if let Some(last) = self.block_time_history.back() {
            let interval = record.timestamp - last.timestamp;
            if interval > (mean_interval * 3.0) as u64 {
                alerts.push(Alert {
                    alert_type: AlertType::TimeBandit, severity: Severity::Medium,
                    message: format!("Block interval {}s (mean {:.1}s)", interval, mean_interval),
                });
            }
        }

        DetectionResult {
            detected: !alerts.is_empty(),
            severity: if alerts.iter().any(|a| a.severity == Severity::High) { Severity::High }
            else if !alerts.is_empty() { Severity::Medium } else { Severity::None },
            alerts,
        }
    }

    pub fn detect_validator_mev_mining(&self, validator: &str, recent_blocks: &[BlockTimeRecord]) -> DetectionResult {
        let mut alerts = Vec::new();
        let proposer_blocks: Vec<_> = recent_blocks.iter().filter(|b| b.proposer == validator).collect();

        if proposer_blocks.len() < 5 {
            return DetectionResult { detected: false, severity: Severity::None, alerts: Vec::new() };
        }

        let empty_count = proposer_blocks.iter().filter(|b| b.num_txs == 0).count();
        let empty_ratio = empty_count as f64 / proposer_blocks.len() as f64;
        if empty_ratio > 0.3 {
            alerts.push(Alert {
                alert_type: AlertType::ValidatorMisbehave, severity: Severity::High,
                message: format!("{}: {:.1}% empty blocks", validator, empty_ratio * 100.0),
            });
        }

        DetectionResult {
            detected: !alerts.is_empty(),
            severity: if alerts.iter().any(|a| a.severity == Severity::High) { Severity::High }
            else { Severity::Medium }, alerts,
        }
    }
}

5.4 MEV 保护策略

5.4.1 批量拍卖机制

pub struct BatchAuction {
    pub batch_interval: u64,
    pub last_batch_height: u64,
    pub pending_orders: Vec<Order>,
    pub min_bid_increment: Uint128,
}

impl BatchAuction {
    pub fn submit_order(&mut self, sender: String, order: Order) -> StdResult<()> {
        self.pending_orders.push(Order { sender, order, submitted_at: order.env_height });
        Ok(())
    }

    pub fn execute_batch(&mut self, current_height: u64) -> StdResult<BatchResult> {
        if current_height - self.last_batch_height < self.batch_interval {
            return Err(StdError::generic_err("Batch not ready"));
        }
        let orders = std::mem::take(&mut self.pending_orders);
        let clearing_price = self.calculate_clearing_price(&orders)?;
        let mut results = Vec::new();
        for order in &orders {
            results.push(self.fill_order(order, clearing_price)?);
        }
        self.last_batch_height = current_height;
        Ok(BatchResult { clearing_price, orders_executed: results.len() as u32,
            total_volume: results.iter().map(|r| r.amount).sum() })
    }

    fn calculate_clearing_price(&self, orders: &[Order]) -> StdResult<Decimal> {
        if orders.is_empty() { return Err(StdError::generic_err("No orders")); }
        let total_buy: Uint128 = orders.iter().filter(|o| o.order_type == OrderType::Buy)
            .map(|o| o.amount).sum();
        let total_sell: Uint128 = orders.iter().filter(|o| o.order_type == OrderType::Sell)
            .map(|o| o.amount).sum();
        if total_buy.is_zero() || total_sell.is_zero() {
            return Err(StdError::generic_err("Insufficient volume"));
        }
        Ok(Decimal::from_ratio(total_buy, total_sell))
    }

    fn fill_order(&self, order: &Order, price: Decimal) -> StdResult<FillResult> {
        Ok(FillResult { trader: order.sender.clone(), original_amount: order.amount,
            filled_amount: order.amount * price, price })
    }
}

6. 链上蜜罐 / Rug Pull 检测

6.1 蜜罐合约检测

特征 正常合约 蜜罐合约
转账函数 对所有用户开放 对非 owner 用户受限
Buy 功能 正常购买 购买正常
Sell/Withdraw 正常卖出 交易失败或返回错误
Gas 消耗 可预测 Sell 时异常高
事件发射 正常 买入发事件,卖出不发
pub struct HoneypotDetector {
    pub test_amounts: Vec<Uint128>,
    pub gas_threshold: u64,
}

impl HoneypotDetector {
    pub fn simulate_trade_test(&self, deps: &DepsMut, contract: &Addr, test_account: &Addr) -> TradeTestResult {
        let mut results = Vec::new();
        for amount in &self.test_amounts {
            let buy_result = self.simulate_buy(deps, contract, test_account, *amount);
            let sell_result = self.simulate_sell(deps, contract, test_account, *amount);
            results.push(TradeTest {
                amount: *amount, buy_success: buy_result.success, buy_gas: buy_result.gas_used,
                sell_success: sell_result.success, sell_gas: sell_result.gas_used,
                sell_error: sell_result.error_message,
            });
        }
        let buy_count = results.iter().filter(|r| r.buy_success).count();
        let sell_count = results.iter().filter(|r| r.sell_success).count();
        TradeTestResult {
            is_honeypot: buy_count > 0 && sell_count == 0,
            confidence: if buy_count > 0 { 0.9 } else { 0.0 },
            details: results,
        }
    }

    fn simulate_buy(&self, _deps: &DepsMut, _contract: &Addr, _account: &Addr, _amount: Uint128) -> SimulateResult {
        SimulateResult { success: true, gas_used: 100_000, error_message: None }
    }

    fn simulate_sell(&self, _deps: &DepsMut, _contract: &Addr, _account: &Addr, _amount: Uint128) -> SimulateResult {
        SimulateResult { success: false, gas_used: 500_000, error_message: Some("Unauthorized".to_string()) }
    }
}

6.2 流动性移除监控

interface LiquidityPool {
  address: string; token0: string; token1: string;
  reserve0: string; reserve1: string;
  totalLiquidity: string; lpTokenSupply: string;
}

interface LiquidityRemoval {
  pool: string; caller: string;
  amount0: string; amount1: string;
  liquidityRemoved: string; remainingLiquidity: string;
  removalPercent: number; txHash: string;
  blockHeight: number; timestamp: number;
}

class LiquidityMonitor {
  private pools: Map<string, LiquidityPool> = new Map();
  private removalThreshold = 50;

  async updatePool(pool: LiquidityPool): Promise<LiquidityRemoval | null> {
    const existing = this.pools.get(pool.address);
    this.pools.set(pool.address, pool);
    if (!existing) return null;

    const existingLp = BigInt(existing.totalLiquidity);
    const currentLp = BigInt(pool.totalLiquidity);
    if (currentLp >= existingLp) return null;

    const removed = existingLp - currentLp;
    const removalPercent = Number(removed * 100n / existingLp);
    if (removalPercent < this.removalThreshold) return null;

    return {
      pool: pool.address, caller: '',
      amount0: (BigInt(existing.reserve0) - BigInt(pool.reserve0)).toString(),
      amount1: (BigInt(existing.reserve1) - BigInt(pool.reserve1)).toString(),
      liquidityRemoved: removed.toString(),
      remainingLiquidity: currentLp.toString(),
      removalPercent, txHash: '', blockHeight: 0, timestamp: Date.now(),
    };
  }
}

6.3 Rug Pull 综合评分

pub struct RugPullScorer {
    pub liquidity_weight: f64, pub age_weight: f64,
    pub holder_weight: f64, pub volume_weight: f64, pub code_weight: f64,
}

impl RugPullScorer {
    pub fn new() -> Self {
        Self { liquidity_weight: 0.30, age_weight: 0.20,
            holder_weight: 0.15, volume_weight: 0.15, code_weight: 0.20 }
    }

    pub fn score(&self, contract: &ScoredContract) -> RugPullScore {
        let mut score = 0.0;
        let mut signals = Vec::new();

        let liquidity_score = if contract.total_liquidity > Uint128::zero() {
            let locked_ratio = Decimal::from_ratio(contract.locked_liquidity, contract.total_liquidity);
            let s = (Decimal::one() - locked_ratio).to_f64();
            if s > 0.5 { signals.push("Low locked liquidity ratio".to_string()); }
            s
        } else { 1.0 };

        let age_days = contract.age_seconds as f64 / 86400.0;
        let age_score = if age_days < 1.0 { signals.push("Contract less than 1 day old".to_string()); 1.0 }
        else if age_days < 7.0 { 0.6 } else if age_days < 30.0 { 0.3 } else { 0.0 };

        let holder_score = if contract.total_holders > 0 {
            let top10_pct = Decimal::from_ratio(contract.top_10_holders_balance, contract.total_supply);
            let s = top10_pct.to_f64();
            if s > 0.8 { signals.push("Top 10 holders own > 80% supply".to_string()); }
            s
        } else { 1.0 };

        score = self.liquidity_weight * liquidity_score + self.age_weight * age_score
            + self.holder_weight * holder_score + self.volume_weight * 0.0 + self.code_weight * contract.code_risk_score;

        let severity = if score >= 0.8 { Severity::Critical }
        else if score >= 0.6 { Severity::High }
        else if score >= 0.4 { Severity::Medium }
        else if score >= 0.2 { Severity::Low } else { Severity::None };

        RugPullScore { contract: contract.address.clone(), score, severity, signals,
            details: RugPullDetails { liquidity_score, age_score, holder_score, volume_score: 0.0, code_score: contract.code_risk_score } }
    }
}

pub struct ScoredContract {
    pub address: String, pub total_liquidity: Uint128, pub locked_liquidity: Uint128,
    pub age_seconds: u64, pub total_supply: Uint128, pub total_holders: u64,
    pub top_10_holders_balance: Uint128, pub avg_daily_volume: Uint128,
    pub current_daily_volume: Uint128, pub code_risk_score: f64,
}

pub struct RugPullScore {
    pub contract: String, pub score: f64, pub severity: Severity,
    pub signals: Vec<String>, pub details: RugPullDetails,
}

pub struct RugPullDetails {
    pub liquidity_score: f64, pub age_score: f64, pub holder_score: f64,
    pub volume_score: f64, pub code_score: f64,
}

7. MSG Chain 安全监控栈

7.1 架构总览

                    MSG Chain 安全监控栈
┌──────────────────────────────────────────────────────────┐
│                   告警响应层                              │
│  Telegram │ Slack │ PagerDuty │ Webhook │ 自动响应脚本     │
├──────────────────────────────────────────────────────────┤
│                   分析引擎层                              │
│  ┌─────────┐ ┌──────────┐ ┌──────────┐ ┌────────────┐   │
│  │ 规则引擎 │ │ 基线分析  │ │ ML 检测  │ │ 威胁情报    │   │
│  └─────────┘ └──────────┘ └──────────┘ └────────────┘   │
├──────────────────────────────────────────────────────────┤
│                   事件处理层                              │
│  ┌─────────┐ ┌──────────┐ ┌──────────┐ ┌────────────┐   │
│  │ 事件过滤 │ │ 协议解码  │ │ 特征提取  │ │ 告警去重    │   │
│  └─────────┘ └──────────┘ └──────────┘ └────────────┘   │
├──────────────────────────────────────────────────────────┤
│                   索引存储层                              │
│  ┌─────────┐ ┌──────────┐ ┌──────────┐ ┌────────────┐   │
│  │ 交易索引 │ │ 事件索引  │ │ 合约状态  │ │ 告警历史    │   │
│  └─────────┘ └──────────┘ └──────────┘ └────────────┘   │
├──────────────────────────────────────────────────────────┤
│                   数据采集层                              │
│  ┌─────────┐ ┌──────────┐ ┌──────────┐ ┌────────────┐   │
│  │ RPC 节点 │ │ Mempool  │ │ Indexer  │ │ 链上查询    │   │
│  └─────────┘ └──────────┘ └──────────┘ └────────────┘   │
├──────────────────────────────────────────────────────────┤
│                   MSG Chain 节点                          │
│             Validators / Full Nodes / RPC Nodes            │
└──────────────────────────────────────────────────────────┘

7.2 组件部署

version: '3.8'
services:
  msg-chain-rpc:
    image: msgchain/node:latest
    ports:
      - "26657:26657"
      - "26656:26656"
    volumes:
      - ./node-data:/root/.msgchain
      - ./config:/root/.msgchain/config
    environment:
      - CHAIN_ID=msg-chain-1
      - MONIKER=security-monitor
    command: ["msgchaind", "start", "--rpc.laddr=tcp://0.0.0.0:26657"]

  security-indexer:
    image: msgchain/security-indexer:latest
    depends_on: [msg-chain-rpc]
    volumes:
      - ./indexer-data:/data
      - ./config/indexer.yaml:/app/config.yaml
    environment:
      - RPC_ENDPOINT=http://msg-chain-rpc:26657
      - START_HEIGHT=0
      - BATCH_SIZE=100
      - STORE_TYPE=postgresql
    command: ["indexer", "start"]

  analysis-engine:
    image: msgchain/security-analysis:latest
    depends_on: [security-indexer]
    ports: ["8080:8080"]
    volumes:
      - ./config/rules.yaml:/app/rules.yaml
      - ./config/baselines:/app/baselines
    environment:
      - INDEXER_GRPC=security-indexer:9090
      - DB_URL=postgresql://postgres:password@postgres:5432/security
      - REDIS_URL=redis://redis:6379
    command: ["analysis-engine", "serve"]

  alert-manager:
    image: msgchain/alert-manager:latest
    depends_on: [analysis-engine]
    ports: ["9093:9093"]
    volumes:
      - ./config/alertmanager.yaml:/etc/alertmanager/config.yml
    environment:
      - SLACK_TOKEN=${SLACK_TOKEN}
      - TELEGRAM_TOKEN=${TELEGRAM_TOKEN}
      - PAGERDUTY_KEY=${PAGERDUTY_KEY}
    command: ["alertmanager"]

  auto-responder:
    image: msgchain/auto-responder:latest
    depends_on: [alert-manager]
    volumes:
      - ./config/responder.yaml:/app/config.yaml
      - ./config/msgchain-key:/app/key.json
    environment:
      - CHAIN_RPC=http://msg-chain-rpc:26657
      - MNEMONIC_FILE=/app/key.json
    command: ["auto-responder", "start"]

  security-dashboard:
    image: grafana/grafana:latest
    depends_on: [postgres, prometheus]
    ports: ["3000:3000"]
    volumes:
      - ./grafana/dashboards:/etc/grafana/provisioning/dashboards
      - ./grafana/datasources:/etc/grafana/provisioning/datasources
    environment:
      - GF_SECURITY_ADMIN_PASSWORD=${GRAFANA_PASSWORD}

  prometheus:
    image: prom/prometheus:latest
    ports: ["9090:9090"]
    volumes:
      - ./prometheus/prometheus.yml:/etc/prometheus/prometheus.yml
      - prometheus-data:/prometheus

  postgres:
    image: postgres:15
    ports: ["5432:5432"]
    volumes:
      - postgres-data:/var/lib/postgresql/data
    environment:
      - POSTGRES_DB=security
      - POSTGRES_PASSWORD=${DB_PASSWORD}

  redis:
    image: redis:7-alpine
    ports: ["6379:6379"]
    volumes:
      - redis-data:/data

volumes:
  prometheus-data:
  postgres-data:
  redis-data:

7.3 数据流

数据流: 交易 → 告警 → 响应的完整链路

1. Mempool 捕获
   MempoolMonitor → PreflightChecker → (通过/拒绝)
   ↓
2. 区块确认
   MSG Chain 节点 → SecurityIndexer → Block + Tx 事件
   ↓
3. 事件处理
   Indexer → EventFilter → ProtocolDecoder → FeatureExtractor
   ↓
4. 分析引擎
   FeatureVector → RuleEngine + BaselineMonitor + MLDetector
   ↓
5. 告警生成
   AnalysisResult → AlertDeduplicator → AlertEvent
   ↓
6. 通知分发
   AlertManager → Slack / Telegram / PagerDuty / Webhook
   ↓
7. 自动响应
   AutoResponder → PauseContract / MultisigTrigger / EmergencyProposal
   ↓
8. 状态更新
   ResponseStatus → Dashboard / AuditLog

7.4 Kafka 事件流

topics:
  - name: security.raw.blocks
    partitions: 4
    replication: 2
    retention.ms: 604800000

  - name: security.raw.transactions
    partitions: 8
    replication: 2
    retention.ms: 604800000

  - name: security.raw.mempool
    partitions: 4
    replication: 2
    retention.ms: 3600000

  - name: security.processed.events
    partitions: 8
    replication: 2
    retention.ms: 2592000000

  - name: security.alerts
    partitions: 4
    replication: 3
    retention.ms: 2592000000

  - name: security.responses
    partitions: 2
    replication: 3
    retention.ms: 2592000000

8. 自动化响应

8.1 自动暂停合约

import { SigningCosmWasmClient } from '@cosmjs/cosmwasm-stargate';
import { DirectSecp256k1HdWallet } from '@cosmjs/proto-signing';

interface PauseAction {
  contractAddress: string;
  reason: string;
  triggeredBy: string;
  alertId: string;
  severity: string;
  autoResume: boolean;
  autoResumeAfterBlocks?: number;
}

interface PauseResult {
  success: boolean;
  txHash: string;
  error?: string;
  action: PauseAction;
}

class AutoPauseResponder {
  private client: SigningCosmWasmClient | null = null;
  private wallet: DirectSecp256k1HdWallet | null = null;
  private config: {
    rpcEndpoint: string; mnemonic: string; chainId: string;
    bech32Prefix: string; minGasPrice: string; maxGas: number;
  };

  constructor(config: typeof AutoPauseResponder.prototype.config) { this.config = config; }

  async init(): Promise<void> {
    this.wallet = await DirectSecp256k1HdWallet.fromMnemonic(this.config.mnemonic, {
      prefix: this.config.bech32Prefix,
    });
    this.client = await SigningCosmWasmClient.connectWithSigner(
      this.config.rpcEndpoint, this.wallet, {
        gasPrice: { denom: 'umsg', amount: this.config.minGasPrice },
      },
    );
  }

  async pauseContract(action: PauseAction): Promise<PauseResult> {
    if (!this.client) return { success: false, txHash: '', error: 'Client not initialized', action };
    const sender = (await this.wallet!.getAccounts())[0].address;
    try {
      const result = await this.client.execute(sender, action.contractAddress,
        { pause: { reason: action.reason, auto_resume: action.autoResume } },
        'auto', undefined, []);
      return { success: result.code === 0, txHash: result.transactionHash, error: result.rawLog, action };
    } catch (err) {
      return { success: false, txHash: '', error: (err as Error).message, action };
    }
  }

  async batchPause(contracts: string[], reason: string, alertId: string): Promise<PauseResult[]> {
    const results: PauseResult[] = [];
    for (const contract of contracts) {
      results.push(await this.pauseContract({
        contractAddress: contract, reason, triggeredBy: 'auto-responder',
        alertId, severity: 'CRITICAL', autoResume: false,
      }));
    }
    return results;
  }
}

8.2 触发多签提案

interface MultisigProposal {
  id: string; title: string; description: string;
  contractAddress: string; actions: ProposalAction[];
  proposer: string; requiredVotes: number;
  deadline: number; status: 'pending' | 'active' | 'approved' | 'executed' | 'rejected';
}

interface ProposalAction {
  type: 'pause' | 'upgrade' | 'emergency_withdraw' | 'update_config' | 'freeze';
  params: Record<string, unknown>;
}

class MultisigTrigger {
  private rpcEndpoint: string;
  private client: CosmWasmClient | null = null;
  private multisigContract: string;

  constructor(rpcEndpoint: string, multisigContract: string) {
    this.rpcEndpoint = rpcEndpoint; this.multisigContract = multisigContract;
  }

  async createEmergencyProposal(proposal: MultisigProposal): Promise<string> {
    if (!this.client) this.client = await CosmWasmClient.connect(this.rpcEndpoint);
    return proposal.id;
  }

  async voteProposal(proposalId: string, voter: string, approve: boolean): Promise<boolean> { return true; }
  async executeProposal(proposalId: string): Promise<boolean> { return true; }

  async getProposalStatus(proposalId: string): Promise<MultisigProposal['status']> {
    const result = await this.client!.queryContractSmart(this.multisigContract, {
      proposal: { proposal_id: proposalId },
    });
    return result.status;
  }
}

8.3 紧急治理提案

pub struct EmergencyGovernanceProposal {
    pub proposer: String, pub proposal_type: EmergencyProposalType,
    pub target_contract: String, pub action: EmergencyAction,
    pub reason: String, pub severity: Severity,
    pub created_at: u64, pub voting_period: u64,
    pub min_approval_percentage: Decimal,
}

pub enum EmergencyProposalType {
    PauseContract, EmergencyUpgrade,
    EmergencyWithdraw { recipient: String },
    UpdateConfig { key: String, value: Binary },
    FreezeFunds { target_account: String },
}

pub enum EmergencyAction {
    Immediate, Delayed { delay_blocks: u64 },
    Scheduled { execute_at_height: u64 },
}

pub struct EmergencyGovernance {
    pub proposal_threshold: Uint128,
    pub voting_period_blocks: u64,
    pub min_approval_pct: Decimal,
    pub emergency_multisig: Vec<Addr>,
    pub required_signatures: u32,
}

impl EmergencyGovernance {
    pub fn create_proposal(&self, deps: DepsMut, env: Env, sender: &Addr,
        proposal_type: EmergencyProposalType, target: String, reason: String) -> StdResult<Response> {
        if !self.emergency_multisig.contains(sender) {
            return Err(StdError::generic_err("Only emergency multisig members can create proposals"));
        }
        let proposal_id = env.block.height;
        let proposal = EmergencyGovernanceProposal {
            proposer: sender.to_string(), proposal_type, target_contract: target,
            action: EmergencyAction::Immediate, reason, severity: Severity::Critical,
            created_at: env.block.time.seconds(), voting_period: self.voting_period_blocks,
            min_approval_percentage: self.min_approval_pct,
        };
        EMERGENCY_PROPOSALS.save(deps.storage, &proposal_id, &proposal)?;
        Ok(Response::new().add_attribute("action", "emergency_proposal_created")
            .add_attribute("proposal_id", proposal_id.to_string()))
    }

    pub fn vote(&self, deps: DepsMut, env: Env, sender: &Addr, proposal_id: u64, approve: bool) -> StdResult<Response> {
        let proposal = EMERGENCY_PROPOSALS.load(deps.storage, &proposal_id)?;
        if !self.emergency_multisig.contains(sender) {
            return Err(StdError::generic_err("Not authorized to vote"));
        }
        if env.block.height > proposal.created_at + proposal.voting_period {
            return Err(StdError::generic_err("Voting period expired"));
        }
        let mut votes = PROPOSAL_VOTES.may_load(deps.storage, &proposal_id)?.unwrap_or_default();
        let voter = sender.to_string();
        if votes.contains(&voter) { return Err(StdError::generic_err("Already voted")); }
        votes.push(voter.clone());

        if approve {
            PROPOSAL_APPROVES.save(deps.storage, &proposal_id, &votes)?;
        } else {
            PROPOSAL_REJECTS.save(deps.storage, &proposal_id, &votes)?;
        }

        let approves = PROPOSAL_APPROVES.may_load(deps.storage, &proposal_id)?.unwrap_or_default();
        let approval_pct = Decimal::from_ratio(Uint128::new(approves.len() as u128),
            Uint128::new(self.emergency_multisig.len() as u128));
        if approval_pct >= self.min_approval_pct {
            return self.execute_proposal(deps, env, proposal_id);
        }
        Ok(Response::new().add_attribute("action", "emergency_vote").add_attribute("voter", voter))
    }

    pub fn execute_proposal(&self, deps: DepsMut, _env: Env, proposal_id: u64) -> StdResult<Response> {
        let proposal = EMERGENCY_PROPOSALS.load(deps.storage, &proposal_id)?;
        let response = match proposal.proposal_type {
            EmergencyProposalType::PauseContract => {
                PAUSED_CONTRACTS.save(deps.storage, &proposal.target_contract, &true)?;
                Response::new().add_attribute("action", "contract_paused")
            }
            EmergencyProposalType::EmergencyWithdraw { recipient } => {
                Response::new().add_attribute("action", "emergency_withdraw")
                    .add_attribute("recipient", recipient)
            }
            _ => Response::new().add_attribute("action", "proposal_executed"),
        };
        EMERGENCY_PROPOSALS.remove(deps.storage, &proposal_id);
        Ok(response)
    }
}

9. 威胁情报

9.1 漏洞数据库订阅

interface VulnerabilityFeed {
  id: string;
  title: string;
  description: string;
  cveId?: string;
  severity: 'CRITICAL' | 'HIGH' | 'MEDIUM' | 'LOW';
  affectedPlatforms: string[];
  affectedVersions: string[];
  publishedAt: string;
  references: string[];
  patchAvailable: boolean;
  exploitAvailable: boolean;
}

interface VulnerabilityDatabase {
  name: string;
  type: 'cve' | 'cosmos_security' | 'cosmwasm_advisory' | 'github_advisory';
  feedUrl: string;
  updateIntervalMs: number;
  lastUpdate?: number;
  vulnerabilities: Map<string, VulnerabilityFeed>;
}

class ThreatIntelligence {
  private databases: VulnerabilityDatabase[] = [];

  addDatabase(db: VulnerabilityDatabase): void {
    this.databases.push(db);
    this.startPolling(db);
  }

  private startPolling(db: VulnerabilityDatabase): void {
    setInterval(async () => {
      try { await this.fetchUpdates(db); } catch (err) {
        console.error(`Failed to fetch ${db.name}:`, err);
      }
    }, db.updateIntervalMs);
  }

  private async fetchUpdates(db: VulnerabilityDatabase): Promise<void> {
    const response = await fetch(db.feedUrl);
    const data = await response.json();
    for (const vuln of data.vulnerabilities || []) {
      if (!db.vulnerabilities.has(vuln.id)) {
        db.vulnerabilities.set(vuln.id, vuln);
        this.notifyVulnerability(vuln);
      }
    }
    db.lastUpdate = Date.now();
  }

  private notifyVulnerability(vuln: VulnerabilityFeed): void {
    if (vuln.severity === 'CRITICAL' || vuln.severity === 'HIGH') {
      console.log(`[ThreatIntel] ${vuln.severity} ${vuln.title}`);
    }
  }

  getAffectedContracts(): string[] {
    const contracts: string[] = [];
    for (const db of this.databases) {
      for (const vuln of db.vulnerabilities.values()) {
        if (vuln.affectedPlatforms.includes('cosmwasm')) {
          // 匹配受影响的合约
        }
      }
    }
    return contracts;
  }

  getVulnerability(id: string): VulnerabilityFeed | undefined {
    for (const db of this.databases) {
      if (db.vulnerabilities.has(id)) return db.vulnerabilities.get(id);
    }
    return undefined;
  }
}

9.2 Cosmos / CosmWasm 安全公告

公告编号 标题 影响 修复版本
CSA-2025-001 wasmd 存储键碰撞漏洞 合约状态污染 wasmd >= 0.45.0
CSA-2025-002 IBC 重放攻击 跨链双花 ibc-go >= 8.2.0
CSA-2025-003 CosmWasm 反序列化拒绝服务 节点崩溃 cosmwasm-vm >= 1.5.2
CSA-2026-001 非预期的 SubMsg gas 退款 Gas 计费绕过 wasmd >= 0.47.1
CSA-2026-002 Dilithium-5 公钥压缩绕过 签名伪造 msgchaind >= 0.12.0

9.3 威胁情报源配置

# threat_intel.yaml
sources:
  - name: cosmos-security-bulletin
    type: cosmos_security
    feed_url: https://github.com/cosmos/security/raw/main/advisories.json
    update_interval: 3600000
    severity_filter: [CRITICAL, HIGH]

  - name: cosmwasm-advisories
    type: cosmwasm_advisory
    feed_url: https://raw.githubusercontent.com/CosmWasm/advisories/main/advisories.json
    update_interval: 3600000
    severity_filter: [CRITICAL, HIGH, MEDIUM]

  - name: nvd-cve-feed
    type: cve
    feed_url: https://services.nvd.nist.gov/rest/json/cves/2.0
    update_interval: 86400000
    keywords: [cosmwasm, cosmos-sdk, wasmd, rust-wasm]

  - name: github-advisory-db
    type: github_advisory
    feed_url: https://api.github.com/advisories
    update_interval: 86400000
    ecosystem: [crates.io, wasm]

alert_on_match:
  - match_type: contract_code_id
    action: scan_all_instances

  - match_type: dependency_version
    action: notify_development_team

  - match_type: platform_version
    action: upgrade_reminder

10. 安全仪表盘设计与指标

10.1 关键指标定义

指标 定义 采集频率 告警阈值
security_tvl_at_risk 高风险合约中锁定的总价值 每区块 > TVL 的 10%
security_pending_attacks 检测到但未确认的攻击数 实时 > 5
security_active_alerts 当前活跃告警数 实时 > 20
security_response_time_p50 告警到响应的中位时间 每小时 > 5 分钟
security_response_time_p99 告警到响应的 P99 时间 每小时 > 30 分钟
security_false_positive_rate 误报率 (7 天滚动) 每日 > 10%
security_mev_extracted_total 检测到的 MEV 提取总量 每日 -
security_honeypot_detected 检测到的蜜罐合约数 每日 -
security_anomaly_score_avg 全网平均异常评分 每区块 > 0.7
// security-metrics.ts
import { Gauge, Counter, Histogram } from 'prom-client';

export const security_tvl_at_risk = new Gauge({
  name: 'msg_security_tvl_at_risk',
  help: '高风险合约中 TVL 总量',
  labelNames: ['severity', 'contract_type'],
});

export const security_pending_attacks = new Gauge({
  name: 'msg_security_pending_attacks',
  help: '待确认攻击数',
  labelNames: ['attack_type'],
});

export const security_active_alerts = new Gauge({
  name: 'msg_security_active_alerts',
  help: '当前活跃告警数',
  labelNames: ['severity', 'rule_id'],
});

export const security_alerts_total = new Counter({
  name: 'msg_security_alerts_total',
  help: '历史告警总数',
  labelNames: ['severity', 'rule_id', 'contract'],
});

export const security_response_time = new Histogram({
  name: 'msg_security_response_time_seconds',
  help: '告警响应时间分布',
  labelNames: ['severity', 'response_type'],
  buckets: [10, 30, 60, 120, 300, 600, 1800, 3600],
});

export const security_auto_responses_total = new Counter({
  name: 'msg_security_auto_responses_total',
  help: '自动化响应执行总数',
  labelNames: ['response_type', 'outcome'],
});

export const security_false_positive_rate = new Gauge({
  name: 'msg_security_false_positive_rate',
  help: '7 天滚动误报率',
  labelNames: ['rule_id'],
});

export const security_mev_extracted_total = new Counter({
  name: 'msg_security_mev_extracted_total',
  help: '检测到的 MEV 提取量',
  labelNames: ['attack_type', 'validator'],
});

export const security_honeypot_detected = new Counter({
  name: 'msg_security_honeypot_detected',
  help: '检测到的蜜罐合约数',
  labelNames: ['risk_level'],
});

export const security_baseline_deviations = new Counter({
  name: 'msg_security_baseline_deviations',
  help: '基线偏差检测数',
  labelNames: ['contract', 'metric'],
});

export const security_preflight_rejected = new Counter({
  name: 'msg_security_preflight_rejected',
  help: 'Preflight 检查拒绝的交易数',
  labelNames: ['reason'],
});

10.2 Grafana 仪表盘

{
  "dashboard": {
    "title": "MSG Chain 安全监控",
    "uid": "msg-chain-security",
    "timezone": "utc",
    "panels": [
      {
        "title": "TVL at Risk",
        "type": "stat",
        "datasource": "Prometheus",
        "targets": [
          {
            "expr": "sum(msg_security_tvl_at_risk) by (severity)",
            "legendFormat": "{{severity}}"
          }
        ]
      },
      {
        "title": "活跃告警数",
        "type": "bargauge",
        "datasource": "Prometheus",
        "targets": [
          {
            "expr": "sum(msg_security_active_alerts) by (severity)",
            "legendFormat": "{{severity}}"
          }
        ]
      },
      {
        "title": "告警趋势 (24h)",
        "type": "timeseries",
        "datasource": "Prometheus",
        "targets": [
          {
            "expr": "rate(msg_security_alerts_total[5m])",
            "legendFormat": "{{severity}}"
          }
        ]
      },
      {
        "title": "响应时间 SLA",
        "type": "stat",
        "datasource": "Prometheus",
        "targets": [
          {
            "expr": "histogram_quantile(0.99, rate(msg_security_response_time_seconds_bucket[24h]))",
            "legendFormat": "P99"
          }
        ]
      },
      {
        "title": "待处理攻击",
        "type": "table",
        "datasource": "PostgreSQL",
        "targets": [
          {
            "rawSql": "SELECT attack_type, count(*), max(severity), min(detected_at) FROM pending_attacks WHERE status='open' GROUP BY attack_type"
          }
        ]
      },
      {
        "title": "合约健康评分",
        "type": "table",
        "datasource": "PostgreSQL",
        "targets": [
          {
            "rawSql": "SELECT contract_address, risk_score, deviation_count, last_alert FROM contract_health ORDER BY risk_score DESC LIMIT 20"
          }
        ]
      },
      {
        "title": "MEV 攻击检测",
        "type": "timeseries",
        "datasource": "Prometheus",
        "targets": [
          {
            "expr": "rate(msg_security_mev_extracted_total[1h])",
            "legendFormat": "{{attack_type}}"
          }
        ]
      },
      {
        "title": "自动响应成功率",
        "type": "gauge",
        "datasource": "Prometheus",
        "targets": [
          {
            "expr": "sum(msg_security_auto_responses_total{outcome='success'}) / sum(msg_security_auto_responses_total) * 100"
          }
        ]
      },
      {
        "title": "误报率 (7d)",
        "type": "timeseries",
        "datasource": "Prometheus",
        "targets": [
          {
            "expr": "msg_security_false_positive_rate",
            "legendFormat": "{{rule_id}}"
          }
        ]
      },
      {
        "title": "Preflight 拒绝统计",
        "type": "piechart",
        "datasource": "Prometheus",
        "targets": [
          {
            "expr": "rate(msg_security_preflight_rejected[24h])",
            "legendFormat": "{{reason}}"
          }
        ]
      }
    ],
    "rows": [
      { "title": "概览", "panels": [0, 1, 2, 3] },
      { "title": "攻击检测", "panels": [4, 5, 6] },
      { "title": "系统状态", "panels": [7, 8, 9] }
    ]
  }
}

10.3 SLA 定义

告警级别 响应时间 (P99) 确认时间 遏制时间 升级策略
P0 (Critical) 1 分钟 5 分钟 15 分钟 立即升级到全体值班
P1 (High) 5 分钟 15 分钟 30 分钟 升级到安全团队
P2 (Medium) 15 分钟 1 小时 4 小时 升级到开发团队
P3 (Low) 1 小时 24 小时 7 天 记录并跟踪
P4 (Info) 24 小时 - - 仅记录

11. 合规监控

11.1 AML/KYC 链上地址标记

interface SanctionedAddress {
  address: string;
  chain: string;
  source: string;
  addedAt: string;
  reason: string;
  riskLevel: 'HIGH' | 'MEDIUM' | 'LOW';
  restrictions: string[];
}

interface AMLConfig {
  sanctionLists: string[];
  updateIntervalMs: number;
  actionOnMatch: 'block' | 'alert' | 'log' | 'escalate';
  whitelistedContracts: string[];
  bech32Prefix: string;
}

class AMLMonitor {
  private sanctions: Set<string> = new Set();
  private config: AMLConfig;

  constructor(config: AMLConfig) { this.config = config; }

  async loadSanctionLists(): Promise<void> {
    for (const url of this.config.sanctionLists) {
      try {
        const response = await fetch(url);
        const data = await response.json();
        for (const entry of data.addresses || []) {
          if (entry.chain === 'cosmos' || !entry.chain) {
            this.sanctions.add(entry.address);
          }
        }
      } catch (err) {
        console.error(`Failed to load sanction list ${url}:`, err);
      }
    }
  }

  checkTransaction(tx: TxEvent): AMLResult {
    const matchedAddresses: string[] = [];

    for (const msg of tx.messages) {
      if (this.sanctions.has(msg.sender)) matchedAddresses.push(msg.sender);
      if (this.sanctions.has(msg.contract)) matchedAddresses.push(msg.contract);
    }

    if (matchedAddresses.length > 0) {
      const action = this.determineAction(matchedAddresses);
      console.log(`[AML] Match found: ${matchedAddresses.join(', ')} action=${action}`);
      return { matched: true, addresses: matchedAddresses, action, txHash: tx.hash };
    }

    return { matched: false, addresses: [], action: 'none', txHash: tx.hash };
  }

  private determineAction(addresses: string[]): string {
    return addresses.length > 0 ? this.config.actionOnMatch : 'none';
  }
}

11.2 制裁地址过滤

# sanctions.yaml
chain: msg-chain-1
bech32_prefix: msg

sanction_lists:
  - name: ofac-sdn
    url: "https://sanctionslistservice.ofac.treas.gov/api/sdn"
    type: ofac
    update_interval: 86400000

  - name: eu-consolidated
    url: "https://webgate.ec.europa.eu/fsd/fsf/api/v1/public/json"
    type: eu_sanctions
    update_interval: 86400000

  - name: msg-chain-banned
    url: "https://msgchain.org/security/sanctions.json"
    type: custom
    update_interval: 3600000

actions:
  on_match_tx: alert_only
  on_match_instantiate: block
  on_match_transfer: escalate_to_admin

exemptions:
  - contract: "msg1xxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxx"
    reason: "Official bridge contract"
  - contract: "msg1yyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyyy"
    reason: "Multisig governance"

12. 案例:监控 + 自动响应实战演练

12.1 场景:重入攻击自动阻断

攻击时间线

时间 事件 监控系统动作
T+0s 攻击者部署恶意合约 Indexer 捕获 instantiate 事件
T+5s 攻击者调用 withdraw Mempool 捕获,Preflight 标记 high_risk
T+6s 交易进入区块 规则引擎匹配 large_transfer + reentrancy_pattern
T+7s 告警触发 AlertManager 发送 Slack 通知
T+8s 系统确认攻击 AutoResponder 发起 pause_contract
T+10s 合约暂停 攻击合约和受害合约同时暂停
T+30s 多签提案自动创建 MultisigTrigger 创建紧急提案
T+60s 安全团队收到通知 PagerDuty 升级告警

监控配置

# 重入攻击检测规则
- id: reentrancy_detection
  name: 重入攻击检测
  description: 同一调用者在同区块内多次调用同一合约
  severity: CRITICAL
  condition:
    type: statistical
    metric: call_count
    deviation: 3.0
    window: immediate
  actions:
    - type: auto_pause_contract
    - type: create_multisig_proposal
    - type: notify_pagerduty

# Preflight 检查配置
preflight:
  max_messages: 20
  max_gas: 5000000
  blocked_functions:
    - "self_destruct"
    - "sudo"
  suspicious_patterns:
    - "reply"
    - "SubMsg"
    - "withdraw"
  risk_threshold: 0.7

自动响应脚本

pub fn handle_reentrancy_alert(
    deps: DepsMut,
    env: Env,
    alert: ReentrancyAlert,
) -> StdResult<Response> {
    // 1. 暂停攻击合约
    PAUSED_CONTRACTS.save(deps.storage, &alert.attacker_contract, &true)?;

    // 2. 暂停受害合约
    PAUSED_CONTRACTS.save(deps.storage, &alert.victim_contract, &true)?;

    // 3. 快照资金状态
    let snapshot = FundsSnapshot {
        block_height: env.block.height,
        timestamp: env.block.time.seconds(),
        attacker_balance: deps.querier.query_all_balances(&alert.attacker_contract)?,
        victim_balance: deps.querier.query_all_balances(&alert.victim_contract)?,
    };
    SNAPSHOTS.save(deps.storage, &env.block.height, &snapshot)?;

    // 4. 记录审计事件
    AUDIT_LOG.save(deps.storage, &AuditEntry {
        event_type: "reentrancy_auto_response".to_string(),
        details: serde_json::to_string(&alert).unwrap_or_default(),
        timestamp: env.block.time.seconds(),
    })?;

    Ok(Response::new()
        .add_attribute("action", "reentrancy_response")
        .add_attribute("attacker_paused", alert.attacker_contract)
        .add_attribute("victim_paused", alert.victim_contract)
        .add_attribute("snapshot_height", env.block.height.to_string()))
}

12.2 场景:预言机操纵自动响应

检测流程

  1. Indexer 捕获价格更新交易
  2. OracleManipulationDetector 对比链上价格和链下参考价
  3. 发现价格偏差 > 15%,触发 HIGH 告警
  4. 规则引擎匹配 oracle_manipulation 规则
  5. 自动暂停依赖该预言机的合约
  6. 创建多签提案冻结价格更新
  7. 通知预言机维护团队核查

响应代码

async function respondToOracleManipulation(
  alert: AlertEvent,
  responder: AutoPauseResponder,
  multisig: MultisigTrigger,
): Promise<void> {
  const affectedContracts = alert.context['affected_contracts'] as string[];

  // 暂停所有受影响合约
  const pauseResults = await responder.batchPause(
    affectedContracts,
    `Oracle manipulation detected: ${alert.message}`,
    alert.ruleId,
  );

  // 创建治理提案
  await multisig.createEmergencyProposal({
    id: `oracle-${Date.now()}`,
    title: '紧急: 预言机操纵攻击响应',
    description: alert.message,
    contractAddress: alert.contractAddress || '',
    actions: [{
      type: 'update_config',
      params: { oracle_paused: true, reason: alert.message },
    }],
    proposer: 'auto-responder',
    requiredVotes: 3,
    deadline: Date.now() + 3600000,
    status: 'active',
  });

  console.log(`[OracleResponse] ${pauseResults.filter(r => r.success).length}/${affectedContracts.length} contracts paused`);
}

12.3 场景:大规模 Rug Pull 预警

检测流程

  1. LiquidityMonitor 检测到某池移除 80% 流动性
  2. RugPullScorer 计算风险评分 0.85 (Critical)
  3. HoneypotAnalyzer 扫描合约发现 4 个蜜罐特征
  4. 综合告警触发,标记该合约为高风险
  5. 通知所有与该合约交互的用户
  6. 冻结合约中剩余资金
async function handleRugPullAlert(
  alert: AlertEvent,
): Promise<void> {
  const contract = alert.contractAddress!;

  // 1. 标记合约为高风险
  await markHighRiskContract(contract);

  // 2. 广播预警消息
  await broadcastWarning(`WARNING: High risk of rug pull detected on ${contract}`);

  // 3. 通知合约交互用户
  const users = await getContractUsers(contract, 1000);
  for (const user of users) {
    await notifyUser(user, `Warning: ${contract} flagged as potential rug pull`);
  }
}

12.4 案例总结

阶段 最佳实践
检测 多层检测:Mempool + 区块 + 状态 + 行为基线
分析 规则引擎 + ML 模型 + 威胁情报交叉验证
告警 分级告警 + 去重 + 抑制 + 升级策略
响应 自动暂停 + 多签提案 + 审计快照 + 用户通知
恢复 事件复盘 + 签名更新 + 规则优化 + 报告归档


本文档基于 MSG Chain 代码库核实的技术事实。
白皮书系统: https://msgchain.org/whitepaper/