dApp Docs/链下预言机网络部署指南
Development reference. Not independently verified for production.

链下预言机网络部署指南

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

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


目录

  1. 预言机在区块链中的角色
  2. MSG Chain 预言机现状
  3. 预言机架构模式
  4. Chainlink 架构借鉴
  5. CosmWasm 预言机合约设计
  6. 链下节点部署
  7. 数据聚合策略
  8. MSG Chain Registry 中的预言机键位
  9. 安全考虑
  10. 与 DeFi 合约集成
  11. 监控与 SLA
  12. 当前 Stub 状态与未来路线图

1. 预言机在区块链中的角色

1.1 区块链的封闭性问题

区块链网络本质上是一个确定性封闭系统。共识引擎要求所有节点对同一组输入产生完全相同的输出。这一设计使得区块链天然无法主动访问外部世界——网络中的每个节点必须独立验证所有状态转换,而外部数据源无法提供这种"可复现的确定性"。

预言机(Oracle)正是为解决这一根本矛盾而生的基础设施组件。它的核心职责是将外部世界的可信数据桥接到区块链的确定性执行环境中。

1.2 价格喂价(Price Feed)

价格喂价是预言机最成熟、使用最广泛的应用场景。在 Cosmos 生态中,几乎每个 DeFi 协议都依赖价格数据来完成其核心逻辑:

借贷协议中的清算机制

借贷协议需要实时跟踪抵押品价值与借款价值的比率。当抵押品价格下跌到清算阈值以下时,协议必须触发清算。如果预言机提供的是过时或操纵过的价格,协议可能提前清算用户(造成用户损失)或延迟清算(造成坏账)。

// 伪代码示例:借贷协议中的清算检查
fn check_liquidation(position: Position, price: Decimal) -> Action {
    let collateral_value = position.collateral_amount * price;
    let loan_value = position.loan_amount;
    let ratio = collateral_value / loan_value;

    if ratio < LIQUIDATION_THRESHOLD {
        Action::Liquidate { position_id: position.id }
    } else {
        Action::Skip
    }
}

AMM 滑点保护

自动做市商(AMM)需要参考外部市场价格来设置合理的滑点保护参数。当链内价格与外部市场出现显著偏离时,AMM 容易遭受三明治攻击(sandwich attack)和闪电贷操纵。

合成资产与衍生品

合成资产协议需要跟踪标的资产的实时价格来铸造、赎回和清算合成资产头寸。每个合成资产对应一个或多个价格源。

稳定币锚定机制

算法稳定币依赖价格预言机来决定是扩张还是收缩供应量。当偏离锚定价格时,协议需要根据预言机报告的价格来执行套利激励机制。

1.3 随机数(Verifiable Random Function)

区块链是确定性的,这意味着单纯的链上随机数生成(如使用区块哈希)可能被验证者操纵。VRF(可验证随机函数)预言机提供了可验证的、不可预测的、无偏的随机数。

应用场景

VRF 的核心数学原理如下:

Prover 拥有私钥 sk 和公钥 pk
输入: alpha (seed)
输出: (beta, pi)
  beta = VRF_hash(sk, alpha)
  pi   = VRF_proof(sk, alpha)

Verifier 验证:
  VRF_verify(pk, alpha, beta, pi) -> true / false

CosmWasm 生态中,VRF 通常通过链下节点生成随机数并提交到链上,链上合约验证 BLS 签名或椭圆曲线证明来完成验证。

1.4 跨链数据(Cross-Chain Data)

跨链预言机负责将一条链上的状态证明(如 IBC 数据包确认、验证者集合变更、特定事件日志)传递到另一条链上。

IBC 预言机

尽管 IBC 本身提供了标准化跨链通信,但在以下场景中仍然需要预言机:

链外数据源上链

1.5 预言机的核心安全属性

无论何种类型,预言机都必须满足以下安全属性:

属性 描述 重要性
可用性 数据始终可访问,不会因单点故障而中断 高
完整性 数据在传输和存储过程中不被篡改 高
时效性 数据反映的是最新状态,而非陈旧数据 高
抗操纵性 单个实体无法影响最终报告值 极高
可验证性 数据消费者可以验证数据的来源和计算过程 高
经济安全性 操纵数据的成本远高于潜在收益 极高

2. MSG Chain 预言机现状

2.1 官方 Stub 状态

截至本文档编写时,MSG Chain 的 Oracle 功能模块处于 Stub(桩代码) 阶段。这意味着:

这是一个规划中的功能,而非已实现的特性。本文档的目的正是在基础设施尚未搭建的阶段,为开发者提供一份完整的架构蓝图与实施指南。

2.2 当前局限性

由于预言机模块处于 Stub 状态,MSG Chain 目前面临以下限制:

无法支持 DeFi 协议

任何需要外部价格数据的 DeFi 协议(借贷、AMM、合成资产、稳定币)都无法安全部署。没有价格预言机意味着:

没有可验证的随机数

链上游戏和 NFT 项目无法获得安全的随机数源。使用区块哈希作为随机数源在 MSG Chain 上同样存在被验证者操纵的风险(虽然该链使用 aBFT 共识,但区块 proposer 仍可在一定程度上影响哈希值)。

跨链数据依赖外部中继

没有原生预言机支持,跨链应用需要依赖第三方中继服务来桥接数据,这不仅引入信任假设,还增加了系统的复杂度和攻击面。

2.3 已有基础

尽管预言机模块尚未实现,MSG Chain 提供了一些可复用的基础设施:

CosmWasm 虚拟机

MSG Chain 完整支持 CosmWasm v1.x,这意味着预言机系统可以由链上智能合约和链下节点协同工作。合约负责聚合逻辑、数据存储和惩罚机制,链下节点负责从外部数据源获取数据并提交到链上。

Registry 模块

MSG Chain Registry 提供了键值对存储机制,可以将预言机节点信息、价格源元数据和合约地址注册到 canonical key 中。Registry 的查询接口支持链上和链下两种访问方式。

IBC 支持

完整的 IBC 支持使得 MSG Chain 可以通过跨链账户和 IBC 数据包与其他 Cosmos 链交换数据。预言机系统可以依赖 IBC 来获取其他链上的价格数据。

2.4 设计目标

在规划和实现 MSG Chain 预言机基础设施时,应遵循以下设计目标:

  1. 去中心化:不存在单点故障或单点控制
  2. 经济安全:攻击成本 > 攻击收益
  3. 模块化:价格源、聚合策略、惩罚逻辑均可插拔
  4. 与 CosmWasm 深度集成:合约接口遵循 CosmWasm 标准
  5. 低延迟:价格更新频率满足 DeFi 需求
  6. 低成本:Gas 消耗和数据更新成本在可接受范围内

3. 预言机架构模式

3.1 中心化预言机(Centralized Oracle)

架构描述

单个实体作为数据提供者,将外部数据签名后提交到链上。智能合约直接信任该单一数据源。

+------------------+       +------------------+
|  外部数据源       |       |                  |
|  (交易所 API)     +------>+  中心化预言机节点  |
|                   |       |                  |
+------------------+       +--------+---------+
                                     |
                                     | 签名提交
                                     |
                            +--------v---------+
                            |                  |
                            |  链上目标合约     |
                            |                  |
                            +------------------+

优点

缺点

适用场景

在 MSG Chain 上的实现

// 中心化预言机的精简实现
pub fn execute_submit_price(
    deps: DepsMut,
    info: MessageInfo,
    price_feed: PriceFeed,
) -> Result<Response, ContractError> {
    // 硬编码的授权提交者
    let authorized_sender = deps.api.addr_validate("msg1authorized...")?;

    if info.sender != authorized_sender {
        return Err(ContractError::Unauthorized {});
    }

    PRICES.save(deps.storage, price_feed.symbol, &price_feed)?;

    Ok(Response::new()
        .add_attribute("action", "price_update")
        .add_attribute("symbol", &price_feed.symbol)
        .add_attribute("price", price_feed.price.to_string()))
}

3.2 去中心化预言机网络(Decentralized Oracle Network)

架构描述

多个独立节点从多个数据源获取数据,各自提交到链上,由链上聚合合约将多个报告整合为一个最终值。

+------------------+       +------------------+
|  外部数据源 A     |       |  节点 1           |
+------------------+       +--------+---------+
                                    |
+------------------+       +--------v---------+       +------------------+
|  外部数据源 B     +------->  节点 2           +------>+  链上聚合合约    |
+------------------+       +--------+---------+       +--------+---------+
                                    |                          |
+------------------+       +--------v---------+       +--------v---------+
|  外部数据源 C     |       |  节点 3           |       |  消费方合约      |
+------------------+       +------------------+       +------------------+

优点

缺点

数据去重问题

当多个节点从不同数据源获取价格时,它们的报告值可能不完全一致。链上聚合合约需要处理这些差异:

// 节点报告结构体
pub struct NodeReport {
    pub node_id: String,
    pub symbol: String,
    pub price: Decimal,
    pub timestamp: u64,
    pub signature: Vec<u8>,
}

适用场景

3.3 基于 Stake 的 Oracle(Stake-based Oracle)

架构描述

在去中心化预言机网络的基础上增加**质押(Staking)和惩罚(Slashing)**机制。节点必须质押一定数量的系统代币才能参与数据提交。如果节点提交的数据被证明是恶意或有偏差的,其质押代币会被部分或全部罚没。

+-------------------+       +-------------------+
|                   |       |  质押要求          |
| 节点注册          +------>+  - 最小质押量      |
| - 申请人地址      |       |  - 锁定期          |
| - 质押数量        |       |  - 惩罚条件        |
| - 公钥            |       +-------------------+
+--------+----------+
         |
         v
+-------------------+
|  链上 Oracle 模块  |
|                   |
|  +-------------+   |
|  | 数据聚合     |   |
|  +-------------+   |
|  +-------------+   |
|  | 信誉评分     |   |
|  +-------------+   |
|  +-------------+   |
|  | 惩罚/奖励    |   |
|  +-------------+   |
+-------------------+

优点

缺点

质押量的确定

质押量需要精算设计。理论上,质押量应满足:

质押量 >= 预期收益 * 惩罚因子 / 攻击概率

在实践中,通常通过治理投票动态调整质押参数。

3.4 混合架构(Hybrid Architecture)

推荐方案

针对 MSG Chain 的实际情况(CosmWasm 支持、Registry 可用、模块化设计),我们推荐采用混合架构:

  1. 数据层:多个链下节点从多个数据源获取数据
  2. 传输层:节点使用 MSG Chain 交易将数据提交到链上
  3. 聚合层:链上 CosmWasm 合约执行数据聚合
  4. 分发层:聚合后的数据通过合约查询接口供其他合约使用
  5. 经济层:基于 Stake 的质押和惩罚机制保障安全性
+--------------------------------------------------+
|                    经济层                           |
|             Stake / Slash / Reward                |
+--------------------------------------------------+
|                    分发层                           |
|          Oracle Queries -> Consumer Contracts      |
+--------------------------------------------------+
|                    聚合层                           |
|          Median / TWAP / Outlier Removal           |
+--------------------------------------------------+
|                    传输层                           |
|         Signed Transactions / IBC Packets          |
+--------------------------------------------------+
|                    数据层                           |
|    CEX API    DEX API    WebSocket    REST         |
+--------------------------------------------------+

各层的 CosmWasm 合约对应关系

层级 合约/模块 功能
聚合层 msg_oracle.wasm 接收提交、聚合数据、存储结果
经济层 msg_oracle_staking.wasm 质押管理、惩罚执行、奖励分配
分发层 msg_oracle_consumer.wasm 消费示例、查询接口封装
Registry MSG Chain Registry 节点注册、元数据存储

Chainlink 是目前最成熟的去中心化预言机网络,其在以太坊上的架构为所有区块链预言机系统提供了参考模版。虽然 MSG Chain 基于 Cosmos SDK + CosmWasm,而非 Ethereum + Solidity,但 Chainlink 的许多设计理念可以直接借鉴。

4.2 去中心化 Oracle 网络(DON)

Chainlink 的去中心化预言机网络由三个角色构成:

Oracle 节点(Node Operators)

数据源适配器(Adapters)

聚合合约(Aggregator Contract)

4.3 聚合合约设计

Chainlink 的聚合合约(FluxAggregator / AccessControlledAggregator)是其架构的核心。以下是其关键设计点的 CosmWasm 适配版本:

报告轮次(Round)

每个价格源维护一个递增的轮次计数。每个轮次中,授权节点可以提交一次报告。当收集到的报告数量达到最小要求时,触发聚合。

#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct Round {
    /// 轮次 ID
    pub round_id: u64,
    /// 该轮次已提交报告的节点数
    pub observation_count: u32,
    /// 所需最少报告数
    pub min_answers: u32,
    /// 聚合后的最终结果
    pub answer: Option<Decimal>,
    /// 该轮次的起始时间
    pub started_at: Timestamp,
    /// 该轮次的结束时间(聚合完成时间)
    pub answered_at: Option<Timestamp>,
}

节点报告验证

合约需要验证报告的签名和节点身份:

pub fn execute_submit_report(
    deps: DepsMut,
    info: MessageInfo,
    round_id: u64,
    price: Decimal,
    signature: Vec<u8>,
) -> Result<Response, ContractError> {
    // 验证节点是否已注册
    let node = NODES.load(deps.storage, &info.sender)?;

    // 验证签名
    let message = encode_price_message(round_id, price);
    let valid = deps.api.secp256k1_verify(
        &message,
        &signature,
        &node.public_key,
    )?;
    if !valid {
        return Err(ContractError::InvalidSignature {});
    }

    // 检查该节点是否已在该轮提交
    let submitted = SUBMISSIONS.may_load(deps.storage, (round_id, &info.sender))?;
    if submitted.is_some() {
        return Err(ContractError::AlreadySubmitted {});
    }

    // 保存报告
    SUBMISSIONS.save(deps.storage, (round_id, &info.sender), &price)?;

    // 检查是否达到最小报告数
    let round = ROUNDS.load(deps.storage, round_id)?;
    if round.observation_count + 1 >= round.min_answers {
        execute_aggregation(deps, round_id)?;
    } else {
        ROUNDS.update(deps.storage, round_id, |mut r| -> Result<_, ContractError> {
            r.observation_count += 1;
            Ok(r)
        })?;
    }

    Ok(Response::new()
        .add_attribute("action", "report")
        .add_attribute("round_id", round_id.to_string())
        .add_attribute("node", info.sender.as_str())
        .add_attribute("price", price.to_string()))
}

Cron 触发机制

在以太坊上,Chainlink 使用 Keeper Network(原 Keepers)来触发聚合。在 MSG Chain 上,有以下几种触发方式:

  1. 交易触发:任何节点提交最终报告时自动触发聚合
  2. 定时触发:链下 cron 作业定期调用合约
  3. 条件触发:当价格波动超过阈值时触发
// 条件触发的简化逻辑
pub fn execute_check_and_trigger(
    deps: DepsMut,
    info: MessageInfo,
    symbol: &str,
) -> Result<Response, ContractError> {
    let current_price = PRICES.load(deps.storage, symbol)?;
    let last_price = LAST_PRICES.load(deps.storage, symbol)?;

    let deviation = (current_price - last_price).abs() / last_price;
    let threshold = THRESHOLDS.load(deps.storage, symbol)?;

    if deviation >= threshold {
        // 触发新的轮次
        let new_round = Round {
            round_id: NEXT_ROUND_ID.load(deps.storage)?,
            observation_count: 0,
            min_answers: MIN_ANSWERS.load(deps.storage)?,
            answer: None,
            started_at: deps.block.time,
            answered_at: None,
        };
        ROUNDS.save(deps.storage, new_round.round_id, &new_round)?;
    }

    Ok(Response::new().add_attribute("action", "check_trigger"))
}

4.4 信誉系统(Reputation System)

Chainlink 引入信誉系统来追踪节点的历史表现。信誉评分影响节点被选中的数据提交的概率,以及节点获得奖励的权重。

信誉指标

指标 描述 计算方式
参与率 节点参与报告轮次的比例 已参与轮次 / 总轮次
准确性 节点报告与聚合结果的偏差 1 - avg(
时效性 节点报告的及时程度 按时提交比例
质押量 节点质押的代币数量 当前质押量

CosmWasm 信誉合约

#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct Reputation {
    pub node: Addr,
    pub total_rounds: u64,
    pub participated_rounds: u64,
    pub accuracy_sum: Decimal,
    pub timeliness_score: Decimal,
    pub stake_amount: Uint128,
    pub last_updated: Timestamp,
}

impl Reputation {
    pub fn participation_rate(&self) -> Decimal {
        if self.total_rounds == 0 {
            return Decimal::zero();
        }
        Decimal::from_ratio(self.participated_rounds, self.total_rounds)
    }

    pub fn accuracy_score(&self) -> Decimal {
        if self.participated_rounds == 0 {
            return Decimal::zero();
        }
        self.accuracy_sum / Decimal::from_uint128(self.participated_rounds.into())
    }

    pub fn composite_score(&self) -> Decimal {
        // 综合评分:参与率 * 0.3 + 准确率 * 0.5 + 时效性 * 0.2
        self.participation_rate() * Decimal::percent(30)
            + self.accuracy_score() * Decimal::percent(50)
            + self.timeliness_score * Decimal::percent(20)
    }
}

4.5 OCR(Off-Chain Reporting)

Chainlink 的 OCR 协议是其最重要的改进之一。它通过链下聚合来大幅降低链上 Gas 成本:

OCR 工作流

  1. 选定的节点通过 P2P 网络交换观察值
  2. 节点在链下运行共识协议(基于网络通信)
  3. 达到共识后,一个节点将聚合结果和一组签名提交到链上
  4. 链上合约只需要验证聚合结果和签名即可,不需要分别处理每个节点的提交
传统模式(无 OCR):
节点1 -> submit(tx1)     \
节点2 -> submit(tx2)      >  N 次链上交易
节点3 -> submit(tx3)     /

OCR 模式:
节点1, 2, 3 P2P 通信 -> 链下聚合 -> 单一提交(tx_agg) -> 1 次链上交易

适配 MSG Chain

OCR 需要节点之间建立 P2P 网络。Cosmos SDK 的 Tendermint P2P 层可用于搭建节点间通信通道。具体实现方式:

// OCR 聚合结果链上验证
pub fn execute_submit_ocr_report(
    deps: DepsMut,
    info: MessageInfo,
    round_id: u64,
    median_price: Decimal,
    observations_count: u8,
    signers: Vec<Addr>,      // 参与签名的节点
    signatures: Vec<Vec<u8>>, // 各节点签名
) -> Result<Response, ContractError> {
    // 验证签名数是否达到阈值
    if signers.len() < OCR_MIN_SIGNERS {
        return Err(ContractError::InsufficientSignatures {});
    }

    // 逐个验证节点签名
    for (i, signer) in signers.iter().enumerate() {
        let node = NODES.load(deps.storage, signer)?;
        let message = encode_ocr_message(round_id, median_price, observations_count);
        let valid = deps.api.secp256k1_verify(
            &message,
            &signatures[i],
            &node.public_key,
        )?;
        if !valid {
            return Err(ContractError::InvalidSignature {});
        }
    }

    // 保存聚合结果
    let round = Round {
        round_id,
        observation_count: observations_count as u32,
        min_answers: 0,
        answer: Some(median_price),
        started_at: Timestamp::default(),
        answered_at: Some(deps.block.time),
    };
    ROUNDS.save(deps.storage, round_id, &round)?;
    PRICES.save(deps.storage, round_id_to_symbol(round_id), &median_price)?;

    Ok(Response::new()
        .add_attribute("action", "ocr_report")
        .add_attribute("round_id", round_id.to_string())
        .add_attribute("price", median_price.to_string())
        .add_attribute("signers", signers.len().to_string()))
}
Chainlink 特性 MSG Chain 适配方式 难度
FluxAggregator CosmWasm 聚合合约 中
OCR 协议 链下节点 P2P + 单次链上提交 高
信誉系统 链上信誉合约 低
Keeper 网络 链下 cron 或条件触发 低
数据源适配器 链下节点插件 中
质押/惩罚 CosmWasm staking 合约 中

5. CosmWasm 预言机合约设计

5.1 合约架构总览

MSG Chain 预言机系统由多个 CosmWasm 合约协同工作:

+-----------------------------------------------------------+
|                     Consumer Contracts                      |
|  (Lending, AMM, Stablecoin, Synthetic Assets)              |
+-----------------------------------------------------------+
          |                    |                    |
          v                    v                    v
+------------------+  +------------------+  +------------------+
|  Oracle Proxy    |  |  Oracle Proxy    |  |  Oracle Proxy    |
|  (查询接口封装)   |  |  (查询接口封装)   |  |  (查询接口封装)   |
+------------------+  +------------------+  +------------------+
          |                    |                    |
          +--------------------+--------------------+
                               |
                               v
+-----------------------------------------------------------+
|                  Oracle Aggregator Contract                |
|  - 接收节点报告                                             |
|  - 执行数据聚合                                             |
|  - 存储历史价格                                             |
|  - 管理轮次状态                                             |
+-----------------------------------------------------------+
          |
          v
+-----------------------------------------------------------+
|                  Oracle Staking Contract                    |
|  - 节点注册                                                |
|  - 质押管理                                                |
|  - 惩罚执行                                                |
|  - 奖励分配                                                |
+-----------------------------------------------------------+

5.2 数据源适配器(Data Source Adapter)

设计原理

每个链下节点内部运行多个数据源适配器。适配器负责:

  1. 连接到外部数据源(CEX API、DEX subgraph、REST API 等)
  2. 获取原始数据
  3. 将数据转换为标准格式
  4. 返回标准化的数值结果

适配器接口定义

// 适配器配置
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct AdapterConfig {
    /// 适配器名称
    pub name: String,
    /// 数据源 URL 模板
    pub url_template: String,
    /// HTTP 方法
    pub method: HttpMethod,
    /// 请求头
    pub headers: Vec<(String, String)>,
    /// JSON 路径表达式(用于提取数值)
    pub json_path: String,
    /// 数据源权重(用于加权聚合)
    pub weight: Option<Decimal>,
    /// 超时时间(毫秒)
    pub timeout_ms: u32,
}

#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub enum HttpMethod {
    Get,
    Post,
}

// 适配器返回的标准输出
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct AdapterResult {
    /// 适配器名称
    pub adapter_name: String,
    /// 提取的价格值
    pub price: Decimal,
    /// 获取数据时的时间戳
    pub observed_at: u64,
    /// 原始响应(仅用于调试)
    pub raw_response: Option<String>,
}

适配器配置示例

适配器 数据源 URL JSON 路径
binance_btc_usdt Binance API https://api.binance.com/api/v3/ticker/price?symbol=BTCUSDT price
coinbase_btc_usd Coinbase API https://api.coinbase.com/v2/prices/BTC-USD/spot data.amount
osmosis_btc Osmosis DEX 通过 LCD 查询 合约查询结果
msg_dex_btc MSG Chain DEX 链上查询 池子价格计算

5.3 聚合逻辑合约(Aggregation Contract)

聚合合约接口设计

#[cw_serde]
pub struct InstantiateMsg {
    /// 系统代币的 denom
    pub native_denom: String,
    /// 最小节点数
    pub min_nodes: u32,
    /// 每个轮次的最小报告数
    pub min_reports_per_round: u32,
    /// 价格偏差阈值(触发新轮次)
    pub deviation_threshold: Decimal,
    /// 心跳间隔(秒,触发新轮次)
    pub heartbeat_seconds: u64,
}

#[cw_serde]
pub enum ExecuteMsg {
    /// 提交价格报告
    SubmitReport {
        symbol: String,
        price: Decimal,
        round_id: u64,
    },
    /// 触发新轮次
    TriggerRound {
        symbol: String,
    },
    /// OCR 批量提交
    SubmitOcrReport {
        symbol: String,
        median_price: Decimal,
        observations_count: u8,
        signers: Vec<String>,
        signatures: Vec<Binary>,
    },
    /// 注册新的价格源
    RegisterPriceFeed {
        symbol: String,
        decimals: u8,
        min_answer: Decimal,
        max_answer: Decimal,
    },
}

#[cw_serde]
pub enum QueryMsg {
    /// 获取当前价格
    GetPrice {
        symbol: String,
    },
    /// 获取历史价格
    GetHistoricalPrice {
        symbol: String,
        round_id: u64,
    },
    /// 获取价格源配置
    GetPriceFeedConfig {
        symbol: String,
    },
    /// 获取最新轮次信息
    GetLatestRound {
        symbol: String,
    },
    /// 列出所有价格源
    ListPriceFeeds {},
}

// 查询响应
#[cw_serde]
pub struct PriceResponse {
    pub symbol: String,
    pub price: Decimal,
    pub decimals: u8,
    pub round_id: u64,
    pub timestamp: Timestamp,
    pub updated_at: Timestamp,
}

核心聚合逻辑实现

const CONFIG: Item<Config> = Item::new("config");
const PRICE_FEEDS: Map<String, PriceFeedConfig> = Map::new("price_feeds");
const PRICES: Map<String, PriceResponse> = Map::new("prices");
const ROUNDS: Map<(String, u64), Round> = Map::new("rounds");
const REPORTS: Map<(String, u64, Addr), Decimal> = Map::new("reports");
const NODE_COUNTS: Map<(String, u64), u32> = Map::new("node_counts");

pub fn execute_submit_report(
    deps: DepsMut,
    info: MessageInfo,
    symbol: String,
    price: Decimal,
    round_id: u64,
) -> Result<Response, ContractError> {
    // 验证节点是否已注册
    let node = NODES.load(deps.storage, &info.sender).map_err(|_| ContractError::NodeNotRegistered)?;
    if !node.active {
        return Err(ContractError::NodeInactive);
    }

    // 验证价格源是否存在
    let feed = PRICE_FEEDS.load(deps.storage, &symbol)?;
    if price < feed.min_answer || price > feed.max_answer {
        return Err(ContractError::PriceOutOfBounds);
    }

    // 验证轮次是否存在
    let round_key = (symbol.clone(), round_id);
    let mut round = ROUNDS.load(deps.storage, &round_key)?;

    // 检查是否已经提交过
    let report_key = (symbol.clone(), round_id, info.sender.clone());
    if REPORTS.has(deps.storage, &report_key) {
        return Err(ContractError::DuplicateReport);
    }

    // 存储报告
    REPORTS.save(deps.storage, &report_key, &price)?;

    // 更新轮次计数
    round.observation_count += 1;
    ROUNDS.save(deps.storage, &round_key, &round)?;

    // 检查是否达到最小报告数,触发聚合
    let config = CONFIG.load(deps.storage)?;
    if round.observation_count >= config.min_reports_per_round && round.answer.is_none() {
        let aggregated = execute_aggregation(deps.storage, &symbol, round_id)?;

        // 更新价格
        PRICES.save(deps.storage, &symbol, &PriceResponse {
            symbol: symbol.clone(),
            price: aggregated,
            decimals: feed.decimals,
            round_id,
            timestamp: round.started_at,
            updated_at: deps.block.time,
        })?;

        // 更新轮次状态
        round.answer = Some(aggregated);
        round.answered_at = Some(deps.block.time);
        ROUNDS.save(deps.storage, &round_key, &round)?;

        emit_price_update_event(deps, &symbol, aggregated, round_id);
    }

    Ok(Response::new()
        .add_attribute("method", "submit_report")
        .add_attribute("symbol", &symbol)
        .add_attribute("round_id", round_id.to_string())
        .add_attribute("node", info.sender.as_str()))
}

fn execute_aggregation(
    storage: &mut dyn Storage,
    symbol: &str,
    round_id: u64,
) -> Result<Decimal, ContractError> {
    // 收集该轮次所有报告
    let reports: Vec<Decimal> = REPORTS
        .range(storage, None, None, Order::Ascending)
        .filter(|r| r.0 .0 == symbol && r.0 .1 == round_id)
        .map(|r| r.1)
        .collect();

    if reports.is_empty() {
        return Err(ContractError::NoReports);
    }

    // 使用中位数聚合
    let mut sorted = reports.clone();
    sorted.sort();
    let median = if sorted.len() % 2 == 0 {
        let mid = sorted.len() / 2;
        (sorted[mid - 1] + sorted[mid]) / Decimal::from_ratio(2u64, 1u64)
    } else {
        sorted[sorted.len() / 2]
    };

    Ok(median)
}

5.4 惩罚机制合约(Slashing Contract)

惩罚条件

条件 惩罚比例 触发方式
提交明显异常价格(与中位数偏差 > 20%) 1% 质押 自动检测
持续提交异常价格(连续 N 轮) 5% 质押 自动检测
双签(同一轮次提交两个不同价格) 10% 质押 自动检测
节点持续离线(超过心跳间隔 N 倍) 0.5% 质押 / 轮 自动检测
治理裁定的恶意行为 治理决定 治理投票

惩罚合约实现

#[cw_serde]
pub enum ExecuteMsg {
    RegisterNode {
        pub_key: Binary,
        initial_stake: Coin,
    },
    UnbondStake {
        amount: Uint128,
    },
    WithdrawUnbonded {},
    /// 提交惩罚证据
    SubmitSlashEvidence {
        node: String,
        evidence_type: SlashType,
        evidence_data: Binary,
    },
}

#[cw_serde]
pub enum SlashType {
    OutlierReport,
    DoubleSign,
    ProlongedDowntime,
    GovernanceAction,
}

pub fn execute_slash(
    deps: DepsMut,
    info: MessageInfo,
    node_addr: Addr,
    evidence_type: SlashType,
) -> Result<Response, ContractError> {
    let mut node = NODES.load(deps.storage, &node_addr)?;
    let slash_percentage = match evidence_type {
        SlashType::OutlierReport => Decimal::percent(1),
        SlashType::DoubleSign => Decimal::percent(10),
        SlashType::ProlongedDowntime => Decimal::percent(1),
        SlashType::GovernanceAction => {
            // 治理投票决定的具体比例
            return Err(ContractError::GovernanceRequired);
        }
    };

    let slashed_amount = (node.stake_amount * slash_percentage).ceil();
    node.stake_amount = node.stake_amount.checked_sub(slashed_amount)?;
    node.slash_count += 1;

    if node.stake_amount < node.min_stake {
        node.active = false;
    }

    NODES.save(deps.storage, &node_addr, &node)?;

    // 将罚没代币转入社区池或销毁
    let burn_msg = BankMsg::Burn {
        amount: vec![Coin {
            denom: deps.querier.query_balance(&node_addr, &CONFIG.load(deps.storage)?.native_denom)?.denom,
            amount: slashed_amount,
        }],
    };

    Ok(Response::new()
        .add_message(burn_msg)
        .add_attribute("action", "slash")
        .add_attribute("node", node_addr.as_str())
        .add_attribute("amount", slashed_amount.to_string())
        .add_attribute("reason", format!("{:?}", evidence_type)))
}

5.5 合约升级机制

预言机合约在部署后可能需要进行功能升级。CosmWasm 支持通过 MigrateMsg 实现合约迁移:

#[cw_serde]
pub enum MigrateMsg {
    /// 升级到新版本的聚合逻辑
    UpgradeAggregation {
        new_min_reports: u32,
        new_deviation_threshold: Decimal,
    },
    /// 升级惩罚参数
    UpgradeSlashing {
        new_outlier_percentage: Decimal,
        new_double_sign_percentage: Decimal,
    },
    /// 添加新的聚合策略
    AddAggregationStrategy {
        strategy_name: String,
        strategy_code_id: u64,
    },
}

5.6 合约存储设计

存储键前缀约定

前缀 存储内容 Map 类型
config 合约配置 Item
price_feeds 价格源配置 Map<String, PriceFeedConfig>
prices 当前价格 Map<String, PriceResponse>
rounds 轮次数据 Map<(String, u64), Round>
reports 节点报告 Map<(String, u64, Addr), Decimal>
nodes 节点信息 Map<Addr, NodeInfo>
reputations 信誉评分 Map<Addr, Reputation>

Gas 优化建议

  1. 使用 prefix 而不是 Map 的 full key 来减少序列化开销
  2. 批量更新时使用 save 而不是逐个 update
  3. 合理设置 Item 和 Map 的选择——不经常变动的配置用 Item
  4. 查询接口使用分页(pagination)来避免单次查询数据量过大
// Gas 优化的批量更新示例
pub fn batch_update_prices(
    deps: DepsMut,
    prices: Vec<(String, Decimal)>,
) -> Result<Response, ContractError> {
    let mut response = Response::new();
    for (symbol, price) in prices {
        PRICES.save(deps.storage, &symbol, &price)?;
        response = response.add_attribute(format!("price_{}", symbol), price.to_string());
    }
    Ok(response)
}

6. 链下节点部署

6.1 节点角色定义

链下 Oracle 节点是预言机网络的核心执行者。每个节点是一个独立的服务进程,负责:

  1. 数据获取:从多个外部数据源获取原始数据
  2. 数据验证:对获取的数据进行基本验证(格式、范围、时效性)
  3. 签名提交:使用节点私钥对数据进行签名,构造并广播交易
  4. Gas 管理:管理账户余额,确保有足够 Gas 提交交易
  5. 监控自检:定期上报健康状态

6.2 节点架构设计

+-------------------------------------------------------------+
|                    Oracle Node Service                        |
|                                                              |
|  +-------------------+  +-------------------+                 |
|  |   Data Fetcher    |  |   Aggregator      |                 |
|  |                   |  |   (可选链下聚合)    |                 |
|  |  - HTTP Client    |  |  - Median Calc     |                 |
|  |  - WebSocket      |  |  - Outlier Detect  |                 |
|  |  - gRPC Client    |  |  - TWAP Calc       |                 |
|  +--------+----------+  +--------+----------+                 |
|           |                      |                            |
|           +----------+-----------+                            |
|                      |                                        |
|             +--------v-----------+                            |
|             |   Report Builder   |                            |
|             |  - 数据编码         |                            |
|             |  - 签名生成         |                            |
|             +--------+-----------+                            |
|                      |                                        |
|             +--------v-----------+                            |
|             |   Tx Broadcaster   |                            |
|             |  - 交易构造         |                            |
|             |  - Gas 估算         |                            |
|             |  - 重试逻辑         |                            |
|             |  - Nonce 管理       |                            |
|             +--------+-----------+                            |
|                      |                                        |
+-------------------------------------------------------------+

6.3 节点注册流程

节点正式加入网络前需要完成注册。注册在链上进行,由 Staking 合约管理:

// 注册信息
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct NodeRegistration {
    /// 节点操作者地址
    pub operator_address: Addr,
    /// 节点公钥(用于签名验证)
    pub public_key: Binary,
    /// 节点描述信息
    pub description: String,
    /// 支持的 feeds 列表
    pub supported_feeds: Vec<String>,
    /// 初始质押量
    pub initial_stake: Coin,
    /// 佣金比例
    pub commission_rate: Decimal,
}

// 链上节点信息
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct NodeInfo {
    pub operator_address: Addr,
    pub public_key: Binary,
    pub description: String,
    pub supported_feeds: Vec<String>,
    pub stake_amount: Uint128,
    pub min_stake: Uint128,
    pub active: bool,
    pub commission_rate: Decimal,
    pub registered_at: Timestamp,
    pub last_heartbeat: Timestamp,
    pub total_rewards: Uint128,
    pub slash_count: u64,
}

注册步骤

  1. 节点操作者生成一个新的 Cosmos 账户(或使用已有账户)
  2. 节点操作者构造注册交易,指定质押数量和公钥
  3. 交易被广播到 MSG Chain,由 Staking 合约处理
  4. 合约验证最小质押量,锁定质押代币,将节点标记为活跃
  5. 节点启动服务,开始获取数据并提交报告

6.4 数据获取层(Data Fetcher)

多数据源获取

链下节点需要从多个独立数据源获取数据,以减少单点依赖。

CEX API 适配器

// 从 Binance 获取 BTC/USDT 价格
async fn fetch_binance_price() -> Result<AdapterResult, NodeError> {
    let url = "https://api.binance.com/api/v3/ticker/price?symbol=BTCUSDT";
    let client = Client::new();
    let resp = client.get(url)
        .header("User-Agent", "msg-oracle-node/1.0")
        .send()
        .await?
        .json::<HashMap<String, String>>()
        .await?;

    let price_str = resp.get("price").ok_or(NodeError::MissingField("price"))?;
    let price = Decimal::from_str(price_str)?;

    Ok(AdapterResult {
        adapter_name: "binance".to_string(),
        price,
        observed_at: current_timestamp(),
        raw_response: Some(format!("{:?}", resp)),
    })
}

DEX 链上查询适配器

// 从 Osmosis 链上查询池子价格
async fn fetch_osmosis_price() -> Result<AdapterResult, NodeError> {
    let lcd_endpoint = "https://lcd.osmosis.zone";
    let pool_id = "1"; // BTC/OSMO 池
    let query_path = format!("/osmosis/gamm/v1beta1/pools/{}", pool_id);

    let client = Client::new();
    let resp = client.get(format!("{}{}", lcd_endpoint, query_path))
        .send()
        .await?
        .json::<serde_json::Value>()
        .await?;

    // 从池子资产计算价格
    let assets = resp["poolAssets"].as_array().ok_or(NodeError::MissingField("poolAssets"))?;
    let price = calculate_price_from_pool(assets)?;

    Ok(AdapterResult {
        adapter_name: "osmosis".to_string(),
        price,
        observed_at: current_timestamp(),
        raw_response: None,
    })
}

数据验证

获取到原始数据后,节点需要进行以下验证:

fn validate_price_data(
    result: &AdapterResult,
    config: &PriceValidationConfig,
) -> Result<(), NodeError> {
    // 1. 检查是否为有效正数
    if result.price <= Decimal::zero() {
        return Err(NodeError::InvalidPrice("price must be positive"));
    }

    // 2. 检查是否在合理范围内
    if result.price < config.min_price || result.price > config.max_price {
        return Err(NodeError::PriceOutOfRange {
            price: result.price,
            min: config.min_price,
            max: config.max_price,
        });
    }

    // 3. 检查与上一轮相比的偏差
    let deviation = (result.price - config.last_price).abs() / config.last_price;
    if deviation > config.max_deviation_per_update {
        return Err(NodeError::ExcessiveDeviation {
            price: result.price,
            last: config.last_price,
            deviation,
        });
    }

    // 4. 检查数据时效性
    let now = current_timestamp();
    if now - result.observed_at > config.max_age_seconds * 1000 {
        return Err(NodeError::StaleData {
            observed_at: result.observed_at,
            now,
        });
    }

    Ok(())
}

6.5 签名与交易提交

签名流程

节点使用私钥对报告数据进行签名,以确保数据的来源可验证:

fn build_and_sign_report(
    node_key: &PrivateKey,
    contract_addr: &str,
    symbol: &str,
    price: Decimal,
    round_id: u64,
    chain_id: &str,
    account_number: u64,
    sequence: u64,
) -> Result<Vec<u8>, NodeError> {
    // 构造签名的消息
    let report_data = ReportData {
        contract: contract_addr.to_string(),
        symbol: symbol.to_string(),
        price: price.to_string(),
        round_id,
        chain_id: chain_id.to_string(),
    };

    // 序列化并哈希
    let encoded = serde_json::to_vec(&report_data)?;
    let hash = sha256(&encoded);

    // 使用 secp256k1 签名
    let signature = node_key.sign(&hash)?;

    // 构造 Cosmos 交易
    let tx = CosmosTx {
        msgs: vec![CosmosMsg::Wasm(WasmMsg::Execute {
            contract_addr: contract_addr.to_string(),
            msg: to_binary(&ExecuteMsg::SubmitReport {
                symbol: symbol.to_string(),
                price,
                round_id,
            })?,
            funds: vec![],
        })],
        fee: Fee {
            amount: vec![Coin {
                denom: "umsg".to_string(),
                amount: Uint128::from(5000u128),
            }],
            gas: 300_000u64,
        },
        signatures: vec![Signature {
            pub_key: node_key.public_key(),
            signature: signature.clone(),
        }],
        memo: format!("oracle-{}", symbol),
        chain_id: chain_id.to_string(),
        account_number,
        sequence,
    };

    // 序列化交易
    let tx_bytes = tx.to_bytes()?;
    Ok(tx_bytes)
}

Gas 管理策略

Gas 管理是链下节点运行的关键挑战。节点需要平衡以下因素:

  1. Gas 价格预估:参考网络当前 Gas 价格,设置适当的 Gas 价格
  2. 重试策略:交易失败时的重试逻辑
  3. Nonce 管理:正确管理账户 nonce,避免交易冲突
  4. 余额监控:确保账户余额充足
struct GasManager {
    client: HttpClient,
    gas_price_multiplier: f64,
    max_retries: u32,
    retry_delay_ms: u64,
}

impl GasManager {
    /// 估算交易 Gas 并提交
    async fn estimate_and_submit(&self, tx: CosmosTx, account: &Account) -> Result<TxResponse, NodeError> {
        // 预估 Gas
        let gas_estimate = self.client.estimate_gas(&tx).await?;

        // Gas 价格 = 基础价格 * 乘数
        let base_gas_price = self.client.get_gas_price().await?;
        let gas_price = base_gas_price * self.gas_price_multiplier;

        // 设置 Gas
        let final_tx = tx.with_gas(gas_estimate, gas_price);

        // 提交交易
        let mut last_error = None;
        for attempt in 0..self.max_retries {
            match self.client.broadcast_tx(&final_tx).await {
                Ok(resp) => return Ok(resp),
                Err(e) => {
                    last_error = Some(e);
                    if attempt < self.max_retries - 1 {
                        sleep(Duration::from_millis(self.retry_delay_ms)).await;
                        // 更新 nonce
                        account.refresh_sequence().await?;
                    }
                }
            }
        }

        Err(NodeError::TxFailed(format!("{:?}", last_error)))
    }
}

6.6 节点配置

配置文件示例

# oracle-node-config.yaml
node:
  operator_address: "msg1operatoraddress..."
  private_key_path: "/etc/oracle/keys/node_secp256k1.key"
  chain_id: "msg-chain-1"
  grpc_endpoint: "https://grpc.msgchain.org:443"
  rpc_endpoint: "https://rpc.msgchain.org:443"
  rest_endpoint: "https://rest.msgchain.org"

contracts:
  oracle_aggregator: "msg1oracleaggregator..."
  oracle_staking: "msg1oraclestaking..."

feeds:
  - symbol: "BTC/USD"
    decimals: 8
    heartbeat_seconds: 30
    deviation_threshold: 0.005
    min_price: 10000
    max_price: 200000
    sources:
      - name: "binance"
        url: "https://api.binance.com/api/v3/ticker/price?symbol=BTCUSDT"
        json_path: "price"
        weight: 0.25
      - name: "coinbase"
        url: "https://api.coinbase.com/v2/prices/BTC-USD/spot"
        json_path: "data.amount"
        weight: 0.25
      - name: "kraken"
        url: "https://api.kraken.com/0/public/Ticker?pair=XBTUSD"
        json_path: "result.XXBTZUSD.c.0"
        weight: 0.25
      - name: "okx"
        url: "https://www.okx.com/api/v5/market/ticker?instId=BTC-USDT"
        json_path: "data.0.last"
        weight: 0.25

  - symbol: "ETH/USD"
    decimals: 8
    heartbeat_seconds: 30
    deviation_threshold: 0.005
    min_price: 500
    max_price: 20000
    sources:
      - name: "binance"
        url: "https://api.binance.com/api/v3/ticker/price?symbol=ETHUSDT"
        json_path: "price"
        weight: 0.25
      - name: "coinbase"
        url: "https://api.coinbase.com/v2/prices/ETH-USD/spot"
        json_path: "data.amount"
        weight: 0.25
      - name: "kraken"
        url: "https://api.kraken.com/0/public/Ticker?pair=ETHUSD"
        json_path: "result.XETHZUSD.c.0"
        weight: 0.25
      - name: "chainlink_eth_usd"
        url: "https://min-api.cryptocompare.com/data/pricemultifull?fsyms=ETH&tsyms=USD"
        json_path: "RAW.ETH.USD.PRICE"
        weight: 0.25

gas:
  gas_price: "1000000000umsg"
  gas_multiplier: 1.5
  max_gas: 500000
  retry_attempts: 3
  retry_delay_ms: 1000

logging:
  level: "info"
  format: "json"
  output: "/var/log/oracle-node.log"

monitoring:
  metrics_port: 9090
  health_check_port: 8080
  pprof_enabled: false

6.7 节点操作

启动命令

# 启动 Oracle 节点
oracle-node --config /etc/oracle/oracle-node-config.yaml

# 启动前检查
oracle-node validate-config --config /etc/oracle/oracle-node-config.yaml

# 测试数据源连接
oracle-node test-source --config /etc/oracle/oracle-node-config.yaml --symbol "BTC/USD"

# 检查节点状态
oracle-node status --operator "msg1operatoraddress..."

Docker 部署

FROM rust:1.70 as builder
WORKDIR /app
COPY . .
RUN cargo build --release --bin oracle-node

FROM debian:bullseye-slim
RUN apt-get update && apt-get install -y ca-certificates tzdata && rm -rf /var/lib/apt/lists/*
COPY --from=builder /app/target/release/oracle-node /usr/local/bin/
COPY config.yaml /etc/oracle/config.yaml
ENTRYPOINT ["oracle-node", "--config", "/etc/oracle/config.yaml"]
# docker-compose.yml
version: '3.8'
services:
  oracle-node-1:
    image: msgchain/oracle-node:latest
    container_name: msg-oracle-1
    volumes:
      - ./config/node1.yaml:/etc/oracle/config.yaml
      - ./keys/node1:/etc/oracle/keys
    ports:
      - "9090:9090"
      - "8080:8080"
    environment:
      - RUST_LOG=info
    restart: unless-stopped

  oracle-node-2:
    image: msgchain/oracle-node:latest
    container_name: msg-oracle-2
    volumes:
      - ./config/node2.yaml:/etc/oracle/config.yaml
      - ./keys/node2:/etc/oracle/keys
    ports:
      - "9091:9090"
      - "8081:8080"
    environment:
      - RUST_LOG=info
    restart: unless-stopped

  oracle-node-3:
    image: msgchain/oracle-node:latest
    container_name: msg-oracle-3
    volumes:
      - ./config/node3.yaml:/etc/oracle/config.yaml
      - ./keys/node3:/etc/oracle/keys
    ports:
      - "9092:9090"
      - "8082:8080"
    environment:
      - RUST_LOG=info
    restart: unless-stopped

6.8 节点安全

私钥管理

推荐方案:
+------------------------------------------------------------+
|                     HSM / KMS                                |
|  (硬件安全模块 或 云 KMS,如 AWS KMS / Azure Key Vault)      |
+------------------------------------------------------------+
                           |
                           v
+------------------------------------------------------------+
|                     Oracle Node                              |
|  - 通过 HSM API 签名,私钥永不离开硬件                        |
|  - 仅持有公钥和 HSM 访问凭证                                  |
+------------------------------------------------------------+

最低要求实践

  1. 私钥加密存储在磁盘上,使用强密码保护
  2. 节点运行环境与私钥存储分离
  3. 定期轮换节点密钥
  4. 监控账户活动,检测异常交易

7. 数据聚合策略

7.1 聚合策略概述

聚合策略决定了如何将多个节点的报告合并为一个权威价格。选择合适的聚合策略对预言机的安全性和准确性至关重要。

策略 抗操纵性 Gas 成本 新鲜度 实现复杂度
中位数 高 低 中 低
均值 低 低 中 低
截断均值 高 低 中 低
TWAP 高 中 低 中
VWAP 中 高 中 高
指数移动平均 中 低 高 低

7.2 中位数聚合(Median)

原理

将 N 个报告值排序后取中间值(N 为奇数时取正中间,N 为偶数时取中间两个的平均值)。

优点

缺点

pub fn compute_median(prices: &[Decimal]) -> Decimal {
    let mut sorted = prices.to_vec();
    sorted.sort();
    let len = sorted.len();
    if len == 0 {
        return Decimal::zero();
    }
    if len % 2 == 0 {
        let mid = len / 2;
        (sorted[mid - 1] + sorted[mid]) / Decimal::from_ratio(2u64, 1u64)
    } else {
        sorted[len / 2]
    }
}

7.3 TWAP(时间加权平均价格)

原理

TWAP 计算一段时间内价格的时间加权平均值。TWAP 天然抵抗瞬时价格操纵,因为攻击者需要在多个连续轮次中维持操纵价格。

计算方法

TWAP = Σ(price_i * interval_i) / Σ(interval_i)

其中:
  price_i    = 第 i 个时间区间的价格
  interval_i = 第 i 个时间区间的持续时间

对预言机节点的存储要求

#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct TwapState {
    /// 计算起始时间
    pub start_time: Timestamp,
    /// 累计价格 * 时间
    pub cumulative: Decimal,
    /// 累计时间
    pub total_seconds: u64,
    /// 更新次数
    pub updates: u64,
}

pub fn update_twap(
    state: &mut TwapState,
    new_price: Decimal,
    current_time: Timestamp,
) -> Decimal {
    let elapsed = (current_time - state.start_time).seconds();
    if elapsed == 0 {
        return new_price;
    }

    // 累计 价格 * 时间
    state.cumulative += new_price * Decimal::from_ratio(elapsed, 1u64);
    state.total_seconds += elapsed;
    state.start_time = current_time;
    state.updates += 1;

    state.cumulative / Decimal::from_ratio(state.total_seconds, 1u64)
}

参数建议

场景 TWAP 窗口
高流动性资产(BTC, ETH) 1-5 分钟
中等流动性资产 5-15 分钟
低流动性资产 15-60 分钟
稳定币锚定检查 1-6 小时

7.4 VWAP(交易量加权平均价格)

原理

VWAP 使用交易量作为权重来计算加权平均价格。交易量越大的交易所,其价格权重越高。

计算方法

VWAP = Σ(price_i * volume_i) / Σ(volume_i)

限制

VWAP 需要节点从数据源获取交易量数据,增加了数据获取复杂度。部分交易所的交易量可能存在虚增问题(刷量)。

#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct VolumeRecord {
    pub exchange: String,
    pub price: Decimal,
    pub volume_24h: Decimal,
}

pub fn compute_vwap(records: &[VolumeRecord]) -> Decimal {
    let total_volume: Decimal = records.iter().map(|r| r.volume_24h).sum();
    if total_volume.is_zero() {
        return Decimal::zero();
    }

    let weighted_sum: Decimal = records
        .iter()
        .map(|r| r.price * r.volume_24h)
        .sum();

    weighted_sum / total_volume
}

7.5 异常值剔除(Outlier Removal)

原理

在聚合之前,先识别并剔除偏离主流价格的异常报告值。这可以防止单个恶意节点(或少数合谋节点)影响聚合结果。

实现算法

Z-Score 方法

pub fn remove_outliers_zscore(
    prices: &mut Vec<Decimal>,
    threshold: f64,
) -> Vec<Decimal> {
    if prices.len() < 3 {
        return prices.clone();
    }

    let mean = prices.iter().sum::<Decimal>()
        / Decimal::from_uint128(prices.len() as u128);
    let variance = prices.iter()
        .map(|p| (p - mean).pow(2))
        .sum::<Decimal>()
        / Decimal::from_uint128(prices.len() as u128);
    let std_dev = variance.sqrt();

    if std_dev.is_zero() {
        return prices.clone();
    }

    prices.retain(|p| {
        let z_score = ((p - mean) / std_dev).abs();
        z_score < Decimal::from_str(&threshold.to_string()).unwrap()
    });

    prices.clone()
}

IQR(四分位距)方法

pub fn remove_outliers_iqr(prices: &[Decimal]) -> Vec<Decimal> {
    let mut sorted = prices.to_vec();
    sorted.sort();

    let n = sorted.len();
    if n < 4 {
        return sorted;
    }

    // 计算 Q1 和 Q3
    let q1_index = n / 4;
    let q3_index = 3 * n / 4;
    let q1 = sorted[q1_index];
    let q3 = sorted[q3_index];
    let iqr = q3 - q1;

    // 定义边界: Q1 - 1.5*IQR, Q3 + 1.5*IQR
    let lower_bound = q1 - iqr * Decimal::from_str("1.5").unwrap();
    let upper_bound = q3 + iqr * Decimal::from_str("1.5").unwrap();

    sorted.into_iter()
        .filter(|p| *p >= lower_bound && *p <= upper_bound)
        .collect()
}

偏差百分比方法

pub fn remove_outliers_deviation(
    prices: &[Decimal],
    max_deviation_pct: Decimal,
) -> Vec<(usize, Decimal)> {
    if prices.is_empty() {
        return vec![];
    }

    // 计算初始中位数
    let mut sorted = prices.to_vec();
    sorted.sort();
    let median = sorted[sorted.len() / 2];

    // 剔除与中位数偏差超过阈值的值
    prices.iter().enumerate()
        .filter(|(_, p)| {
            let deviation = (p - median).abs() / median;
            deviation <= max_deviation_pct
        })
        .map(|(i, p)| (i, *p))
        .collect()
}

7.6 组合策略

推荐在实际部署中采用多层组合策略:

原始报告值
       |
       v
+------------------+
| 阶段 1:          |
| 异常值剔除       |  (基于 IQR 或偏差百分比)
+------------------+
       |
       v
+------------------+
| 阶段 2:          |
| 时间加权         |  (可选 TWAP)
+------------------+
       |
       v
+------------------+
| 阶段 3:          |
| 中位数聚合       |  (或截断均值)
+------------------+
       |
       v
   最终价格
pub fn compute_final_price(
    reports: &[NodeReport],
    twap_state: &mut TwapState,
    config: &AggregationConfig,
) -> Result<Decimal, ContractError> {
    // 阶段 1: 提取价格值
    let mut prices: Vec<Decimal> = reports.iter().map(|r| r.price).collect();

    // 阶段 2: 异常值剔除
    prices = remove_outliers_deviation(&prices, config.max_deviation_pct);

    if prices.len() < config.min_valid_prices as usize {
        return Err(ContractError::InsufficientValidPrices);
    }

    // 阶段 3: 中位数聚合
    let median_price = compute_median(&prices);

    // 阶段 4: 可选 TWAP 平滑
    if config.use_twap {
        Ok(update_twap(twap_state, median_price, config.current_time))
    } else {
        Ok(median_price)
    }
}

7.7 价格新鲜度检查

聚合合约还需要确保价格不是过时的:

pub fn enforce_freshness(
    report_timestamps: &[u64],
    max_age_seconds: u64,
    current_time: u64,
) -> Result<(), ContractError> {
    let min_timestamp = current_time.saturating_sub(max_age_seconds * 1000);

    // 检查足够多的报告是新鲜的
    let fresh_count = report_timestamps
        .iter()
        .filter(|&&ts| ts >= min_timestamp)
        .count();

    if fresh_count < report_timestamps.len() / 2 {
        return Err(ContractError::StaleReports {
            fresh: fresh_count,
            total: report_timestamps.len(),
        });
    }

    Ok(())
}

8. MSG Chain Registry 中的预言机键位

8.1 MSG Chain Registry 概述

MSG Chain Registry 是基于 Cosmos SDK KVStore 构建的键值对存储系统。它允许将任意结构化数据注册到 canonical key 中,供链上和链下查询。预言机系统可以利用 Registry 来存储:

8.2 Registry 键位设计

命名空间约定

所有预言机关联的 Registry 键位使用 oracle/ 作为顶层命名空间:

oracle/<category>/<identifier>

推荐的键位结构

oracle/
├── feeds/
│   ├── btc-usd          -> PriceFeedMetadata
│   ├── eth-usd          -> PriceFeedMetadata
│   ├── msg-usd          -> PriceFeedMetadata
│   └── atom-usd         -> PriceFeedMetadata
├── nodes/
│   ├── msg1node1...     -> NodeRegistration
│   ├── msg1node2...     -> NodeRegistration
│   └── msg1node3...     -> NodeRegistration
├── contracts/
│   ├── aggregator       -> ContractAddress
│   ├── staking          -> ContractAddress
│   └── proxy            -> ContractAddress
├── config/
│   ├── global           -> OracleGlobalConfig
│   ├── slash_params     -> SlashParameters
│   └── reward_params    -> RewardParameters
├── status/
│   ├── heartbeat        -> HeartbeatStatus
│   └── rounds           -> LastRoundInfo
└── version              -> OracleVersion

8.3 Registry 数据模型

PriceFeedMetadata

#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct PriceFeedMetadata {
    /// 价格源符号
    pub symbol: String,
    /// 显示名称
    pub name: String,
    /// 精度
    pub decimals: u8,
    /// 心跳间隔(秒)
    pub heartbeat_seconds: u64,
    /// 偏差阈值
    pub deviation_threshold: Decimal,
    /// 最小价格
    pub min_price: Decimal,
    /// 最大价格
    pub max_price: Decimal,
    /// 关联的聚合合约地址
    pub aggregator_contract: String,
    /// 是否激活
    pub active: bool,
    /// 管理多签地址
    pub admin: String,
}

NodeRegistration

#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct NodeRegistration {
    /// 操作者地址
    pub operator_address: String,
    /// 节点公钥
    pub public_key_hex: String,
    /// 节点名称
    pub name: String,
    /// 描述
    pub description: String,
    /// 支持的 feeds
    pub supported_feeds: Vec<String>,
    /// 质押量
    pub stake_amount: String,
    /// 最小质押量
    pub min_stake: String,
    /// 是否活跃
    pub active: bool,
    /// 注册时间戳
    pub registered_at: String,
    /// 最后心跳
    pub last_heartbeat: String,
    /// 端点 URL(链下)
    pub endpoint_url: Option<String>,
}

OracleGlobalConfig

#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct OracleGlobalConfig {
    /// 最小活跃节点数
    pub min_active_nodes: u32,
    /// 每个轮次最小报告数
    pub min_reports_per_round: u32,
    /// 每个轮次最大报告数
    pub max_reports_per_round: u32,
    /// 默认心跳间隔(秒)
    pub default_heartbeat: u64,
    /// 默认偏差阈值
    pub default_deviation_threshold: Decimal,
    /// 奖励比例
    pub reward_rate: Decimal,
    /// 最大惩罚比例
    pub max_slash_rate: Decimal,
    /// 争议解决时间(秒)
    pub dispute_resolution_time: u64,
    /// 版本
    pub version: String,
}

8.4 Registry 查询方式

链下查询

# 通过 REST API 查询价格源元数据
curl https://rest.msgchain.org/msgchain/registry/v1/key/oracle/feeds/btc-usd

# 通过 CLI 查询
msgchaind q registry key oracle/feeds/btc-usd

# 查询所有 feeds
msgchaind q registry prefix oracle/feeds/

链上查询(CosmWasm)

// 从 CosmWasm 合约中查询 Registry
pub fn query_feed_from_registry(
    deps: Deps,
    symbol: &str,
) -> Result<PriceFeedMetadata, ContractError> {
    let registry_key = format!("oracle/feeds/{}", symbol);

    // 通过 Registry 模块查询
    let registry_query = RegistryQuery {
        key: registry_key,
    };

    let result: PriceFeedMetadata = deps.querier.query(&QueryRequest::Custom(
        MsgchainQuery::Registry(registry_query),
    ))?;

    Ok(result)
}

8.5 Registry 更新流程

Registry 数据的更新需要遵循权限控制:

#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub enum RegistryUpdateMsg {
    /// 注册新的价格源
    RegisterFeed {
        metadata: PriceFeedMetadata,
        admin_signature: Binary,
    },
    /// 更新价格源配置
    UpdateFeed {
        symbol: String,
        metadata: PriceFeedMetadata,
        admin_signature: Binary,
    },
    /// 注册节点
    RegisterNode {
        registration: NodeRegistration,
    },
    /// 更新节点心跳
    UpdateHeartbeat {
        node_address: String,
        timestamp: String,
    },
    /// 更新合约地址
    UpdateContract {
        contract_type: String,
        address: String,
    },
}

权限模型

操作 权限要求 备注
注册 Feed 治理通过或管理员多签 需要 Committee 投票
更新 Feed 管理员多签 Precompile 管理员
注册节点 节点自身 同时需要质押
更新心跳 节点自身 自动更新
更新合约 治理通过 重大变更需要治理
更新全局配置 治理通过 参数调整

8.6 Registry 与合约的协同

启动连接流程

1. 通过治理提案在 Registry 中注册 oracle 相关键位
2. 部署 Oracle Aggregator 合约,将合约地址写入 Registry
3. 部署 Oracle Staking 合约,写入 Registry
4. 节点通过 Registry 查询合约地址
5. 节点注册到 Staking 合约,同时更新 Registry 中的节点信息
6. 消费者合约通过 Registry 查找 Aggregator 合约地址
7. 消费者合约调用 Aggregator 查询价格

9. 安全考虑

9.1 数据源多样性

重要性

单一数据源是预言机系统最常见的失败模式。如果所有节点都从同一个 API 获取数据,那么该 API 的故障或数据异常会导致整个预言机系统输出错误价格。

最佳实践

  1. 每个节点使用多个独立数据源:每个节点至少从 3 个独立数据源获取数据
  2. 节点间数据源异质化:不同节点应该使用不同的数据源组合
  3. 数据源类型多样化:同时包含 CEX 和 DEX 数据源
  4. 避免公共依赖:两个看似独立的数据源可能在底层共享相同的数据提供商
// 节点配置中的多数据源示例
let source_configs = vec![
    SourceConfig {
        name: "binance",
        url: "https://api.binance.com/api/v3/ticker/price?symbol=BTCUSDT",
        weight: 0.3,
    },
    SourceConfig {
        name: "coinbase",
        url: "https://api.coinbase.com/v2/prices/BTC-USD/spot",
        weight: 0.3,
    },
    SourceConfig {
        name: "kraken",
        url: "https://api.kraken.com/0/public/Ticker?pair=XBTUSD",
        weight: 0.2,
    },
    SourceConfig {
        name: "bybit",
        url: "https://api.bybit.com/v5/market/tickers?category=spot&symbol=BTCUSDT",
        weight: 0.2,
    },
];

// 节点内聚合(节点自身从中位数计算)
let prices: Vec<Decimal> = fetch_all_sources(&source_configs).await?;
let node_price = compute_median(&prices);

9.2 节点廉洁保证

质押机制

节点必须质押原生代币才能参与预言机网络。质押量决定了节点的"权重"和潜在损失。

质押设计原则

激励相容条件:
  诚实行为的期望收益 >= 恶意行为的期望收益

即:
  奖励 * 诚实概率 >= (操纵收益 - 罚没损失) * 攻击成功概率

转化为质押要求:
  质押量 >= (操纵收益 * 攻击成功概率) / 罚没比例

锁定期

节点申请解绑质押后,需要经过一个锁定期才能取回代币。锁定期为发现和惩罚该节点在锁定期前的不良行为预留了时间窗口。

pub const UNBONDING_PERIOD: u64 = 21 * 24 * 60 * 60; // 21 天(秒)

pub struct UnbondingEntry {
    pub node: Addr,
    pub amount: Uint128,
    pub start_time: Timestamp,
    pub completion_time: Timestamp,
}

9.3 争议解决

争议场景

当某个数据消费者认为预言机报告的价格存在异常时,可以发起争议。

争议流程

1. 争议发起者提交争议证据(需支付争议押金)
2. 进入争议窗口期(如 24 小时)
3. 收集各方证据
4. 治理投票或仲裁委员会裁决
5. 如有恶意行为,执行惩罚
6. 如争议无效,没收发起者押金

链上争议合约

#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct Dispute {
    /// 争议 ID
    pub id: u64,
    /// 发起者
    pub challenger: Addr,
    /// 被争议的节点
    pub accused_node: Addr,
    /// 争议类型
    pub dispute_type: DisputeType,
    /// 争议涉及的轮次
    pub round_id: u64,
    /// 争议涉及的 symbol
    pub symbol: String,
    /// 证据数据
    pub evidence: Binary,
    /// 状态
    pub status: DisputeStatus,
    /// 创建时间
    pub created_at: Timestamp,
    /// 解决时间
    pub resolved_at: Option<Timestamp>,
    /// 争议押金
    pub deposit: Coin,
}

#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub enum DisputeType {
    /// 价格明显偏离外部市场
    PriceManipulation,
    /// 节点提交的数据与实际不符
    DataMismatch,
    /// 节点在多个轮次中行为异常
    PatternManipulation,
}

#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub enum DisputeStatus {
    Open,
    UnderInvestigation,
    ResolvedSlash,
    ResolvedDismissed,
}

9.4 惩罚削权(Slashing)

自动惩戒

某些类型的违规行为可以在链上自动检测并处罚:

pub fn check_and_slash_outliers(
    deps: DepsMut,
    symbol: &str,
    round_id: u64,
    aggregated_price: Decimal,
    deviation_threshold: Decimal,
) -> Result<Vec<Addr>, ContractError> {
    let mut slashed_nodes = vec![];

    // 遍历该轮次所有报告
    let reports = REPORTS.range(
        deps.storage,
        Some((symbol.to_string(), round_id, Addr::unchecked(""))),
        Some((symbol.to_string(), round_id + 1, Addr::unchecked(""))),
        Order::Ascending,
    );

    for result in reports {
        let ((sym, rid, node_addr), price) = result?;
        let deviation = (price - aggregated_price).abs() / aggregated_price;

        if deviation > deviation_threshold {
            // 执行自动惩罚
            let mut node = NODES.load(deps.storage, &node_addr)?;
            let slash_amount = node.stake_amount * Decimal::percent(1);
            node.stake_amount = node.stake_amount.checked_sub(slash_amount)?;
            node.slash_count += 1;

            if node.stake_amount < node.min_stake {
                node.active = false;
            }

            NODES.save(deps.storage, &node_addr, &node)?;
            slashed_nodes.push(node_addr);
        }
    }

    Ok(slashed_nodes)
}

累积惩罚

允许对重复违规者施加逐渐增加的惩罚:

pub fn calculate_slash_rate(slash_count: u64, base_rate: Decimal) -> Decimal {
    // 每次违规惩罚比例增加 50%
    base_rate * Decimal::from_ratio(3u128.pow(slash_count as u32), 2u128.pow(slash_count as u32))
}
违规次数 惩罚比例
1 1.0%
2 1.5%
3 2.25%
4 3.375%
5 5.062%

9.5 闪电贷攻击防护

攻击原理

攻击者通过闪电贷借入大量资金,在流动性较低的 DEX 上操纵价格,同时触发依赖该价格的预言机和协议。

防护措施

  1. 使用 TWAP:瞬时价格操纵无法影响 TWAP
  2. 设置最小报告数:要求足够多的节点和数源报告
  3. 价格偏差限制:单次更新的价格变化不能超过阈值
  4. 数据源中包含非 DEX 价格:CEX 价格不易被链上闪电贷操纵
pub fn validate_price_update(
    new_price: Decimal,
    last_price: Decimal,
    max_deviation_pct: Decimal,
    min_sources: u32,
    actual_sources: u32,
) -> Result<(), ContractError> {
    // 检查价格偏差
    let deviation = (new_price - last_price).abs() / last_price;
    if deviation > max_deviation_pct {
        return Err(ContractError::ExcessiveDeviation {
            deviation,
            max: max_deviation_pct,
        });
    }

    // 检查数据源数量
    if actual_sources < min_sources {
        return Err(ContractError::InsufficientSources {
            actual: actual_sources,
            required: min_sources,
        });
    }

    Ok(())
}

9.6 治理后门与安全委员会

紧急暂停

在发生安全事件时,需要能够紧急暂停预言机系统:

#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct EmergencyConfig {
    /// 是否暂停
    pub paused: bool,
    /// 暂停原因
    pub reason: Option<String>,
    /// 暂停者
    pub paused_by: Option<Addr>,
    /// 暂停时间
    pub paused_at: Option<Timestamp>,
}

pub fn execute_pause(
    deps: DepsMut,
    info: MessageInfo,
    reason: String,
) -> Result<Response, ContractError> {
    // 验证调用者是否为安全委员会成员
    let committee = COMMITTEE.load(deps.storage)?;
    if !committee.contains(&info.sender) {
        return Err(ContractError::Unauthorized);
    }

    PAUSE.save(deps.storage, &EmergencyConfig {
        paused: true,
        reason: Some(reason),
        paused_by: Some(info.sender),
        paused_at: Some(deps.block.time),
    })?;

    Ok(Response::new()
        .add_attribute("action", "pause")
        .add_attribute("reason", &reason))
}

pub fn execute_unpause(
    deps: DepsMut,
    info: MessageInfo,
) -> Result<Response, ContractError> {
    // 治理多签解暂停
    let admin = ADMIN.load(deps.storage)?;
    if info.sender != admin {
        return Err(ContractError::Unauthorized);
    }

    PAUSE.save(deps.storage, &EmergencyConfig {
        paused: false,
        reason: None,
        paused_by: None,
        paused_at: None,
    })?;

    Ok(Response::new().add_attribute("action", "unpause"))
}

9.7 经济安全模型总结

+-------------------------------------------------------------+
|                    攻击成本计算                                |
|                                                              |
|  攻击成本 =                                                    |
|    控制节点所需质押量                                           |
|    + 操纵数据源的成本                                          |
|    + 被惩罚后损失的信誉                                       |
|    + 遭到法律追诉的风险                                       |
|                                                              |
|  攻击收益 =                                                    |
|    从 DeFi 协议中提取的价值                                     |
|    - 交易费用                                                  |
|    - 滑点损失                                                  |
|                                                              |
|  安全条件: 攻击成本 > 攻击收益                                   |
+-------------------------------------------------------------+

通过以下方式确保攻击成本最大化:

  1. 高质押要求:控制足够数量的节点需要大量资本
  2. 数据源多样性:同时操纵多个独立数据源极其困难
  3. TWAP 窗口:需要长时间维持操纵价格
  4. 锁定期:作恶后无法立即提取质押
  5. 争议机制:提供额外的威慑

10. 与 DeFi 合约集成

10.1 集成模式

预言机的最终价值体现在 DeFi 协议的集成使用中。以下介绍如何在 MSG Chain 上的常见 DeFi 协议中集成预言机。

10.2 借贷协议集成

清算触发

借贷协议是最依赖预言机的 DeFi 组件。其核心逻辑是维持抵押率在安全阈值以上:

// 借贷协议中的预言机查询
use msg_oracle::msg::QueryMsg as OracleQuery;
use msg_oracle::state::PriceResponse;

pub fn query_user_health(
    deps: Deps,
    user: Addr,
    oracle_address: Addr,
) -> Result<HealthResponse, ContractError> {
    let user_positions = POSITIONS.load(deps.storage, &user)?;
    let mut total_collateral_value = Uint128::zero();
    let mut total_debt_value = Uint128::zero();

    for position in &user_positions {
        // 查询预言机获取当前价格
        let price: PriceResponse = deps.querier.query_wasm_smart(
            oracle_address.clone(),
            &OracleQuery::GetPrice {
                symbol: position.collateral_symbol.clone(),
            },
        )?;

        let collateral_value = position.collateral_amount * price.price;
        total_collateral_value += collateral_value;

        let debt_value = position.debt_amount; // 假设债务以稳定币计价
        total_debt_value += debt_value;
    }

    let health_ratio = if total_debt_value.is_zero() {
        Decimal::from_ratio(u128::MAX, 1u128)
    } else {
        Decimal::from_ratio(total_collateral_value, total_debt_value)
    };

    Ok(HealthResponse {
        user,
        collateral_value: total_collateral_value,
        debt_value: total_debt_value,
        health_ratio,
        is_healthy: health_ratio >= LIQUIDATION_THRESHOLD,
    })
}

清算执行

pub fn execute_liquidate(
    deps: DepsMut,
    info: MessageInfo,
    user: Addr,
    oracle_address: Addr,
) -> Result<Response, ContractError> {
    // 查询用户健康状态
    let health = query_user_health(deps.as_ref(), user.clone(), oracle_address)?;
    if health.is_healthy {
        return Err(ContractError::PositionHealthy {
            ratio: health.health_ratio,
        });
    }

    // 执行清算
    let liquidator_reward = health.collateral_value * LIQUIDATOR_REWARD_RATE;
    // ... 清算逻辑

    Ok(Response::new()
        .add_attribute("action", "liquidate")
        .add_attribute("user", user.as_str())
        .add_attribute("liquidator", info.sender.as_str())
        .add_attribute("reward", liquidator_reward.to_string()))
}

10.3 AMM 滑点保护

动态滑点

AMM 可以使用预言机价格来设置动态滑点保护,减少三明治攻击损失:

pub fn compute_dynamic_slippage(
    oracle_price: Decimal,
    amm_price: Decimal,
    trade_size: Uint128,
    pool_liquidity: Uint128,
) -> Decimal {
    // 基础滑点 = 交易大小 / 流动性
    let base_slippage = Decimal::from_ratio(trade_size, pool_liquidity);

    // 价格偏离度 = |预言机价格 - AMM 价格| / AMM 价格
    let price_deviation = (oracle_price - amm_price).abs() / amm_price;

    // 动态滑点 = 基础滑点 + 价格偏离度 * 系数
    let dynamic_slippage = base_slippage + price_deviation * Decimal::percent(50);

    // 最小滑点保护
    min_slippage.max(dynamic_slippage)
}

价格偏离限制

pub fn enforce_price_consistency(
    oracle_price: Decimal,
    amm_price: Decimal,
    max_deviation: Decimal,
) -> Result<(), ContractError> {
    let deviation = (oracle_price - amm_price).abs() / amm_price;
    if deviation > max_deviation {
        return Err(ContractError::ExcessivePriceDeviation {
            oracle_price,
            amm_price,
            deviation,
            max: max_deviation,
        });
    }
    Ok(())
}

10.4 稳定币锚定机制

锚定检查

算法稳定币需要持续监控价格偏离度,并调整供应量:

pub fn check_peg_stability(
    oracle_price: Decimal,
    target_price: Decimal,
    peg_threshold: Decimal,
) -> PegStatus {
    let deviation = (oracle_price - target_price).abs() / target_price;

    if deviation <= peg_threshold {
        PegStatus::Stable
    } else if oracle_price > target_price {
        PegStatus::AbovePeg { deviation }
    } else {
        PegStatus::BelowPeg { deviation }
    }
}

pub fn execute_peg_rebalance(
    deps: DepsMut,
    oracle_address: Addr,
) -> Result<Response, ContractError> {
    // 查询预言机价格
    let price: PriceResponse = deps.querier.query_wasm_smart(
        oracle_address,
        &OracleQuery::GetPrice {
            symbol: "STABLE/USD".to_string(),
        },
    )?;

    let peg_status = check_peg_stability(price.price, Decimal::one(), Decimal::percent(1));

    match peg_status {
        PegStatus::AbovePeg { deviation } => {
            // 价格高于锚定 -> 增发稳定币/减少抵押
            expand_supply(deps, deviation)?;
        }
        PegStatus::BelowPeg { deviation } => {
            // 价格低于锚定 -> 回收稳定币/增加抵押
            contract_supply(deps, deviation)?;
        }
        PegStatus::Stable => {
            // 在锚定范围内,不做操作
        }
    }

    Ok(Response::new()
        .add_attribute("action", "rebalance")
        .add_attribute("peg_status", format!("{:?}", peg_status)))
}

10.5 合成资产协议

铸造与赎回

合成资产协议的核心是保持合成资产与标的资产的价格锚定:

pub fn execute_mint_synthetic(
    deps: DepsMut,
    info: MessageInfo,
    oracle_address: Addr,
    symbol: String,
    collateral_amount: Uint128,
) -> Result<Response, ContractError> {
    // 查询标的资产价格
    let price: PriceResponse = deps.querier.query_wasm_smart(
        oracle_address,
        &OracleQuery::GetPrice {
            symbol: format!("{}/USD", symbol),
        },
    )?;

    // 计算可铸造数量
    let mint_amount = calculate_mint_amount(
        collateral_amount,
        price.price,
        COLLATERALIZATION_RATIO,
    )?;

    // 锁定抵押品
    // 铸造合成资产
    // ...

    Ok(Response::new()
        .add_attribute("action", "mint_synthetic")
        .add_attribute("symbol", &symbol)
        .add_attribute("amount", mint_amount.to_string()))
}

10.6 Oracle Proxy 合约

设计目的

为了避免每个 DeFi 协议都直接依赖 Oracle Aggregator 合约(升级时需要更新所有消费者),引入 Proxy 合约作为中间层:

#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct ProxyConfig {
    /// 当前使用的 Aggregator 地址
    pub aggregator: Addr,
    /// 待切换的 Aggregator 地址
    pub pending_aggregator: Option<Addr>,
    /// 延迟时间(治理更新后的切换延迟)
    pub timelock_seconds: u64,
    /// 切换生效时间
    pub effective_time: Option<Timestamp>,
}

pub fn execute_update_aggregator(
    deps: DepsMut,
    info: MessageInfo,
    new_aggregator: Addr,
) -> Result<Response, ContractError> {
    let admin = ADMIN.load(deps.storage)?;
    if info.sender != admin {
        return Err(ContractError::Unauthorized);
    }

    let mut config = CONFIG.load(deps.storage)?;
    config.pending_aggregator = Some(new_aggregator);
    config.effective_time = Some(
        deps.block.time.plus_seconds(config.timelock_seconds),
    );
    CONFIG.save(deps.storage, &config)?;

    Ok(Response::new()
        .add_attribute("action", "schedule_update")
        .add_attribute("new_aggregator", new_aggregator.as_str())
        .add_attribute("effective_time", config.effective_time.unwrap().to_string()))
}

pub fn execute_query_price(
    deps: Deps,
    symbol: String,
) -> Result<PriceResponse, ContractError> {
    let config = CONFIG.load(deps.storage)?;

    // 查询当前 Aggregator
    let price: PriceResponse = deps.querier.query_wasm_smart(
        config.aggregator.clone(),
        &OracleQuery::GetPrice { symbol },
    )?;

    Ok(price)
}

Proxy 合约的使用流程:

DeFi 协议 -> Oracle Proxy -> Oracle Aggregator V1
                          -> Oracle Aggregator V2 (切换后)

10.7 集成清单

步骤 描述 验证方式
1 部署 Oracle Aggregator 合约 合约地址注册到 Registry
2 部署 Oracle Proxy 合约 Proxy 指向 Aggregator
3 注册节点并激活 至少 3 个节点在线
4 配置价格源 BTC/USD, ETH/USD 等
5 等待首个轮次完成 查询到首笔价格
6 验证价格准确性 与外部市场对比
7 集成到 DeFi 协议 Proxy 地址写入协议配置
8 端到端测试 模拟各种市场场景
9 设置监控告警 配置 SLA 监控
10 逐步提高安全参数 从低安全模式升级

11. 监控与 SLA

11.1 SLA 指标

预言机网络的 SLA(服务等级协议)从以下维度衡量:

指标 描述 目标值 告警阈值
数据新鲜度 最新价格与当前时间的差值 < 30 秒 > 60 秒
节点参与率 每轮次实际参与节点数/总节点数 > 80% < 60%
价格更新延迟 从数据源到链上确认的时间 < 15 秒 > 30 秒
响应延迟 合约查询的返回时间 < 100ms > 500ms
节点存活率 24 小时内节点在线时间 > 99.9% < 99%
异常值率 被剔除的异常报告占比 < 1% > 5%
节点分歧度 节点间报告的偏差 < 0.1% > 0.5%

11.2 数据新鲜度监控

链上新鲜度检查

pub fn query_price_freshness(
    deps: Deps,
    symbol: &str,
    max_age_seconds: u64,
) -> Result<FreshnessStatus, ContractError> {
    let price = PRICES.load(deps.storage, symbol)?;
    let current_time = deps.block.time;
    let age = current_time.minus_nanos(price.updated_at.nanos());

    Ok(FreshnessStatus {
        symbol: symbol.to_string(),
        price: price.price,
        updated_at: price.updated_at,
        age_seconds: age.seconds(),
        is_fresh: age.seconds() <= max_age_seconds,
    })
}

链下 Prometheus 指标

// 链下节点暴露的 Prometheus 指标
lazy_static! {
    static ref PRICE_FRESHNESS: GaugeVec = register_gauge_vec!(
        "oracle_price_freshness_seconds",
        "Time since last price update in seconds",
        &["symbol"]
    ).unwrap();

    static ref NODE_PARTICIPATION: Gauge = register_gauge!(
        "oracle_node_participation_rate",
        "Ratio of active nodes participating in rounds"
    ).unwrap();

    static ref REPORT_LATENCY: Histogram = register_histogram!(
        "oracle_report_latency_seconds",
        "Latency of price report submission",
        vec![0.1, 0.5, 1.0, 2.0, 5.0, 10.0, 30.0]
    ).unwrap();

    static ref NODE_HEALTH: GaugeVec = register_gauge_vec!(
        "oracle_node_health",
        "Health status of oracle nodes (1 = healthy, 0 = unhealthy)",
        &["node_address"]
    ).unwrap();
}

11.3 节点 Uptime 监控

心跳机制

每个节点定期在链上提交心跳,证明自己处于活跃状态:

pub fn execute_heartbeat(
    deps: DepsMut,
    info: MessageInfo,
) -> Result<Response, ContractError> {
    let mut node = NODES.load(deps.storage, &info.sender)?;
    node.last_heartbeat = deps.block.time;
    NODES.save(deps.storage, &info.sender, &node)?;

    Ok(Response::new()
        .add_attribute("action", "heartbeat")
        .add_attribute("node", info.sender.as_str())
        .add_attribute("time", deps.block.time.to_string()))
}

Uptime 跟踪

pub fn check_node_downtime(
    deps: Deps,
    max_missed_heartbeats: u64,
    heartbeat_interval: u64,
) -> Result<Vec<Addr>, ContractError> {
    let mut down_nodes = vec![];
    let current_time = deps.block.time;

    for result in NODES.range(deps.storage, None, None, Order::Ascending) {
        let (addr, node) = result?;
        if !node.active {
            continue;
        }

        let elapsed = current_time.minus_nanos(node.last_heartbeat.nanos());
        let missed_beats = elapsed.seconds() / heartbeat_interval;

        if missed_beats >= max_missed_heartbeats {
            down_nodes.push(addr);
        }
    }

    Ok(down_nodes)
}

11.4 响应延迟监控

端到端延迟测量

外部数据源 API 调用时间
    + 节点内部处理时间
    + 交易签名时间
    + 交易广播延迟
    + 链上 mempool 等待时间
    + 区块包含和确认时间
    + 合约执行时间
    = 端到端延迟

延迟分解

典型延迟分布(预估):

 数据获取:     50-200ms   (HTTP API 调用)
 数据处理:     5-20ms     (JSON 解析、验证)
 交易签名:     1-5ms      (本地签名)
 交易广播:     50-200ms   (RPC 提交)
 等待入块:     1-6s       (Tendermint 共识)
 合约执行:     10-50ms    (Wasm 执行)

 端到端总计:   1.2-6.5s   (取决于网络状况)

11.5 告警规则

Prometheus Alertmanager 配置示例

groups:
  - name: oracle_alerts
    rules:
      - alert: PriceStale
        expr: oracle_price_freshness_seconds > 60
        for: 30s
        labels:
          severity: critical
        annotations:
          summary: "Price feed {{ $labels.symbol }} is stale"
          description: "Price for {{ $labels.symbol }} has not been updated for {{ $value }} seconds"

      - alert: NodeDown
        expr: oracle_node_health == 0
        for: 5m
        labels:
          severity: warning
        annotations:
          summary: "Oracle node {{ $labels.node_address }} is down"
          description: "Node has been unresponsive for more than 5 minutes"

      - alert: LowParticipation
        expr: oracle_node_participation_rate < 0.6
        for: 10m
        labels:
          severity: warning
        annotations:
          summary: "Low node participation rate"
          description: "Participation rate is {{ $value | humanizePercentage }}"

      - alert: HighReportLatency
        expr: histogram_quantile(0.95, rate(oracle_report_latency_seconds_bucket[5m])) > 5
        for: 5m
        labels:
          severity: warning
        annotations:
          summary: "High report latency"
          description: "95th percentile report latency is {{ $value }} seconds"

      - alert: NodeDeviation
        expr: oracle_node_deviation > 0.005
        for: 1h
        labels:
          severity: info
        annotations:
          summary: "Node {{ $labels.node_address }} shows high deviation"
          description: "Deviation is {{ $value | humanizePercentage }}"

11.6 Grafana 监控面板

面板布局

+-------------------------------------------------------------+
|  预言机网络概览                                                |
|  +-----------+ +-----------+ +-----------+ +-----------+      |
|  | 活跃节点数 | | 价格源数   | | 平均新鲜度 | | 参与率    |      |
|  +-----------+ +-----------+ +-----------+ +-----------+      |
|                                                              |
|  价格源最新价格                                                |
|  +-------------+ +-------------+ +-------------+              |
|  | BTC/USD     | | ETH/USD     | | MSG/USD     |              |
|  | $58,234     | | $3,124      | | $0.0523     |              |
|  | (2s 前)     | | (1s 前)     | | (5s 前)     |              |
|  +-------------+ +-------------+ +-------------+              |
|                                                              |
|  价格新鲜度时间序列                                             |
|  +-----------------------------------------------------+     |
|  | 图表: 各价格源最新更新距现在的时间                         |     |
|  +-----------------------------------------------------+     |
|                                                              |
|  节点健康状态表                                                |
|  +-----+----------+--------+---------+--------+-------+      |
|  | 节点 | 地址     | 状态   | 质押量   | 心跳   | 异常  |      |
|  +-----+----------+--------+---------+--------+-------+      |
|  | 1   | msg1... | Active | 100000  | 2s ago | 0     |      |
|  | 2   | msg1... | Active | 100000  | 3s ago | 0     |      |
|  | 3   | msg1... | Active | 100000  | 1s ago | 0     |      |
|  +-----+----------+--------+---------+--------+-------+      |
+-------------------------------------------------------------+

11.7 SLA 报告

定期 SLA 报告模板

# 预言机网络 SLA 报告

周期: 2026-07-01 至 2026-07-07

## 概览
- 总运行时间: 168 小时
- 计划内停机: 0 小时
- 计划外停机: 0.2 小时
- 可用性: 99.88%

## 价格源
| 价格源   | 更新次数 | 平均更新间隔 | 最大更新间隔 | 新鲜度达标率 |
|----------|---------|-------------|-------------|-------------|
| BTC/USD  | 20160   | 29.8s       | 45.2s       | 99.5%       |
| ETH/USD  | 20160   | 29.8s       | 42.1s       | 99.6%       |
| MSG/USD  | 10080   | 59.7s       | 78.3s       | 98.2%       |

## 节点
| 节点 ID   | 参与率 | 准确率 | 延迟(p50) | 延迟(p99) |
|-----------|--------|--------|-----------|-----------|
| node-1    | 99.8%  | 99.97% | 1.2s      | 4.5s      |
| node-2    | 99.5%  | 99.95% | 1.4s      | 5.1s      |
| node-3    | 99.7%  | 99.98% | 1.1s      | 4.2s      |

## 事件
- 2026-07-03 14:22 UTC: Binance API 短暂不可用 (30s),节点自动切换至备用数据源
- 2026-07-05 08:15 UTC: 节点 node-2 因网络维护离线 12 分钟

## 改进建议
- MSG/USD 价格源更新间隔偏高,建议增加一个数据源

12. 当前 Stub 状态与未来路线图

12.1 Stub 状态说明

REST 端点状态

MSG Chain 当前的 /agent/v1/oracle/price 端点状态:

GET https://rest.msgchain.org/agent/v1/oracle/price?symbol=BTC/USD

响应 (当前 Stub):
{
  "status": "stub",
  "message": "Oracle module is not yet implemented",
  "symbol": "BTC/USD",
  "price": null,
  "timestamp": null
}

可用功能

功能 状态 备注
REST 端点路由 ✅ 已注册 仅返回占位响应
链上合约 ❌ 未部署 需要 CosmWasm 开发
链下节点 ❌ 未开发 需要节点软件
Registry 键位 ❌ 未注册 需要治理提案
SDK/CLI 集成 ❌ 未实现 需要客户端更新

当前可以开始的工作

即使 Oracle 模块是 Stub,以下工作可以立即开始:

  1. 合约开发:使用 CosmWasm 开发 Oracle 聚合合约,在本地测试网验证
  2. 节点原型:开发链下节点软件,模拟数据提交流程
  3. Registry 设计:完成键位设计提案,提交给社区讨论
  4. 安全模型设计:完成经济安全分析和参数计算
  5. 集成测试:编写 DeFi 协议与 Oracle 的集成测试用例

12.2 路线图

阶段一:基础架构(预计 1-2 个月)

目标: 搭建最小可行预言机网络

任务:
  [ ] 设计并开发 Oracle Aggregator CosmWasm 合约 v1
  [ ] 设计并开发 Oracle Staking CosmWasm 合约 v1
  [ ] 开发链下节点原型(单数据源、中位数聚合)
  [ ] 在 MSG Chain testnet 上部署测试版本
  [ ] 基础监控面板搭建
  [ ] Registry 键位注册

交付物:
  - 在 testnet 上运行的 3 节点预言机网络
  - 至少 2 个价格源(BTC/USD, ETH/USD)
  - 完整的合约源码和部署脚本
  - 节点操作文档

阶段二:安全强化(预计 2-3 个月)

目标: 提升系统安全性和可靠性

任务:
  [ ] 实现 TWAP 聚合策略
  [ ] 实现异常值自动剔除
  [ ] 实现节点信誉系统
  [ ] 实现争议解决合约
  [ ] 实现自动惩罚机制
  [ ] 多数据源支持(每个节点 3+ 数据源)
  [ ] 链下 OCR 协议实现
  [ ] 安全审计(内部 + 第三方)

交付物:
  - 通过安全审计的合约代码
  - 支持 8 个以上价格源
  - 节点数量扩展到 7+
  - 完整的 SLA 监控系统

阶段三:生态集成(预计 3-4 个月)

目标: 与 MSG Chain 上的 DeFi 生态深度集成

任务:
  [ ] 发布 Oracle Proxy 合约
  [ ] 与借贷协议集成
  [ ] 与 AMM 协议集成
  [ ] 合成资产协议支持
  [ ] VRF 随机数预言机
  [ ] 跨链价格源(通过 IBC)
  [ ] 开发者 SDK 发布
  [ ] 资助和激励计划启动

交付物:
  - MSG Chain DeFi 生态中的预言机标准
  - 开发者文档和示例代码
  - 集成 SDK 和 CLI 工具
  - 生态系统资助提案

阶段四:生产就绪(预计 4-6 个月)

目标: 主网生产环境全面运行

任务:
  [ ] 主网 Oracle Aggregator 合约部署
  [ ] 主网 Registry 键位更新
  [ ] 正式节点操作者激励计划
  [ ] 与外部审计师完成最终审计
  [ ] 发布 SLA 承诺
  [ ] 社区治理参数投票
  [ ] 持续优化 Gas 成本

交付物:
  - 主网上完全运行的预言机网络
  - 10+ 独立节点操作者
  - 20+ 价格源
  - 3+ 集成 DeFi 协议
  - 99.9% 以上可用性

12.3 已知限制和风险

技术风险

风险 概率 影响 缓解措施
CosmWasm Gas 限制导致复杂聚合失败 中 高 OCR 链下聚合、优化合约代码
节点数量不足影响去中心化 低 中 激励计划、降低参与门槛
区块空间竞争导致价格更新延迟 中 中 优先级机制、动态 Gas
数据源 API 变更导致中断 高 高 多数据源冗余、快速响应更新

经济风险

风险 概率 影响 缓解措施
质押量不足以威慑攻击 中 高 定期安全评估、参数调整
节点合谋操纵价格 低 极高 数据源多样性、异常检测、争议机制
闪电贷结合价格操纵 中 高 TWAP、最小报告数、CEX 数据源
治理攻击(参数篡改) 低 高 时间锁、多签治理、安全委员会

运营风险

风险 概率 影响 缓解措施
节点运营者退出 中 中 锁定期、平滑过渡、备用节点池
节点私钥泄露 低 高 HSM/KMS、多签操作、密钥轮换
区块链网络升级不兼容 中 中 合约迁移机制、充分测试

12.4 社区参与

贡献方式

沟通渠道


附录

A.1 术语表

术语 英文 定义
预言机 Oracle 将外部数据桥接到区块链的组件
价格喂价 Price Feed 持续更新的链上价格数据流
聚合合约 Aggregator Contract 收集并整合多个节点报告的智能合约
链下报告 Off-Chain Reporting 节点在链下协商后统一提交的机制
质押 Staking 节点锁定代币作为行为担保
惩罚 Slashing 对不良行为节点罚没质押代币
信誉评分 Reputation Score 追踪节点历史表现的综合评分
时间加权平均价格 TWAP 按时间加权的平均价格
交易量加权平均价格 VWAP 按交易量加权的平均价格
可验证随机函数 VRF 可公开验证的安全随机数生成
女巫攻击 Sybil Attack 创建多个虚假身份破坏网络
三明治攻击 Sandwich Attack 通过前置和后置交易套利
数据源适配器 Adapter 连接外部数据源的标准化组件

A.2 参考资源

A.3 配置变更记录

版本 日期 变更内容 作者
v0.1.0 2026-07-08 初始架构设计 Oracle WG

A.4 免责声明

本文档描述的预言机网络架构和组件目前处于规划阶段。MSG Chain 的 /agent/v1/oracle/price 端点在本文档撰写时返回 Stub 响应,尚未实现实际的预言机功能。本文档的内容仅作为技术参考和规划蓝图,不构成对任何功能可用性的承诺或保证。

实际部署时应参考最新的官方文档和技术公告。涉及主网资产和协议的操作应在充分测试和审计后进行。


本文档由 MSG Chain Oracle 工作组维护
如有问题或建议,请提交 GitHub Issue 或参与社区讨论