链上安全监控与威胁检测指南
数据来源:MSG Chain 代码库核实
主网状态: No-Go — 当前 MSGChain 主网裁决为 No-Go,以下内容反映代码实际状态,不代表生产可用。
目录
- 链上威胁建模
- 实时交易监控架构
- 合约行为基线
- 告警规则引擎
- MEV 攻击检测
- 链上蜜罐 / Rug Pull 检测
- MSG Chain 安全监控栈
- 自动化响应
- 威胁情报
- 安全仪表盘设计与指标
- 合规监控
- 案例:监控 + 自动响应实战演练
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 上重入攻击特征:
- 攻击合约调用受害合约的
withdraw/transfer函数 - 受害合约使用
SubMsg发送资金但未先更新内部状态 - 攻击合约在
replyhandler 中再次调用withdraw - 循环利用直到 Gas 耗尽或达到区块 Gas 上限
检测指标:
| 指标 | 正常范围 | 异常信号 |
|---|---|---|
| 单区块内同一合约调用次数 | 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 上闪电贷攻击特征:
- 单笔交易中包含大量跨合约
SubMsg调用链 - 交易中间有价格操纵操作,交易结束前价格恢复
- 资金流动形成闭环(借出 → 操纵 → 获利 → 归还)
- 通常涉及 AMM 池、借贷合约和价格预言机
检测模型:
检测逻辑:
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函数缺少额外的权限检查sudo入口点被非验证者调用- 角色/权限存储可被篡改
- Agent API 的 Constitution 检查被绕过
检测规则:
| 权限操作 | 调用者 | 风险等级 |
|---|---|---|
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 ¤t_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 场景:预言机操纵自动响应
检测流程
- Indexer 捕获价格更新交易
- OracleManipulationDetector 对比链上价格和链下参考价
- 发现价格偏差 > 15%,触发 HIGH 告警
- 规则引擎匹配
oracle_manipulation规则 - 自动暂停依赖该预言机的合约
- 创建多签提案冻结价格更新
- 通知预言机维护团队核查
响应代码
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 预警
检测流程
- LiquidityMonitor 检测到某池移除 80% 流动性
- RugPullScorer 计算风险评分 0.85 (Critical)
- HoneypotAnalyzer 扫描合约发现 4 个蜜罐特征
- 综合告警触发,标记该合约为高风险
- 通知所有与该合约交互的用户
- 冻结合约中剩余资金
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/
