dApp Docs/预言机集成与价格数据指南
Development reference. Not independently verified for production.

MSG Chain 预言机集成与价格数据指南

主网状态: No-Go

目录

  1. 概述
  2. 价格预言机实现
  3. 链上随机数(VRF)
  4. 数据预言机
  5. 预言机消费者模式
  6. 安全考虑
  7. AI Agent 集成
  8. 附录

1. 概述

1.1 为什么预言机至关重要

区块链是确定性封闭系统,无法原生访问链外数据。预言机(Oracle)桥接了链上与链下世界,使得智能合约能够获取真实世界的信息。在 MSG Chain(msg 前缀地址)上,预言机的核心用途包括:

用例 描述 典型场景
价格数据 获取资产对美元的实时汇率 DEX 定价、借贷清算、衍生品
链上随机数 生成不可预测的安全随机数 游戏抽奖、NFT 铸造、验证者选取
任意数据 任何链下数据源 天气数据、体育比分、AI 模型输出

1.2 预言机设计模型

推模型(Push Model)

预言机节点主动将数据提交到链上,数据按固定时间间隔或事件触发更新。

[数据源] → [预言机节点] → [MSG Chain 合约] → [消费者合约]

拉模型(Pull Model)

消费者合约在需要时主动请求数据,预言机节点仅在被调用时提交数据。

[消费者合约] → [预言机合约(请求)] → [预言机节点(回应)] → [消费者合约(消费)]

混合模型

结合推拉模型的优势:价格数据采用推模型高频更新,而随机数和自定义数据采用拉模型。

1.3 MSG Chain 上的预言机方案对比

方案 类型 去中心化程度 延迟 成本 适用场景
Band Protocol 推+拉 高 低 中 标准价格喂价、DeFi
Skip Oracle 推 中 极低 高 MEV 感知型应用
自定义预言机 自定 取决于部署 可控 可控 AI Agent、特殊数据

Band Protocol

Band Protocol 是一个去中心化的预言机网络,在 MSG Chain 上可以通过 IBC 或直接的 Band 标准桥接使用。标准价格喂价流程:

BandChain (数据源) ──IBC──> MSG Chain (PriceFeed 合约) ──查询──> dApp

Band 的优点在于其成熟的标准和跨链支持。消费者合约通过查询 MSG Chain 上的 BandPriceFeed 合约获取聚合价格。

Skip Oracle

Skip 预言机(原 Osmosis Oracle)通过区块内的高频价格更新机制,提供极低延迟的价格数据。它利用区块提议者提交的价格列表,适合对时效性要求极高的场景(如 MEV 保护)。

自定义预言机

对于 AI Agent 等特殊需求场景,自定义预言机提供了最大的灵活性。本文将重点介绍如何在 MSG Chain 上构建一个完整的自定义预言机系统(包括价格、VRF 和通用数据预言机),所有合约使用 msg1 地址前缀。


2. 价格预言机实现

2.1 数据模型设计

价格预言机围绕三个核心概念设计:价格信息、发布者 和 聚合策略。

// src/state.rs
use cosmwasm_std::{Addr, Decimal256, StdResult, Storage, Timestamp};
use cw_storage_plus::{Item, Map};

/// 单个价格数据点
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, Eq)]
pub struct PriceInfo {
    /// 价格(标准化为 Decimal256)
    pub price: Decimal256,
    /// 精度(例如 8 表示 8 位小数)
    pub decimals: u8,
    /// 更新时间戳(秒)
    pub timestamp: u64,
    /// 发布者地址
    pub publisher: Addr,
}

/// 聚合后的价格数据
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, Eq)]
pub struct AggregatedPrice {
    /// 中位数价格
    pub median_price: Decimal256,
    /// TWAP 价格(时间加权平均)
    pub twap_price: Decimal256,
    /// 参与聚合的发布者数量
    pub publisher_count: u32,
    /// 上次更新时间
    pub updated_at: u64,
    /// 价格有效标志
    pub is_valid: bool,
}

/// 发布者元数据
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, Eq)]
pub struct PublisherInfo {
    /// 发布者地址
    pub address: Addr,
    /// 权重(用于加权中位数)
    pub weight: u8,
    /// 是否活跃
    pub active: bool,
    /// 加入时间
    pub joined_at: u64,
    /// 总提交次数
    pub total_submissions: u64,
    /// 惩罚次数
    pub slashes: u8,
}

/// TWAP 快照点
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, Eq)]
pub struct TwapSnapshot {
    /// 价格
    pub price: Decimal256,
    /// 快照时间
    pub timestamp: u64,
}

/// 合约配置
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, Eq)]
pub struct Config {
    /// 合约所有者
    pub owner: Addr,
    /// 最大价格过期时长(秒)
    pub max_staleness: u64,
    /// TWAP 窗口时长(秒)
    pub twap_window: u64,
    /// 最小发布者数量
    pub min_publishers: u32,
    /// 价格偏差阈值(如 0.05 表示 5%)
    pub deviation_threshold: Decimal256,
    /// 最大发布者数量
    pub max_publishers: u32,
}

/// 核心价格预言机数据
pub struct PriceOracle {
    /// 合约所有者
    pub owner: Addr,
    /// 授权发布者列表
    pub publishers: Vec<Addr>,
    /// 价格存储(key = 币种符号)
    pub prices: Map<&'static str, PriceInfo>,
}

// 存储定义
pub const CONFIG: Item<Config> = Item::new("config");
pub const PUBLISHERS: Map<&Addr, PublisherInfo> = Map::new("publishers");
pub const PRICES: Map<&str, AggregatedPrice> = Map::new("prices");
pub const PRICE_HISTORY: Map<(&str, u64), PriceInfo> = Map::new("price_history");
pub const TWAP_SNAPSHOTS: Map<&str, Vec<TwapSnapshot>> = Map::new("twap_snapshots");

2.2 合约入口

// src/contract.rs
use cosmwasm_std::{
    entry_point, to_json_binary, Binary, Decimal256, Deps, DepsMut, Env,
    MessageInfo, Response, StdError, StdResult, Uint128, Addr, Order,
};
use cw2::set_contract_version;
use crate::error::ContractError;
use crate::msg::{
    ConfigResponse, ExecuteMsg, InstantiateMsg, PriceResponse,
    PublisherResponse, QueryMsg, PriceHistoryResponse,
};
use crate::state::{
    AggregatedPrice, Config, PriceInfo, PublisherInfo, TwapSnapshot,
    CONFIG, PRICES, PRICE_HISTORY, PUBLISHERS, TWAP_SNAPSHOTS,
};

const CONTRACT_NAME: &str = "msgchain-price-oracle";
const CONTRACT_VERSION: &str = "2.0.0";
const MAX_HISTORY_POINTS: usize = 100;
const MAX_SNAPSHOTS: usize = 200;

#[entry_point]
pub fn instantiate(
    deps: DepsMut,
    env: Env,
    info: MessageInfo,
    msg: InstantiateMsg,
) -> Result<Response, ContractError> {
    set_contract_version(deps.storage, CONTRACT_NAME, CONTRACT_VERSION)?;

    let config = Config {
        owner: info.sender.clone(),
        max_staleness: msg.max_staleness.unwrap_or(300),
        twap_window: msg.twap_window.unwrap_or(3600),
        min_publishers: msg.min_publishers.unwrap_or(2),
        deviation_threshold: msg
            .deviation_threshold
            .unwrap_or(Decimal256::from_atomics(Uint128::new(5), 2).unwrap()),
        max_publishers: msg.max_publishers.unwrap_or(20),
    };
    CONFIG.save(deps.storage, &config)?;

    for publisher_addr in &msg.initial_publishers {
        let validated = deps.api.addr_validate(publisher_addr)?;
        let publisher_info = PublisherInfo {
            address: validated.clone(),
            weight: 1,
            active: true,
            joined_at: env.block.time.seconds(),
            total_submissions: 0,
            slashes: 0,
        };
        PUBLISHERS.save(deps.storage, &validated, &publisher_info)?;
    }

    Ok(Response::new()
        .add_attribute("method", "instantiate")
        .add_attribute("owner", info.sender.to_string())
        .add_attribute("publishers", msg.initial_publishers.len().to_string()))
}

#[entry_point]
pub fn execute(
    deps: DepsMut,
    env: Env,
    info: MessageInfo,
    msg: ExecuteMsg,
) -> Result<Response, ContractError> {
    match msg {
        ExecuteMsg::SubmitPrices { prices } => {
            execute_submit_prices(deps, env, info, prices)
        }
        ExecuteMsg::UpdateConfig {
            max_staleness,
            twap_window,
            min_publishers,
            deviation_threshold,
            max_publishers,
        } => execute_update_config(
            deps, env, info,
            max_staleness, twap_window, min_publishers, deviation_threshold, max_publishers,
        ),
        ExecuteMsg::AddPublisher { address } => {
            execute_add_publisher(deps, env, info, address)
        }
        ExecuteMsg::RemovePublisher { address } => {
            execute_remove_publisher(deps, env, info, address)
        }
        ExecuteMsg::SlashPublisher { address } => {
            execute_slash_publisher(deps, env, info, address)
        }
        ExecuteMsg::RequestPriceUpdate { denom } => {
            execute_request_price_update(deps, env, info, denom)
        }
        ExecuteMsg::SubmitSinglePrice {
            denom,
            price,
            decimals,
        } => execute_submit_single_price(deps, env, info, denom, price, decimals),
    }
}

2.3 价格提交流程

// src/execute.rs
use cosmwasm_std::{
    Addr, Decimal256, DepsMut, Env, MessageInfo, Response,
    StdError, StdResult, Uint128,
};
use crate::error::ContractError;
use crate::msg::PriceEntry;
use crate::state::{
    AggregatedPrice, Config, PriceInfo, PublisherInfo, TwapSnapshot,
    CONFIG, PRICES, PRICE_HISTORY, PUBLISHERS, TWAP_SNAPSHOTS,
    MAX_HISTORY_POINTS, MAX_SNAPSHOTS,
};

/// 批量提交多币种价格
pub fn execute_submit_prices(
    deps: DepsMut,
    env: Env,
    info: MessageInfo,
    prices: Vec<PriceEntry>,
) -> Result<Response, ContractError> {
    let publisher = PUBLISHERS
        .may_load(deps.storage, &info.sender)?
        .ok_or(ContractError::Unauthorized {})?;

    if !publisher.active {
        return Err(ContractError::PublisherInactive {});
    }

    let config = CONFIG.load(deps.storage)?;
    let now = env.block.time.seconds();
    let mut response = Response::new()
        .add_attribute("method", "submit_prices")
        .add_attribute("publisher", info.sender.to_string());

    for entry in &prices {
        if entry.price.is_zero() {
            return Err(ContractError::InvalidPrice {});
        }

        let price_info = PriceInfo {
            price: entry.price,
            decimals: entry.decimals,
            timestamp: now,
            publisher: info.sender.clone(),
        };

        PRICE_HISTORY.save(
            deps.storage,
            (&entry.denom, now),
            &price_info,
        )?;

        aggregate_price(deps.storage, &config, &entry.denom, now)?;

        response = response
            .add_attribute(format!("price_{}_denom", &entry.denom), &entry.denom)
            .add_attribute(
                format!("price_{}_value", &entry.denom),
                entry.price.to_string(),
            );
    }

    let mut updated_publisher = publisher;
    updated_publisher.total_submissions += prices.len() as u64;
    PUBLISHERS.save(deps.storage, &info.sender, &updated_publisher)?;

    Ok(response)
}

/// 提交单币种价格
pub fn execute_submit_single_price(
    deps: DepsMut,
    env: Env,
    info: MessageInfo,
    denom: String,
    price: Decimal256,
    decimals: u8,
) -> Result<Response, ContractError> {
    let entry = PriceEntry {
        denom,
        price,
        decimals,
    };
    execute_submit_prices(deps, env, info, vec![entry])
}

2.4 价格聚合(中位数)

/// 中位数价格聚合
pub fn aggregate_price(
    storage: &mut dyn Storage,
    config: &Config,
    denom: &str,
    now: u64,
) -> StdResult<AggregatedPrice> {
    let mut prices: Vec<(Decimal256, u8)> = Vec::new();

    let publishers: Vec<PublisherInfo> = PUBLISHERS
        .range(storage, None, None, Order::Ascending)
        .filter_map(|r| r.ok())
        .map(|(_, p)| p)
        .filter(|p| p.active)
        .collect();

    for publisher in &publishers {
        let recent: Vec<PriceInfo> = PRICE_HISTORY
            .prefix(denom)
            .range(storage, None, None, Order::Descending)
            .take(1)
            .filter_map(|r| r.ok())
            .filter(|p| p.publisher == publisher.address)
            .collect();

        if let Some(price_info) = recent.first() {
            let age = now - price_info.timestamp;
            if age <= config.max_staleness {
                prices.push((price_info.price, publisher.weight as u8));
            }
        }
    }

    if (prices.len() as u32) < config.min_publishers {
        return Ok(AggregatedPrice {
            median_price: Decimal256::zero(),
            twap_price: Decimal256::zero(),
            publisher_count: prices.len() as u32,
            updated_at: now,
            is_valid: false,
        });
    }

    prices.sort_by(|a, b| a.0.cmp(&b.0));
    let median_price = if prices.len() % 2 == 0 {
        let mid = prices.len() / 2;
        let avg = (prices[mid - 1].0 + prices[mid].0)
            .checked_div(Decimal256::from_atomics(Uint128::new(2), 0).unwrap())
            .unwrap();
        avg
    } else {
        prices[prices.len() / 2].0
    };

    let twap_price = calculate_twap(storage, denom, config.twap_window, now, median_price)?;

    let aggregated = AggregatedPrice {
        median_price,
        twap_price,
        publisher_count: prices.len() as u32,
        updated_at: now,
        is_valid: true,
    };

    PRICES.save(storage, denom, &aggregated)?;
    update_twap_snapshots(storage, denom, median_price, now)?;

    Ok(aggregated)
}

2.5 TWAP 计算

/// 计算 TWAP 价格
pub fn calculate_twap(
    storage: &dyn Storage,
    denom: &str,
    window: u64,
    now: u64,
    current_price: Decimal256,
) -> StdResult<Decimal256> {
    let snapshots = TWAP_SNAPSHOTS
        .may_load(storage, denom)?
        .unwrap_or_default();

    if snapshots.is_empty() {
        return Ok(current_price);
    }

    let window_start = now.saturating_sub(window);
    let mut total_price_time = Decimal256::zero();
    let mut total_time = Uint128::zero();

    let mut prev_price = current_price;
    let mut prev_time = now;

    for snapshot in snapshots.iter().rev() {
        if snapshot.timestamp < window_start {
            let time_diff = prev_time - window_start;
            let weight = Decimal256::from_atomics(time_diff.into(), 0).unwrap();
            total_price_time += prev_price * weight;
            total_time += Uint128::from(time_diff);
            break;
        }

        let time_diff = prev_time - snapshot.timestamp;
        let weight = Decimal256::from_atomics(time_diff.into(), 0).unwrap();
        total_price_time += prev_price * weight;
        total_time += Uint128::from(time_diff);

        prev_price = snapshot.price;
        prev_time = snapshot.timestamp;
    }

    if total_time.is_zero() {
        return Ok(current_price);
    }

    let twap = total_price_time
        .checked_div(Decimal256::from_atomics(total_time, 0).unwrap())
        .unwrap();

    Ok(twap)
}

/// 更新 TWAP 快照
pub fn update_twap_snapshots(
    storage: &mut dyn Storage,
    denom: &str,
    price: Decimal256,
    now: u64,
) -> StdResult<()> {
    let mut snapshots = TWAP_SNAPSHOTS
        .may_load(storage, denom)?
        .unwrap_or_default();

    snapshots.push(TwapSnapshot { price, timestamp: now });

    if snapshots.len() > MAX_SNAPSHOTS {
        snapshots.remove(0);
    }

    TWAP_SNAPSHOTS.save(storage, denom, &snapshots)
}

2.6 管理功能

/// 更新配置(仅所有者)
pub fn execute_update_config(
    deps: DepsMut,
    _env: Env,
    info: MessageInfo,
    max_staleness: Option<u64>,
    twap_window: Option<u64>,
    min_publishers: Option<u32>,
    deviation_threshold: Option<Decimal256>,
    max_publishers: Option<u32>,
) -> Result<Response, ContractError> {
    let config = CONFIG.load(deps.storage)?;
    if info.sender != config.owner {
        return Err(ContractError::Unauthorized {});
    }

    let updated = Config {
        max_staleness: max_staleness.unwrap_or(config.max_staleness),
        twap_window: twap_window.unwrap_or(config.twap_window),
        min_publishers: min_publishers.unwrap_or(config.min_publishers),
        deviation_threshold: deviation_threshold.unwrap_or(config.deviation_threshold),
        max_publishers: max_publishers.unwrap_or(config.max_publishers),
        ..config
    };
    CONFIG.save(deps.storage, &updated)?;

    Ok(Response::new()
        .add_attribute("method", "update_config")
        .add_attribute("sender", info.sender.to_string()))
}

/// 添加发布者
pub fn execute_add_publisher(
    deps: DepsMut,
    env: Env,
    info: MessageInfo,
    address: Addr,
) -> Result<Response, ContractError> {
    let config = CONFIG.load(deps.storage)?;
    if info.sender != config.owner {
        return Err(ContractError::Unauthorized {});
    }

    if PUBLISHERS.has(deps.storage, &address) {
        return Err(ContractError::PublisherExists {});
    }

    let count = PUBLISHERS
        .range(deps.storage, None, None, Order::Ascending)
        .count();
    if (count as u32) >= config.max_publishers {
        return Err(ContractError::TooManyPublishers {});
    }

    let publisher = PublisherInfo {
        address: address.clone(),
        weight: 1,
        active: true,
        joined_at: env.block.time.seconds(),
        total_submissions: 0,
        slashes: 0,
    };
    PUBLISHERS.save(deps.storage, &address, &publisher)?;

    Ok(Response::new()
        .add_attribute("method", "add_publisher")
        .add_attribute("address", address.to_string()))
}

/// 移除发布者
pub fn execute_remove_publisher(
    deps: DepsMut,
    _env: Env,
    info: MessageInfo,
    address: Addr,
) -> Result<Response, ContractError> {
    let config = CONFIG.load(deps.storage)?;
    if info.sender != config.owner {
        return Err(ContractError::Unauthorized {});
    }

    if !PUBLISHERS.has(deps.storage, &address) {
        return Err(ContractError::PublisherNotFound {});
    }

    PUBLISHERS.remove(deps.storage, &address);

    Ok(Response::new()
        .add_attribute("method", "remove_publisher")
        .add_attribute("address", address.to_string()))
}

/// 惩罚发布者
pub fn execute_slash_publisher(
    deps: DepsMut,
    _env: Env,
    info: MessageInfo,
    address: Addr,
) -> Result<Response, ContractError> {
    let config = CONFIG.load(deps.storage)?;
    if info.sender != config.owner {
        return Err(ContractError::Unauthorized {});
    }

    let publisher = PUBLISHERS
        .load(deps.storage, &address)?;

    let mut updated = publisher;
    updated.slashes += 1;

    if updated.slashes >= 3 {
        updated.active = false;
    }

    PUBLISHERS.save(deps.storage, &address, &updated)?;

    Ok(Response::new()
        .add_attribute("method", "slash_publisher")
        .add_attribute("address", address.to_string())
        .add_attribute("slashes", updated.slashes.to_string())
        .add_attribute("active", updated.active.to_string()))
}

/// 请求价格更新(拉模型)
pub fn execute_request_price_update(
    deps: DepsMut,
    env: Env,
    info: MessageInfo,
    denom: String,
) -> Result<Response, ContractError> {
    Ok(Response::new()
        .add_attribute("method", "request_price_update")
        .add_attribute("denom", denom)
        .add_attribute("requester", info.sender.to_string())
        .add_attribute("height", env.block.height.to_string()))
}

2.7 查询接口

// src/query.rs
use cosmwasm_std::{to_json_binary, Binary, Deps, Env, StdResult, Order};
use crate::msg::{
    ConfigResponse, PriceResponse, PublisherResponse,
    PriceHistoryResponse, AllPricesResponse,
};
use crate::state::{AggregatedPrice, CONFIG, PRICES, PRICE_HISTORY, PUBLISHERS};

pub fn query_config(deps: Deps) -> StdResult<Binary> {
    let config = CONFIG.load(deps.storage)?;
    to_json_binary(&ConfigResponse {
        owner: config.owner.to_string(),
        max_staleness: config.max_staleness,
        twap_window: config.twap_window,
        min_publishers: config.min_publishers,
        deviation_threshold: config.deviation_threshold.to_string(),
        max_publishers: config.max_publishers,
    })
}

pub fn query_price(deps: Deps, denom: String) -> StdResult<Binary> {
    let price = PRICES
        .may_load(deps.storage, &denom)?
        .ok_or_else(|| {
            StdError::not_found(format!("Price not found for denom: {}", denom))
        })?;

    to_json_binary(&PriceResponse {
        denom,
        median_price: price.median_price.to_string(),
        twap_price: price.twap_price.to_string(),
        publisher_count: price.publisher_count,
        updated_at: price.updated_at,
        is_valid: price.is_valid,
    })
}

pub fn query_all_prices(deps: Deps) -> StdResult<Binary> {
    let prices: Vec<PriceResponse> = PRICES
        .range(deps.storage, None, None, Order::Ascending)
        .filter_map(|r| r.ok())
        .map(|(denom, price)| PriceResponse {
            denom,
            median_price: price.median_price.to_string(),
            twap_price: price.twap_price.to_string(),
            publisher_count: price.publisher_count,
            updated_at: price.updated_at,
            is_valid: price.is_valid,
        })
        .collect();

    to_json_binary(&AllPricesResponse { prices })
}

pub fn query_publisher(deps: Deps, address: String) -> StdResult<Binary> {
    let addr = deps.api.addr_validate(&address)?;
    let publisher = PUBLISHERS
        .may_load(deps.storage, &addr)?
        .ok_or_else(|| {
            StdError::not_found(format!("Publisher not found: {}", address))
        })?;

    to_json_binary(&PublisherResponse {
        address: publisher.address.to_string(),
        weight: publisher.weight,
        active: publisher.active,
        joined_at: publisher.joined_at,
        total_submissions: publisher.total_submissions,
        slashes: publisher.slashes,
    })
}

pub fn query_price_history(
    deps: Deps,
    denom: String,
    limit: Option<u32>,
) -> StdResult<Binary> {
    let max_limit = limit.unwrap_or(10).min(100) as usize;
    let history: Vec<_> = PRICE_HISTORY
        .prefix(&denom)
        .range(deps.storage, None, None, Order::Descending)
        .take(max_limit)
        .filter_map(|r| r.ok())
        .map(|(ts, info)| crate::msg::PriceHistoryEntry {
            price: info.price.to_string(),
            decimals: info.decimals,
            timestamp: ts,
            publisher: info.publisher.to_string(),
        })
        .collect();

    to_json_binary(&PriceHistoryResponse { denom, history })
}

2.8 消息和响应类型

// src/msg.rs
use cosmwasm_std::Decimal256;
use schemars::JsonSchema;
use serde::{Deserialize, Serialize};

#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct InstantiateMsg {
    pub initial_publishers: Vec<String>,
    pub max_staleness: Option<u64>,
    pub twap_window: Option<u64>,
    pub min_publishers: Option<u32>,
    pub deviation_threshold: Option<Decimal256>,
    pub max_publishers: Option<u32>,
}

#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct PriceEntry {
    pub denom: String,
    pub price: Decimal256,
    pub decimals: u8,
}

#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub enum ExecuteMsg {
    SubmitPrices { prices: Vec<PriceEntry> },
    UpdateConfig {
        max_staleness: Option<u64>,
        twap_window: Option<u64>,
        min_publishers: Option<u32>,
        deviation_threshold: Option<Decimal256>,
        max_publishers: Option<u32>,
    },
    AddPublisher { address: String },
    RemovePublisher { address: String },
    SlashPublisher { address: String },
    RequestPriceUpdate { denom: String },
    SubmitSinglePrice {
        denom: String,
        price: Decimal256,
        decimals: u8,
    },
}

#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub enum QueryMsg {
    Config {},
    Price { denom: String },
    AllPrices {},
    Publisher { address: String },
    PriceHistory { denom: String, limit: Option<u32> },
}

#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct ConfigResponse {
    pub owner: String,
    pub max_staleness: u64,
    pub twap_window: u64,
    pub min_publishers: u32,
    pub deviation_threshold: String,
    pub max_publishers: u32,
}

#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct PriceResponse {
    pub denom: String,
    pub median_price: String,
    pub twap_price: String,
    pub publisher_count: u32,
    pub updated_at: u64,
    pub is_valid: bool,
}

#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct AllPricesResponse {
    pub prices: Vec<PriceResponse>,
}

#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct PublisherResponse {
    pub address: String,
    pub weight: u8,
    pub active: bool,
    pub joined_at: u64,
    pub total_submissions: u64,
    pub slashes: u8,
}

#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct PriceHistoryEntry {
    pub price: String,
    pub decimals: u8,
    pub timestamp: u64,
    pub publisher: String,
}

#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct PriceHistoryResponse {
    pub denom: String,
    pub history: Vec<PriceHistoryEntry>,
}

2.9 错误类型

// src/error.rs
use cosmwasm_std::StdError;
use thiserror::Error;

#[derive(Error, Debug, PartialEq)]
pub enum ContractError {
    #[error("{0}")]
    Std(#[from] StdError),

    #[error("Unauthorized")]
    Unauthorized {},

    #[error("Publisher not found")]
    PublisherNotFound {},

    #[error("Publisher already exists")]
    PublisherExists {},

    #[error("Publisher is inactive")]
    PublisherInactive {},

    #[error("Too many publishers")]
    TooManyPublishers {},

    #[error("Invalid price: must be positive")]
    InvalidPrice {},

    #[error("Price too old")]
    PriceTooOld {},

    #[error("Not enough publishers for aggregation")]
    InsufficientPublishers {},

    #[error("Price deviation exceeds threshold")]
    PriceDeviationExceeded {},

    #[error("Random request not found")]
    RandomRequestNotFound {},

    #[error("Random request already fulfilled")]
    RandomRequestAlreadyFulfilled {},
}

2.10 测试

// src/tests.rs
#[cfg(test)]
mod tests {
    use super::*;
    use cosmwasm_std::testing::{
        mock_dependencies, mock_env, mock_info,
    };
    use cosmwasm_std::{
        from_json, Addr, Coin, Decimal256, Timestamp, Uint128,
    };
    use crate::contract::{execute, instantiate, query};
    use crate::msg::{
        ExecuteMsg, InstantiateMsg, PriceEntry, QueryMsg, PriceResponse,
    };

    const MSG_PUBLISHER_1: &str = "msg1publisher1address";
    const MSG_PUBLISHER_2: &str = "msg1publisher2address";
    const MSG_PUBLISHER_3: &str = "msg1publisher3address";
    const MSG_OWNER: &str = "msg1owneraddress";

    fn create_valid_price(price: &str) -> Decimal256 {
        Decimal256::from_atomics(
            Uint128::from(price.parse::<u128>().unwrap()),
            6,
        ).unwrap()
    }

    fn setup_contract() -> (
        cosmwasm_std::OwnedDeps<
            cosmwasm_std::MemoryStorage,
            cosmwasm_std::testing::MockApi,
            cosmwasm_std::testing::MockQuerier,
        >, Env,
    ) {
        let mut deps = mock_dependencies();
        let env = mock_env();

        let msg = InstantiateMsg {
            initial_publishers: vec![
                MSG_PUBLISHER_1.to_string(),
                MSG_PUBLISHER_2.to_string(),
                MSG_PUBLISHER_3.to_string(),
            ],
            max_staleness: Some(300),
            twap_window: Some(3600),
            min_publishers: Some(2),
            deviation_threshold: Some(Decimal256::percent(5)),
            max_publishers: Some(20),
        };

        let info = mock_info(MSG_OWNER, &[]);
        instantiate(deps.as_mut(), env.clone(), info, msg).unwrap();
        (deps, env)
    }

    fn submit_prices(
        deps: &mut cosmwasm_std::OwnedDeps<
            cosmwasm_std::MemoryStorage,
            cosmwasm_std::testing::MockApi,
            cosmwasm_std::testing::MockQuerier,
        >,
        env: &Env,
        publisher: &str,
        prices: Vec<(&str, &str, u8)>,
    ) {
        let entries: Vec<PriceEntry> = prices
            .into_iter()
            .map(|(denom, price, decimals)| PriceEntry {
                denom: denom.to_string(),
                price: create_valid_price(price),
                decimals,
            })
            .collect();

        let msg = ExecuteMsg::SubmitPrices { prices: entries };
        let info = mock_info(publisher, &[]);
        execute(deps.as_mut(), env.clone(), info, msg).unwrap();
    }

    fn query_price(
        deps: &cosmwasm_std::OwnedDeps<
            cosmwasm_std::MemoryStorage,
            cosmwasm_std::testing::MockApi,
            cosmwasm_std::testing::MockQuerier,
        >,
        denom: &str,
    ) -> PriceResponse {
        let msg = QueryMsg::Price { denom: denom.to_string() };
        let bin = query(deps.as_ref(), mock_env(), msg).unwrap();
        from_json(&bin).unwrap()
    }

    #[test]
    fn test_instantiate() {
        let (deps, _) = setup_contract();
        let config: ConfigResponse =
            from_json(query(deps.as_ref(), mock_env(), QueryMsg::Config {}).unwrap()).unwrap();
        assert_eq!(config.owner, MSG_OWNER);
        assert_eq!(config.max_staleness, 300);
        assert_eq!(config.min_publishers, 2);
    }

    #[test]
    fn test_submit_and_query_price() {
        let (mut deps, env) = setup_contract();

        submit_prices(&mut deps, &env, MSG_PUBLISHER_1, vec![("MSG", "1500000", 6)]);
        submit_prices(&mut deps, &env, MSG_PUBLISHER_2, vec![("MSG", "1520000", 6)]);
        submit_prices(&mut deps, &env, MSG_PUBLISHER_3, vec![("MSG", "1490000", 6)]);

        let resp = query_price(&deps, "MSG");
        assert!(resp.is_valid);
        assert_eq!(resp.publisher_count, 3);
    }

    #[test]
    fn test_median_calculation() {
        let prices = vec![
            Decimal256::from_atomics(Uint128::new(100), 0).unwrap(),
            Decimal256::from_atomics(Uint128::new(200), 0).unwrap(),
            Decimal256::from_atomics(Uint128::new(300), 0).unwrap(),
        ];

        let mut sorted = prices.clone();
        sorted.sort();
        let median = sorted[1];
        assert_eq!(median, Decimal256::from_atomics(Uint128::new(200), 0).unwrap());
    }

    #[test]
    fn test_median_even_count() {
        let prices = vec![
            Decimal256::from_atomics(Uint128::new(100), 0).unwrap(),
            Decimal256::from_atomics(Uint128::new(200), 0).unwrap(),
            Decimal256::from_atomics(Uint128::new(300), 0).unwrap(),
            Decimal256::from_atomics(Uint128::new(400), 0).unwrap(),
        ];

        let mut sorted = prices.clone();
        sorted.sort();
        let mid = sorted.len() / 2;
        let avg = (sorted[mid - 1] + sorted[mid])
            .checked_div(Decimal256::from_atomics(Uint128::new(2), 0).unwrap())
            .unwrap();
        assert_eq!(avg, Decimal256::from_atomics(Uint128::new(250), 0).unwrap());
    }

    #[test]
    fn test_stale_price_rejected() {
        let (mut deps, mut env) = setup_contract();

        submit_prices(&mut deps, &env, MSG_PUBLISHER_1, vec![("MSG", "1500000", 6)]);

        env.block.time = Timestamp::from_seconds(env.block.time.seconds() + 600);

        submit_prices(&mut deps, &env, MSG_PUBLISHER_2, vec![("MSG", "1510000", 6)]);
        submit_prices(&mut deps, &env, MSG_PUBLISHER_3, vec![("MSG", "1490000", 6)]);

        let resp = query_price(&deps, "MSG");
        assert_eq!(resp.publisher_count, 2);
    }

    #[test]
    fn test_insufficient_publishers() {
        let (mut deps, env) = setup_contract();

        submit_prices(&mut deps, &env, MSG_PUBLISHER_1, vec![("MSG", "1500000", 6)]);

        let resp = query_price(&deps, "MSG");
        assert!(!resp.is_valid);
    }

    #[test]
    fn test_unauthorized_submission() {
        let (mut deps, env) = setup_contract();
        let unauthorized = "msg1unauthorized";

        let entry = PriceEntry {
            denom: "MSG".to_string(),
            price: create_valid_price("1500000"),
            decimals: 6,
        };

        let msg = ExecuteMsg::SubmitPrices { prices: vec![entry] };
        let info = mock_info(unauthorized, &[]);
        let result = execute(deps.as_mut(), env.clone(), info, msg);
        assert!(result.is_err());
    }

    #[test]
    fn test_add_publisher() {
        let (mut deps, env) = setup_contract();
        let new_publisher = "msg1newpublisher";

        let msg = ExecuteMsg::AddPublisher { address: new_publisher.to_string() };
        let info = mock_info(MSG_OWNER, &[]);
        execute(deps.as_mut(), env.clone(), info, msg).unwrap();

        let query_msg = QueryMsg::Publisher { address: new_publisher.to_string() };
        let bin = query(deps.as_ref(), env.clone(), query_msg).unwrap();
        let resp: PublisherResponse = from_json(&bin).unwrap();
        assert!(resp.active);
    }

    #[test]
    fn test_remove_publisher() {
        let (mut deps, env) = setup_contract();

        let msg = ExecuteMsg::RemovePublisher { address: MSG_PUBLISHER_1.to_string() };
        let info = mock_info(MSG_OWNER, &[]);
        execute(deps.as_mut(), env.clone(), info, msg).unwrap();

        let query_msg = QueryMsg::Publisher { address: MSG_PUBLISHER_1.to_string() };
        let result = query(deps.as_ref(), env.clone(), query_msg);
        assert!(result.is_err());
    }

    #[test]
    fn test_slash_publisher() {
        let (mut deps, env) = setup_contract();

        let msg = ExecuteMsg::SlashPublisher { address: MSG_PUBLISHER_2.to_string() };
        let info = mock_info(MSG_OWNER, &[]);

        execute(deps.as_mut(), env.clone(), info.clone(), msg.clone()).unwrap();

        let query_msg = QueryMsg::Publisher { address: MSG_PUBLISHER_2.to_string() };
        let bin = query(deps.as_ref(), env.clone(), query_msg).unwrap();
        let resp: PublisherResponse = from_json(&bin).unwrap();
        assert_eq!(resp.slashes, 1);
        assert!(resp.active);
    }

    #[test]
    fn test_slash_deactivates_after_three() {
        let (mut deps, env) = setup_contract();

        let msg = ExecuteMsg::SlashPublisher { address: MSG_PUBLISHER_2.to_string() };
        let info = mock_info(MSG_OWNER, &[]);

        for _ in 0..3 {
            execute(deps.as_mut(), env.clone(), info.clone(), msg.clone()).unwrap();
        }

        let query_msg = QueryMsg::Publisher { address: MSG_PUBLISHER_2.to_string() };
        let bin = query(deps.as_ref(), env.clone(), query_msg).unwrap();
        let resp: PublisherResponse = from_json(&bin).unwrap();
        assert_eq!(resp.slashes, 3);
        assert!(!resp.active);
    }

    #[test]
    fn test_twap_calculation() {
        let (mut deps, mut env) = setup_contract();

        submit_prices(&mut deps, &env, MSG_PUBLISHER_1, vec![("MSG", "1000000", 6)]);
        submit_prices(&mut deps, &env, MSG_PUBLISHER_2, vec![("MSG", "1000000", 6)]);

        env.block.time = Timestamp::from_seconds(env.block.time.seconds() + 100);

        submit_prices(&mut deps, &env, MSG_PUBLISHER_1, vec![("MSG", "2000000", 6)]);
        submit_prices(&mut deps, &env, MSG_PUBLISHER_2, vec![("MSG", "2000000", 6)]);

        let resp = query_price(&deps, "MSG");
        let twap: Decimal256 = resp.twap_price.parse().unwrap();
        let one = Decimal256::one();
        let two = Decimal256::from_atomics(Uint128::new(2), 0).unwrap();
        assert!(twap > one);
        assert!(twap < two);
    }

    #[test]
    fn test_multiple_denoms() {
        let (mut deps, env) = setup_contract();

        let prices = vec![
            ("MSG", "1500000", 6u8),
            ("USDC", "1000000", 6u8),
            ("ATOM", "10000000", 6u8),
            ("OSMO", "800000", 6u8),
        ];

        for (denom, price, decimals) in &prices {
            submit_prices(&mut deps, &env, MSG_PUBLISHER_1, vec![(denom, price, *decimals)]);
            submit_prices(&mut deps, &env, MSG_PUBLISHER_2, vec![(denom, price, *decimals)]);
            submit_prices(&mut deps, &env, MSG_PUBLISHER_3, vec![(denom, price, *decimals)]);
        }

        let query_msg = QueryMsg::AllPrices {};
        let bin = query(deps.as_ref(), env.clone(), query_msg).unwrap();
        let resp: AllPricesResponse = from_json(&bin).unwrap();
        assert_eq!(resp.prices.len(), 4);
    }

    #[test]
    fn test_price_history() {
        let (mut deps, env) = setup_contract();

        for i in 0..5 {
            let price_val = 1500000 + i * 10000;
            let price_str = price_val.to_string();
            submit_prices(&mut deps, &env, MSG_PUBLISHER_1, vec![("MSG", &price_str, 6)]);
            submit_prices(&mut deps, &env, MSG_PUBLISHER_2, vec![("MSG", &price_str, 6)]);
        }

        let query_msg = QueryMsg::PriceHistory { denom: "MSG".to_string(), limit: Some(3) };
        let bin = query(deps.as_ref(), env.clone(), query_msg).unwrap();
        let resp: PriceHistoryResponse = from_json(&bin).unwrap();
        assert_eq!(resp.history.len(), 3);
    }

    #[test]
    fn test_update_config_unauthorized() {
        let (mut deps, env) = setup_contract();
        let stranger = "msg1stranger";

        let msg = ExecuteMsg::UpdateConfig {
            max_staleness: Some(600),
            twap_window: None,
            min_publishers: None,
            deviation_threshold: None,
            max_publishers: None,
        };
        let info = mock_info(stranger, &[]);
        let result = execute(deps.as_mut(), env.clone(), info, msg);
        assert!(result.is_err());
    }

    #[test]
    fn test_price_with_different_decimals() {
        let (mut deps, env) = setup_contract();

        submit_prices(&mut deps, &env, MSG_PUBLISHER_1, vec![
            ("MSG", "1500000", 6),
            ("BTC", "5000000000", 8),
        ]);
        submit_prices(&mut deps, &env, MSG_PUBLISHER_2, vec![
            ("MSG", "1510000", 6),
            ("BTC", "5010000000", 8),
        ]);

        let msg_resp = query_price(&deps, "MSG");
        let btc_resp = query_price(&deps, "BTC");
        assert!(msg_resp.is_valid);
        assert!(btc_resp.is_valid);
    }

    #[test]
    fn test_submit_single_price() {
        let (mut deps, env) = setup_contract();

        let msg = ExecuteMsg::SubmitSinglePrice {
            denom: "ETH".to_string(),
            price: create_valid_price("250000000"),
            decimals: 8,
        };
        let info = mock_info(MSG_PUBLISHER_1, &[]);
        execute(deps.as_mut(), env.clone(), info, msg).unwrap();

        let msg2 = ExecuteMsg::SubmitSinglePrice {
            denom: "ETH".to_string(),
            price: create_valid_price("251000000"),
            decimals: 8,
        };
        let info2 = mock_info(MSG_PUBLISHER_2, &[]);
        execute(deps.as_mut(), env.clone(), info2, msg2).unwrap();

        let resp = query_price(&deps, "ETH");
        assert!(resp.is_valid);
    }

    #[test]
    fn test_zero_price_rejected() {
        let (mut deps, env) = setup_contract();

        let msg = ExecuteMsg::SubmitSinglePrice {
            denom: "MSG".to_string(),
            price: Decimal256::zero(),
            decimals: 6,
        };
        let info = mock_info(MSG_PUBLISHER_1, &[]);
        let result = execute(deps.as_mut(), env.clone(), info, msg);
        assert!(result.is_err());
    }

    #[test]
    fn test_request_price_update() {
        let (mut deps, env) = setup_contract();

        let msg = ExecuteMsg::RequestPriceUpdate { denom: "MSG".to_string() };
        let info = mock_info("msg1consumer", &[]);
        let resp = execute(deps.as_mut(), env.clone(), info, msg).unwrap();

        assert_eq!(resp.attributes[0].value, "request_price_update");
        assert_eq!(resp.attributes[1].value, "MSG");
    }
}

3. 链上随机数(VRF)

3.1 VRF 数据模型

// src/vrf_state.rs
use cosmwasm_std::{Addr, Binary, Storage};
use cw_storage_plus::{Item, Map};

/// 随机数请求
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, Eq)]
pub struct RandomRequest {
    pub height: u64,
    pub requester: Addr,
    pub request_type: RandomnessType,
    pub fulfilled: bool,
    pub seed: Option<Binary>,
    pub timestamp: u64,
    pub callback_addr: Option<Addr>,
    pub callback_msg: Option<Binary>,
}

/// 随机数响应
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, Eq)]
pub struct RandomResponse {
    pub request_id: u64,
    pub randomness: Binary,
    pub proof: Binary,
    pub timestamp: u64,
    pub source: RandomnessSource,
}

/// 随机数类型
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, Eq, JsonSchema)]
pub enum RandomnessType {
    General,
    Lottery,
    NftMint,
    ValidatorSelection,
}

/// 随机数来源
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, Eq, JsonSchema)]
pub enum RandomnessSource {
    Dar,
    Vrf,
}

/// VRF 配置
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, Eq)]
pub struct VrfConfig {
    pub owner: Addr,
    pub vrf_publishers: Vec<Addr>,
    pub dar_enabled: bool,
    pub max_request_timeout: u64,
    pub callback_gas_limit: u64,
}

pub const VRF_CONFIG: Item<VrfConfig> = Item::new("vrf_config");
pub const RANDOM_REQUESTS: Map<u64, RandomRequest> = Map::new("random_requests");
pub const RANDOM_RESPONSES: Map<u64, RandomResponse> = Map::new("random_responses");
pub const REQUEST_NONCE: Item<u64> = Item::new("request_nonce");
pub const DAR_SEED: Item<Binary> = Item::new("dar_seed");

3.2 VRF 合约入口

// src/vrf_contract.rs
use cosmwasm_std::{
    entry_point, to_json_binary, Binary, Deps, DepsMut, Env,
    MessageInfo, Response, StdError, StdResult,
};
use sha2::{Digest, Sha256};
use crate::error::ContractError;
use crate::vrf_msg::{ExecuteMsg, InstantiateMsg, QueryMsg, RandomnessResponse, RandomRequestResponse};
use crate::vrf_state::{
    RandomRequest, RandomResponse, RandomnessSource, RandomnessType,
    VrfConfig, DAR_SEED, RANDOM_REQUESTS, RANDOM_RESPONSES, REQUEST_NONCE, VRF_CONFIG,
};

/// 从区块哈希生成 DAR 随机数
pub fn derive_dar_randomness(env: &Env) -> StdResult<Binary> {
    let mut hasher = Sha256::new();
    hasher.update(env.block.hash.as_bytes());
    hasher.update(env.block.height.to_le_bytes());
    hasher.update(env.block.time.nanos().to_le_bytes());
    let result = hasher.finalize();
    Ok(Binary::from(result.to_vec()))
}

#[entry_point]
pub fn instantiate(
    deps: DepsMut,
    env: Env,
    info: MessageInfo,
    msg: InstantiateMsg,
) -> Result<Response, ContractError> {
    let config = VrfConfig {
        owner: info.sender.clone(),
        vrf_publishers: msg
            .initial_vrf_publishers
            .iter()
            .map(|a| deps.api.addr_validate(a))
            .collect::<StdResult<Vec<_>>>()?,
        dar_enabled: msg.dar_enabled.unwrap_or(true),
        max_request_timeout: msg.max_request_timeout.unwrap_or(100),
        callback_gas_limit: msg.callback_gas_limit.unwrap_or(500_000),
    };
    VRF_CONFIG.save(deps.storage, &config)?;
    REQUEST_NONCE.save(deps.storage, &0u64)?;

    let dar_seed = derive_dar_randomness(&env)?;
    DAR_SEED.save(deps.storage, &dar_seed)?;

    Ok(Response::new()
        .add_attribute("method", "instantiate_vrf")
        .add_attribute("owner", info.sender.to_string()))
}

#[entry_point]
pub fn execute(
    deps: DepsMut,
    env: Env,
    info: MessageInfo,
    msg: ExecuteMsg,
) -> Result<Response, ContractError> {
    match msg {
        ExecuteMsg::RequestRandomness { request_type, seed, callback_addr, callback_msg } => {
            execute_request_randomness(deps, env, info, request_type, seed, callback_addr, callback_msg)
        }
        ExecuteMsg::FulfillRandomness { request_id, randomness, proof } => {
            execute_fulfill_randomness(deps, env, info, request_id, randomness, proof)
        }
        ExecuteMsg::UpdateConfig { dar_enabled, max_request_timeout, callback_gas_limit } => {
            execute_vrf_update_config(deps, env, info, dar_enabled, max_request_timeout, callback_gas_limit)
        }
        ExecuteMsg::AddVrfPublisher { address } => {
            execute_add_vrf_publisher(deps, env, info, address)
        }
        ExecuteMsg::RemoveVrfPublisher { address } => {
            execute_remove_vrf_publisher(deps, env, info, address)
        }
    }
}

3.3 请求随机数

/// 请求随机数
pub fn execute_request_randomness(
    deps: DepsMut,
    env: Env,
    info: MessageInfo,
    request_type: Option<RandomnessType>,
    seed: Option<Binary>,
    callback_addr: Option<String>,
    callback_msg: Option<Binary>,
) -> Result<Response, ContractError> {
    let mut nonce = REQUEST_NONCE.load(deps.storage)?;
    nonce += 1;
    let request_id = nonce;
    REQUEST_NONCE.save(deps.storage, &nonce)?;

    let callback = callback_addr
        .map(|addr| deps.api.addr_validate(&addr))
        .transpose()?;

    let request = RandomRequest {
        height: env.block.height,
        requester: info.sender,
        request_type: request_type.unwrap_or(RandomnessType::General),
        fulfilled: false,
        seed,
        timestamp: env.block.time.seconds(),
        callback_addr: callback,
        callback_msg,
    };

    RANDOM_REQUESTS.save(deps.storage, request_id, &request)?;

    let config = VRF_CONFIG.load(deps.storage)?;
    let mut response = Response::new()
        .add_attribute("method", "request_randomness")
        .add_attribute("request_id", request_id.to_string())
        .add_attribute("requester", request.requester.to_string())
        .add_attribute("height", env.block.height.to_string());

    if config.dar_enabled {
        let dar_randomness = derive_dar_randomness(&env)?;
        let combined = combine_randomness(
            &dar_randomness,
            &request_id.to_le_bytes(),
            request.seed.as_deref(),
        );

        let dar_response = RandomResponse {
            request_id,
            randomness: combined.clone(),
            proof: Binary::from(b"dar_proof"),
            timestamp: env.block.time.seconds(),
            source: RandomnessSource::Dar,
        };

        RANDOM_RESPONSES.save(deps.storage, request_id, &dar_response)?;

        let mut fulfilled_request = request;
        fulfilled_request.fulfilled = true;
        RANDOM_REQUESTS.save(deps.storage, request_id, &fulfilled_request)?;

        response = response
            .add_attribute("randomness_source", "dar")
            .add_attribute("randomness", combined.to_base64())
            .add_attribute("fulfilled", "true");

        if let (Some(callback_addr), Some(callback_msg)) =
            (request.callback_addr, request.callback_msg)
        {
            response = response.add_submessage(cosmwasm_std::SubMsg::new(
                cosmwasm_std::CosmosMsg::Wasm(cosmwasm_std::WasmMsg::Execute {
                    contract_addr: callback_addr.to_string(),
                    msg: callback_msg,
                    funds: vec![],
                }),
            ));
        }
    } else {
        response = response
            .add_attribute("randomness_source", "vrf_pending")
            .add_attribute("fulfilled", "false");
    }

    Ok(response)
}

3.4 履行随机数(VRF 回退)

/// 履行随机数请求(由授权的 VRF 发布者调用)
pub fn execute_fulfill_randomness(
    deps: DepsMut,
    env: Env,
    info: MessageInfo,
    request_id: u64,
    randomness: Binary,
    proof: Binary,
) -> Result<Response, ContractError> {
    let config = VRF_CONFIG.load(deps.storage)?;

    if !config.vrf_publishers.contains(&info.sender) {
        return Err(ContractError::Unauthorized {});
    }

    let request = RANDOM_REQUESTS.load(deps.storage, request_id)?;

    if request.fulfilled {
        return Err(ContractError::RandomRequestAlreadyFulfilled {});
    }

    let proof_valid = verify_vrf_proof(&randomness, &proof, &info.sender);
    if !proof_valid {
        return Err(StdError::generic_err("Invalid VRF proof").into());
    }

    let response = RandomResponse {
        request_id,
        randomness: randomness.clone(),
        proof: proof.clone(),
        timestamp: env.block.time.seconds(),
        source: RandomnessSource::Vrf,
    };
    RANDOM_RESPONSES.save(deps.storage, request_id, &response)?;

    let mut fulfilled_request = request;
    fulfilled_request.fulfilled = true;
    RANDOM_REQUESTS.save(deps.storage, request_id, &fulfilled_request)?;

    let mut res = Response::new()
        .add_attribute("method", "fulfill_randomness")
        .add_attribute("request_id", request_id.to_string())
        .add_attribute("source", "vrf")
        .add_attribute("publisher", info.sender.to_string());

    if let (Some(callback_addr), Some(callback_msg)) =
        (request.callback_addr.clone(), request.callback_msg.clone())
    {
        let exec_msg = cosmwasm_std::WasmMsg::Execute {
            contract_addr: callback_addr.to_string(),
            msg: callback_msg,
            funds: vec![],
        };
        res = res.add_message(exec_msg);
    }

    Ok(res)
}

3.5 随机数组合与 VRF 验证

// src/vrf_crypto.rs
use sha2::{Digest, Sha256};
use cosmwasm_std::{Addr, Binary};
use hmac::{Hmac, Mac};

/// 组合多个随机数源
pub fn combine_randomness(
    dar_randomness: &Binary,
    request_id_bytes: &[u8],
    seed: Option<&Binary>,
) -> Binary {
    let mut hasher = Sha256::new();
    hasher.update(dar_randomness.as_slice());
    hasher.update(request_id_bytes);
    if let Some(s) = seed {
        hasher.update(s.as_slice());
    }
    let result = hasher.finalize();
    Binary::from(result.to_vec())
}

/// VRF 验证(HMAC-SHA256)
pub fn verify_vrf_proof(
    randomness: &Binary,
    proof: &Binary,
    publisher: &Addr,
) -> bool {
    type HmacSha256 = Hmac<Sha256>;
    let key = publisher.as_bytes();
    let mac = HmacSha256::new_from_slice(key).unwrap();
    // 在实际实现中应包含完整的 VRF 验证逻辑
    proof.as_slice() == mac.finalize().into_bytes().as_slice()
}

/// 从随机数生成范围内整数
pub fn random_number_in_range(randomness: &Binary, min: u64, max: u64) -> u64 {
    let range = max - min + 1;
    let bytes = randomness.as_slice();
    let mut buf = [0u8; 8];
    let len = bytes.len().min(8);
    buf[..len].copy_from_slice(&bytes[..len]);
    let value = u64::from_le_bytes(buf);
    let max_valid = u64::MAX - (u64::MAX % range);
    if value <= max_valid {
        min + (value % range)
    } else {
        let mut hasher = Sha256::new();
        hasher.update(bytes);
        let new_bytes = hasher.finalize();
        random_number_in_range(&Binary::from(new_bytes.to_vec()), min, max)
    }
}

/// 随机布尔值
pub fn random_bool(randomness: &Binary, true_probability: f64) -> bool {
    let value = random_number_in_range(randomness, 0, 10000);
    let threshold = (true_probability * 10000.0) as u64;
    value < threshold
}

#[cfg(test)]
mod vrf_crypto_tests {
    use super::*;

    #[test]
    fn test_random_number_in_range() {
        let randomness = Binary::from(&[1u8; 32]);
        let num = random_number_in_range(&randomness, 0, 100);
        assert!(num <= 100);

        let mut sum = 0u64;
        for i in 0..1000 {
            let r = Binary::from(&[(i as u8); 32]);
            sum += random_number_in_range(&r, 1, 6);
        }
        let avg = sum as f64 / 1000.0;
        assert!((avg - 3.5).abs() < 0.5);
    }

    #[test]
    fn test_random_bool() {
        let randomness = Binary::from(&[2u8; 32]);
        assert!(random_bool(&randomness, 1.0));
        assert!(!random_bool(&randomness, 0.0));
    }

    #[test]
    fn test_vrf_proof_verification() {
        let publisher = Addr::unchecked("msg1vrfpublisher");
        let randomness = Binary::from(&[3u8; 32]);
        type HmacSha256 = Hmac<Sha256>;
        let mut mac = HmacSha256::new_from_slice(publisher.as_bytes()).unwrap();
        mac.update(randomness.as_slice());
        let proof = Binary::from(mac.finalize().into_bytes().to_vec());
        assert!(verify_vrf_proof(&randomness, &proof, &publisher));

        let wrong_publisher = Addr::unchecked("msg1wrong");
        assert!(!verify_vrf_proof(&randomness, &proof, &wrong_publisher));
    }
}

3.6 VRF 查询与测试

// src/vrf_query.rs
use cosmwasm_std::{to_json_binary, Binary, Deps, StdResult};
use crate::vrf_msg::{QueryMsg, RandomnessResponse, RandomRequestResponse};
use crate::vrf_state::{RANDOM_REQUESTS, RANDOM_RESPONSES, VRF_CONFIG};

pub fn query_randomness(deps: Deps, request_id: u64) -> StdResult<Binary> {
    let response = RANDOM_RESPONSES.load(deps.storage, request_id)?;
    to_json_binary(&RandomnessResponse {
        request_id: response.request_id,
        randomness: response.randomness.to_base64(),
        proof: response.proof.to_base64(),
        timestamp: response.timestamp,
        source: format!("{:?}", response.source),
    })
}

pub fn query_random_request(deps: Deps, request_id: u64) -> StdResult<Binary> {
    let request = RANDOM_REQUESTS.load(deps.storage, request_id)?;
    to_json_binary(&RandomRequestResponse {
        request_id,
        height: request.height,
        requester: request.requester.to_string(),
        fulfilled: request.fulfilled,
        request_type: format!("{:?}", request.request_type),
        timestamp: request.timestamp,
    })
}

// src/vrf_msg.rs
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct InstantiateMsg {
    pub initial_vrf_publishers: Vec<String>,
    pub dar_enabled: Option<bool>,
    pub max_request_timeout: Option<u64>,
    pub callback_gas_limit: Option<u64>,
}

#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub enum ExecuteMsg {
    RequestRandomness {
        request_type: Option<RandomnessType>,
        seed: Option<Binary>,
        callback_addr: Option<String>,
        callback_msg: Option<Binary>,
    },
    FulfillRandomness {
        request_id: u64,
        randomness: Binary,
        proof: Binary,
    },
    UpdateConfig {
        dar_enabled: Option<bool>,
        max_request_timeout: Option<u64>,
        callback_gas_limit: Option<u64>,
    },
    AddVrfPublisher { address: String },
    RemoveVrfPublisher { address: String },
}

#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub enum QueryMsg {
    Randomness { request_id: u64 },
    RandomRequest { request_id: u64 },
}

#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct RandomnessResponse {
    pub request_id: u64,
    pub randomness: String,
    pub proof: String,
    pub timestamp: u64,
    pub source: String,
}

#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct RandomRequestResponse {
    pub request_id: u64,
    pub height: u64,
    pub requester: String,
    pub fulfilled: bool,
    pub request_type: String,
    pub timestamp: u64,
}

3.7 VRF 测试

#[cfg(test)]
mod vrf_tests {
    use super::*;
    use cosmwasm_std::testing::{mock_dependencies, mock_env, mock_info};
    use cosmwasm_std::{from_json, Addr, Timestamp};

    const VRF_PUBLISHER: &str = "msg1vrfpublisher";
    const VRF_OWNER: &str = "msg1vrfowner";

    fn setup_vrf() -> (
        cosmwasm_std::OwnedDeps<
            cosmwasm_std::MemoryStorage,
            cosmwasm_std::testing::MockApi,
            cosmwasm_std::testing::MockQuerier,
        >, Env,
    ) {
        let mut deps = mock_dependencies();
        let env = mock_env();
        let msg = InstantiateMsg {
            initial_vrf_publishers: vec![VRF_PUBLISHER.to_string()],
            dar_enabled: Some(true),
            max_request_timeout: Some(100),
            callback_gas_limit: Some(500_000),
        };
        let info = mock_info(VRF_OWNER, &[]);
        instantiate(deps.as_mut(), env.clone(), info, msg).unwrap();
        (deps, env)
    }

    #[test]
    fn test_request_randomness_dar() {
        let (mut deps, env) = setup_vrf();
        let msg = ExecuteMsg::RequestRandomness {
            request_type: Some(RandomnessType::General),
            seed: None, callback_addr: None, callback_msg: None,
        };
        let info = mock_info("msg1consumer", &[]);
        let resp = execute(deps.as_mut(), env.clone(), info, msg).unwrap();
        let fulfilled = resp.attributes.iter().find(|a| a.key == "fulfilled");
        assert_eq!(fulfilled.unwrap().value, "true");

        let query_msg = QueryMsg::Randomness { request_id: 1 };
        let bin = query(deps.as_ref(), env.clone(), query_msg).unwrap();
        let rand_resp: RandomnessResponse = from_json(&bin).unwrap();
        assert_eq!(rand_resp.request_id, 1);
    }

    #[test]
    fn test_request_randomness_vrf_fallback() {
        let (mut deps, env) = setup_vrf();
        let msg = ExecuteMsg::UpdateConfig {
            dar_enabled: Some(false), max_request_timeout: None, callback_gas_limit: None,
        };
        let info = mock_info(VRF_OWNER, &[]);
        execute(deps.as_mut(), env.clone(), info, msg).unwrap();

        let req_msg = ExecuteMsg::RequestRandomness {
            request_type: Some(RandomnessType::Lottery),
            seed: None, callback_addr: None, callback_msg: None,
        };
        let info = mock_info("msg1consumer", &[]);
        execute(deps.as_mut(), env.clone(), info, req_msg).unwrap();

        let randomness = Binary::from(&[42u8; 32]);
        type HmacSha256 = Hmac<Sha256>;
        let mut mac = HmacSha256::new_from_slice(VRF_PUBLISHER.as_bytes()).unwrap();
        mac.update(randomness.as_slice());
        let proof = Binary::from(mac.finalize().into_bytes().to_vec());

        let fulfill_msg = ExecuteMsg::FulfillRandomness {
            request_id: 1, randomness: randomness.clone(), proof,
        };
        let info = mock_info(VRF_PUBLISHER, &[]);
        execute(deps.as_mut(), env.clone(), info, fulfill_msg).unwrap();

        let query_msg = QueryMsg::Randomness { request_id: 1 };
        let bin = query(deps.as_ref(), env.clone(), query_msg).unwrap();
        let rand_resp: RandomnessResponse = from_json(&bin).unwrap();
        assert_eq!(rand_resp.source, "Vrf");
    }

    #[test]
    fn test_dar_derivation() {
        let env = mock_env();
        let randomness = derive_dar_randomness(&env).unwrap();
        assert_eq!(randomness.len(), 32);

        let randomness2 = derive_dar_randomness(&env).unwrap();
        assert_eq!(randomness, randomness2);
    }

    #[test]
    fn test_vrf_callback() {
        let (mut deps, env) = setup_vrf();
        let callback_addr = "msg1callbackcontract";
        let callback_msg = Binary::from(b"{\"action\":\"process_result\"}");

        let msg = ExecuteMsg::RequestRandomness {
            request_type: Some(RandomnessType::General),
            seed: None,
            callback_addr: Some(callback_addr.to_string()),
            callback_msg: Some(callback_msg.clone()),
        };
        let info = mock_info("msg1consumer", &[]);
        let resp = execute(deps.as_mut(), env.clone(), info, msg).unwrap();
        assert!(resp.messages.len() > 0);
    }

    #[test]
    fn test_double_fulfill_rejected() {
        let (mut deps, env) = setup_vrf();
        let config_msg = ExecuteMsg::UpdateConfig {
            dar_enabled: Some(false), max_request_timeout: None, callback_gas_limit: None,
        };
        let info = mock_info(VRF_OWNER, &[]);
        execute(deps.as_mut(), env.clone(), info, config_msg).unwrap();

        let req_msg = ExecuteMsg::RequestRandomness {
            request_type: None, seed: None, callback_addr: None, callback_msg: None,
        };
        let info = mock_info("msg1consumer", &[]);
        execute(deps.as_mut(), env.clone(), info, req_msg).unwrap();

        let randomness = Binary::from(&[1u8; 32]);
        type HmacSha256 = Hmac<Sha256>;
        let mut mac = HmacSha256::new_from_slice(VRF_PUBLISHER.as_bytes()).unwrap();
        mac.update(randomness.as_slice());
        let proof = Binary::from(mac.finalize().into_bytes().to_vec());

        let fulfill_msg = ExecuteMsg::FulfillRandomness {
            request_id: 1, randomness: randomness.clone(), proof: proof.clone(),
        };
        let info = mock_info(VRF_PUBLISHER, &[]);
        execute(deps.as_mut(), env.clone(), info.clone(), fulfill_msg.clone()).unwrap();

        let result = execute(deps.as_mut(), env.clone(), info, fulfill_msg);
        assert!(result.is_err());
    }

    #[test]
    fn test_unauthorized_vrf_fulfill() {
        let (mut deps, env) = setup_vrf();
        let config_msg = ExecuteMsg::UpdateConfig {
            dar_enabled: Some(false), max_request_timeout: None, callback_gas_limit: None,
        };
        let info = mock_info(VRF_OWNER, &[]);
        execute(deps.as_mut(), env.clone(), info, config_msg).unwrap();

        let req_msg = ExecuteMsg::RequestRandomness {
            request_type: None, seed: None, callback_addr: None, callback_msg: None,
        };
        let info = mock_info("msg1consumer", &[]);
        execute(deps.as_mut(), env.clone(), info, req_msg).unwrap();

        let randomness = Binary::from(&[1u8; 32]);
        let proof = Binary::from(b"fake_proof");
        let fulfill_msg = ExecuteMsg::FulfillRandomness { request_id: 1, randomness, proof };
        let info = mock_info("msg1unauthorized", &[]);
        let result = execute(deps.as_mut(), env.clone(), info, fulfill_msg);
        assert!(result.is_err());
    }

    #[test]
    fn test_random_request_query() {
        let (mut deps, env) = setup_vrf();
        let msg = ExecuteMsg::RequestRandomness {
            request_type: Some(RandomnessType::Lottery),
            seed: Some(Binary::from(b"custom_seed")),
            callback_addr: None, callback_msg: None,
        };
        let info = mock_info("msg1consumer", &[]);
        execute(deps.as_mut(), env.clone(), info, msg).unwrap();

        let query_msg = QueryMsg::RandomRequest { request_id: 1 };
        let bin = query(deps.as_ref(), env.clone(), query_msg).unwrap();
        let resp: RandomRequestResponse = from_json(&bin).unwrap();
        assert!(resp.fulfilled);
    }

    #[test]
    fn test_random_number_integration() {
        let (mut deps, env) = setup_vrf();
        let msg = ExecuteMsg::RequestRandomness {
            request_type: Some(RandomnessType::Lottery),
            seed: None, callback_addr: None, callback_msg: None,
        };
        let info = mock_info("msg1consumer", &[]);
        execute(deps.as_mut(), env.clone(), info, msg).unwrap();

        let query_msg = QueryMsg::Randomness { request_id: 1 };
        let bin = query(deps.as_ref(), env.clone(), query_msg).unwrap();
        let resp: RandomnessResponse = from_json(&bin).unwrap();
        let random_bytes = Binary::from_base64(&resp.randomness).unwrap();
        let dice = random_number_in_range(&random_bytes, 1, 6);
        assert!(dice >= 1 && dice <= 6);
    }
}

4. 数据预言机

4.1 数据预言机架构

// src/data_oracle/state.rs
use cosmwasm_std::{Addr, Binary, Storage};
use cw_storage_plus::{Item, Map};

#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, Eq)]
pub struct DataSource {
    pub source_id: String,
    pub name: String,
    pub signers: Vec<Addr>,
    pub min_signatures: u8,
    pub active: bool,
    pub reputation: u8,
    pub total_submissions: u64,
    pub failed_submissions: u64,
}

#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, Eq)]
pub struct DataPoint {
    pub data_id: String,
    pub source_id: String,
    pub data_type: DataType,
    pub value: Binary,
    pub decimals: u8,
    pub timestamp: u64,
    pub signatures: Vec<Signature>,
    pub verified: bool,
}

#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, Eq, JsonSchema)]
pub enum DataType {
    Numeric, String, Json, Boolean, Bytes,
}

#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, Eq)]
pub struct Signature {
    pub signer: Addr,
    pub signature: Binary,
    pub timestamp: u64,
}

#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, Eq)]
pub struct DataOracleConfig {
    pub owner: Addr,
    pub max_staleness: u64,
    pub max_sources: u32,
}

pub const DATA_CONFIG: Item<DataOracleConfig> = Item::new("data_config");
pub const DATA_SOURCES: Map<&str, DataSource> = Map::new("data_sources");
pub const DATA_POINTS: Map<&str, DataPoint> = Map::new("data_points");
pub const DATA_HISTORY: Map<(&str, u64), DataPoint> = Map::new("data_history");

4.2 数据预言机合约

// src/data_oracle/contract.rs
use cosmwasm_std::{
    entry_point, to_json_binary, Binary, Deps, DepsMut, Env,
    MessageInfo, Response, StdError, StdResult,
};
use sha2::{Digest, Sha256};
use crate::error::ContractError;
use crate::data_oracle::msg::{ExecuteMsg, InstantiateMsg, QueryMsg};
use crate::data_oracle::state::{
    DataOracleConfig, DataPoint, DataSource, DataType, Signature,
    DATA_CONFIG, DATA_HISTORY, DATA_POINTS, DATA_SOURCES,
};

#[entry_point]
pub fn instantiate(
    deps: DepsMut,
    _env: Env,
    info: MessageInfo,
    msg: InstantiateMsg,
) -> Result<Response, ContractError> {
    let config = DataOracleConfig {
        owner: info.sender.clone(),
        max_staleness: msg.max_staleness.unwrap_or(600),
        max_sources: msg.max_sources.unwrap_or(50),
    };
    DATA_CONFIG.save(deps.storage, &config)?;
    Ok(Response::new()
        .add_attribute("method", "instantiate_data_oracle")
        .add_attribute("owner", info.sender.to_string()))
}

#[entry_point]
pub fn execute(
    deps: DepsMut,
    env: Env,
    info: MessageInfo,
    msg: ExecuteMsg,
) -> Result<Response, ContractError> {
    match msg {
        ExecuteMsg::RegisterDataSource { source_id, name, signers, min_signatures } => {
            execute_register_data_source(deps, env, info, source_id, name, signers, min_signatures)
        }
        ExecuteMsg::SubmitData { data_id, source_id, data_type, value, decimals, signatures } => {
            execute_submit_data(deps, env, info, data_id, source_id, data_type, value, decimals, signatures)
        }
        ExecuteMsg::UpdateDataSource { source_id, name, active } => {
            execute_update_data_source(deps, env, info, source_id, name, active)
        }
    }
}

4.3 数据提交与验证

/// 注册数据源
pub fn execute_register_data_source(
    deps: DepsMut,
    _env: Env,
    info: MessageInfo,
    source_id: String,
    name: String,
    signers: Vec<String>,
    min_signatures: u8,
) -> Result<Response, ContractError> {
    let config = DATA_CONFIG.load(deps.storage)?;
    if info.sender != config.owner {
        return Err(ContractError::Unauthorized {});
    }
    if DATA_SOURCES.has(deps.storage, &source_id) {
        return Err(StdError::generic_err("Data source already exists").into());
    }
    let validated_signers: Vec<Addr> = signers
        .iter()
        .map(|s| deps.api.addr_validate(s))
        .collect::<StdResult<_>>()?;

    let source = DataSource {
        source_id: source_id.clone(),
        name,
        signers: validated_signers,
        min_signatures,
        active: true,
        reputation: 100,
        total_submissions: 0,
        failed_submissions: 0,
    };
    DATA_SOURCES.save(deps.storage, &source_id, &source)?;
    Ok(Response::new()
        .add_attribute("method", "register_data_source")
        .add_attribute("source_id", source_id))
}

/// 提交数据
pub fn execute_submit_data(
    deps: DepsMut,
    env: Env,
    _info: MessageInfo,
    data_id: String,
    source_id: String,
    data_type: DataType,
    value: Binary,
    decimals: u8,
    signatures: Vec<Signature>,
) -> Result<Response, ContractError> {
    let source = DATA_SOURCES.load(deps.storage, &source_id)?;
    if !source.active {
        return Err(StdError::generic_err("Data source is inactive").into());
    }
    if (signatures.len() as u8) < source.min_signatures {
        return Err(StdError::generic_err(format!(
            "Insufficient signatures: got {}, need {}", signatures.len(), source.min_signatures
        )).into());
    }

    let data_hash = hash_data(&data_id, &source_id, &value, decimals);
    let mut valid_signatures = 0u8;
    for sig in &signatures {
        if source.signers.contains(&sig.signer) {
            if verify_data_signature(&data_hash, &sig.signature, &sig.signer) {
                valid_signatures += 1;
            }
        }
    }

    if valid_signatures < source.min_signatures {
        let mut updated_source = source;
        updated_source.failed_submissions += 1;
        if updated_source.reputation > 0 {
            updated_source.reputation -= 5;
        }
        DATA_SOURCES.save(deps.storage, &source_id, &updated_source)?;
        return Err(StdError::generic_err(format!(
            "Signature verification failed: valid={}, required={}", valid_signatures, source.min_signatures
        )).into());
    }

    let data_point = DataPoint {
        data_id: data_id.clone(),
        source_id: source_id.clone(),
        data_type,
        value: value.clone(),
        decimals,
        timestamp: env.block.time.seconds(),
        signatures,
        verified: true,
    };
    DATA_POINTS.save(deps.storage, &data_id, &data_point)?;
    DATA_HISTORY.save(deps.storage, (&data_id, env.block.time.seconds()), &data_point)?;

    let mut updated_source = source;
    updated_source.total_submissions += 1;
    DATA_SOURCES.save(deps.storage, &source_id, &updated_source)?;

    Ok(Response::new()
        .add_attribute("method", "submit_data")
        .add_attribute("data_id", data_id)
        .add_attribute("source_id", source_id)
        .add_attribute("verified", "true"))
}

4.4 签名验证

// src/data_oracle/crypto.rs
use cosmwasm_std::{Addr, Binary};
use sha2::{Digest, Sha256};
use hmac::{Hmac, Mac};

/// 哈希数据(用于签名验证)
pub fn hash_data(data_id: &str, source_id: &str, value: &Binary, decimals: u8) -> Vec<u8> {
    let mut hasher = Sha256::new();
    hasher.update(data_id.as_bytes());
    hasher.update(source_id.as_bytes());
    hasher.update(value.as_slice());
    hasher.update([decimals]);
    hasher.finalize().to_vec()
}

/// 验证数据签名(HMAC-SHA256)
pub fn verify_data_signature(data_hash: &[u8], signature: &Binary, signer: &Addr) -> bool {
    type HmacSha256 = Hmac<Sha256>;
    let key = signer.as_bytes();
    let mut mac = HmacSha256::new_from_slice(key).unwrap();
    mac.update(data_hash);
    mac.finalize().into_bytes().as_slice() == signature.as_slice()
}

// src/data_oracle/msg.rs
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct InstantiateMsg {
    pub max_staleness: Option<u64>,
    pub max_sources: Option<u32>,
}

#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub enum ExecuteMsg {
    RegisterDataSource {
        source_id: String,
        name: String,
        signers: Vec<String>,
        min_signatures: u8,
    },
    SubmitData {
        data_id: String,
        source_id: String,
        data_type: DataType,
        value: Binary,
        decimals: u8,
        signatures: Vec<SignatureJson>,
    },
    UpdateDataSource {
        source_id: String,
        name: Option<String>,
        active: Option<bool>,
    },
}

#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct SignatureJson {
    pub signer: String,
    pub signature: String,
    pub timestamp: u64,
}

#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub enum QueryMsg {
    DataPoint { data_id: String },
    DataSource { source_id: String },
}

#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct DataPointResponse {
    pub data_id: String,
    pub source_id: String,
    pub data_type: String,
    pub value: String,
    pub decimals: u8,
    pub timestamp: u64,
    pub verified: bool,
}

#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct DataSourceResponse {
    pub source_id: String,
    pub name: String,
    pub active: bool,
    pub reputation: u8,
    pub total_submissions: u64,
}

5. 预言机消费者模式

5.1 消费者合约

// src/consumer/state.rs
use cosmwasm_std::{Addr, Decimal256};
use cw_storage_plus::{Item, Map};

#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, Eq)]
pub struct ConsumerConfig {
    pub owner: Addr,
    pub price_oracle_addr: Addr,
    pub vrf_addr: Option<Addr>,
    pub data_oracle_addr: Option<Addr>,
    pub max_price_age: u64,
    pub price_tolerance: Decimal256,
    pub fallback_to_twap: bool,
}

#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, Eq)]
pub struct CachedPrice {
    pub price: Decimal256,
    pub timestamp: u64,
    pub source: PriceSource,
}

#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, Eq, JsonSchema)]
pub enum PriceSource { Oracle, Twap, Fallback }

#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, Eq)]
pub struct DecisionLog {
    pub id: u64,
    pub decision_type: String,
    pub data_used: String,
    pub result: String,
    pub timestamp: u64,
}

pub const CONSUMER_CONFIG: Item<ConsumerConfig> = Item::new("consumer_config");
pub const CACHED_PRICES: Map<&str, CachedPrice> = Map::new("cached_prices");
pub const DECISION_HISTORY: Map<u64, DecisionLog> = Map::new("decision_history");
pub const DECISION_NONCE: Item<u64> = Item::new("decision_nonce");

// src/consumer/msg.rs
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct InstantiateMsg {
    pub price_oracle_addr: String,
    pub vrf_addr: Option<String>,
    pub data_oracle_addr: Option<String>,
    pub max_price_age: Option<u64>,
    pub price_tolerance: Option<Decimal256>,
    pub fallback_to_twap: Option<bool>,
}

#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub enum ExecuteMsg {
    ConsumePrice { denom: String },
    ConsumeRandomness { seed: Option<Binary> },
    ConsumeData { data_id: String },
    ExecuteTradeDecision { denom: String, threshold: Decimal256 },
    UpdateConsumerConfig {
        max_price_age: Option<u64>,
        price_tolerance: Option<Decimal256>,
        fallback_to_twap: Option<bool>,
    },
}

#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub enum QueryMsg {
    Config {},
    CachedPrice { denom: String },
    DecisionHistory { limit: Option<u32> },
}

5.2 安全消费价格

// src/consumer/contract.rs
use cosmwasm_std::{
    entry_point, to_json_binary, Binary, Decimal256, Deps, DepsMut,
    Env, MessageInfo, Response, StdError, StdResult, WasmMsg,
    WasmQuery, QueryRequest,
};
use crate::error::ContractError;
use crate::consumer::msg::{ExecuteMsg, InstantiateMsg, QueryMsg};
use crate::consumer::state::{
    CachedPrice, ConsumerConfig, DecisionLog, PriceSource,
    CACHED_PRICES, CONSUMER_CONFIG, DECISION_HISTORY, DECISION_NONCE,
};

/// 消费价格的核心函数
pub fn execute_consume_price(
    deps: DepsMut,
    env: Env,
    _info: MessageInfo,
    denom: String,
) -> Result<Response, ContractError> {
    let config = CONSUMER_CONFIG.load(deps.storage)?;

    let price_response: crate::msg::PriceResponse = deps.querier.query(
        &QueryRequest::Wasm(WasmQuery::Smart {
            contract_addr: config.price_oracle_addr.to_string(),
            msg: to_json_binary(&crate::msg::QueryMsg::Price {
                denom: denom.clone(),
            })?,
        }),
    )?;

    let now = env.block.time.seconds();

    if now - price_response.updated_at > config.max_price_age {
        if config.fallback_to_twap {
            let twap_price: Decimal256 = price_response.twap_price.parse().map_err(|_| {
                StdError::generic_err("Invalid TWAP price format")
            })?;

            if !twap_price.is_zero() {
                CACHED_PRICES.save(deps.storage, &denom, &CachedPrice {
                    price: twap_price,
                    timestamp: now,
                    source: PriceSource::Twap,
                })?;

                return Ok(Response::new()
                    .add_attribute("method", "consume_price")
                    .add_attribute("denom", &denom)
                    .add_attribute("source", "twap")
                    .add_attribute("price", twap_price.to_string()));
            }
        }
        return Err(ContractError::PriceTooOld {});
    }

    if !price_response.is_valid {
        return Err(StdError::generic_err("Price is invalid").into());
    }

    let median_price: Decimal256 = price_response.median_price.parse().map_err(|_| {
        StdError::generic_err("Invalid median price format")
    })?;

    CACHED_PRICES.save(deps.storage, &denom, &CachedPrice {
        price: median_price,
        timestamp: now,
        source: PriceSource::Oracle,
    })?;

    Ok(Response::new()
        .add_attribute("method", "consume_price")
        .add_attribute("denom", denom)
        .add_attribute("source", "oracle")
        .add_attribute("price", median_price.to_string())
        .add_attribute("updated_at", price_response.updated_at.to_string()))
}

/// 消费随机数
pub fn execute_consume_randomness(
    deps: DepsMut,
    env: Env,
    info: MessageInfo,
    seed: Option<Binary>,
) -> Result<Response, ContractError> {
    let config = CONSUMER_CONFIG.load(deps.storage)?;
    let vrf_addr = config.vrf_addr
        .ok_or_else(|| StdError::generic_err("VRF contract not configured"))?;

    let request_msg = crate::vrf_msg::ExecuteMsg::RequestRandomness {
        request_type: Some(crate::vrf_state::RandomnessType::General),
        seed,
        callback_addr: Some(env.contract.address.to_string()),
        callback_msg: Some(to_json_binary(&ExecuteMsg::ConsumeRandomness { seed: None })?),
    };

    let msg = WasmMsg::Execute {
        contract_addr: vrf_addr.to_string(),
        msg: to_json_binary(&request_msg)?,
        funds: vec![],
    };

    Ok(Response::new()
        .add_message(msg)
        .add_attribute("method", "request_randomness")
        .add_attribute("requester", info.sender.to_string()))
}

5.3 交易决策与回退

/// 基于价格执行交易决策
pub fn execute_trade_decision(
    deps: DepsMut,
    env: Env,
    info: MessageInfo,
    denom: String,
    threshold: Decimal256,
) -> Result<Response, ContractError> {
    let consume_result = execute_consume_price(deps.branch(), env.clone(), info.clone(), denom.clone());

    let price = match consume_result {
        Ok(_) => CACHED_PRICES.load(deps.storage, &denom)?.price,
        Err(_) => {
            let cached = CACHED_PRICES.may_load(deps.storage, &denom)?
                .ok_or_else(|| StdError::generic_err("No cached price available"))?;

            let config = CONSUMER_CONFIG.load(deps.storage)?;
            let age = env.block.time.seconds() - cached.timestamp;
            if age > config.max_price_age * 2 {
                return execute_conservative_fallback(deps, env, denom, threshold);
            }
            cached.price
        }
    };

    let decision = if price > threshold { "BUY" }
    else if price < threshold { "SELL" }
    else { "HOLD" };

    let mut nonce = DECISION_NONCE.load(deps.storage)?;
    nonce += 1;
    DECISION_NONCE.save(deps.storage, &nonce)?;

    DECISION_HISTORY.save(deps.storage, nonce, &DecisionLog {
        id: nonce,
        decision_type: "trade".to_string(),
        data_used: format!("price={}", price),
        result: decision.to_string(),
        timestamp: env.block.time.seconds(),
    })?;

    Ok(Response::new()
        .add_attribute("method", "trade_decision")
        .add_attribute("denom", denom)
        .add_attribute("price", price.to_string())
        .add_attribute("threshold", threshold.to_string())
        .add_attribute("decision", decision))
}

/// 保守回退策略
fn execute_conservative_fallback(
    deps: DepsMut,
    env: Env,
    denom: String,
    _threshold: Decimal256,
) -> Result<Response, ContractError> {
    let mut nonce = DECISION_NONCE.load(deps.storage)?;
    nonce += 1;
    DECISION_NONCE.save(deps.storage, &nonce)?;

    DECISION_HISTORY.save(deps.storage, nonce, &DecisionLog {
        id: nonce,
        decision_type: "trade_fallback".to_string(),
        data_used: "conservative".to_string(),
        result: "HOLD".to_string(),
        timestamp: env.block.time.seconds(),
    })?;

    Ok(Response::new()
        .add_attribute("method", "conservative_fallback")
        .add_attribute("denom", denom)
        .add_attribute("decision", "HOLD")
        .add_attribute("reason", "price_unavailable"))
}

6. 安全考虑

6.1 发布者信誉系统

// src/security/reputation.rs
use cosmwasm_std::{Addr, Decimal256, StdResult, Storage};
use crate::state::{PublisherInfo, PUBLISHERS};

#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, Eq)]
pub struct ReputationScore {
    pub publisher: Addr,
    pub base_score: u8,
    pub accuracy_score: u8,
    pub uptime_score: u8,
    pub total_score: u8,
}

/// 计算发布者的综合信誉得分
pub fn calculate_reputation(
    storage: &dyn Storage,
    publisher: &Addr,
    deviation: Decimal256,
    max_allowed_deviation: Decimal256,
) -> StdResult<ReputationScore> {
    let info = PUBLISHERS.load(storage, publisher)?;

    let base_score = if info.active {
        100u8.saturating_sub(info.slashes * 20)
    } else { 0 };

    let accuracy_score = if deviation <= max_allowed_deviation {
        let ratio = deviation / max_allowed_deviation;
        100u8 - ((ratio * Decimal256::from_ratio(40u128, 1u128))
            .to_uint_floor().try_into().unwrap_or(0u8))
    } else { 50 };

    let uptime_score = if info.total_submissions > 1000 { 100 }
    else if info.total_submissions > 500 { 80 }
    else if info.total_submissions > 100 { 60 }
    else { 40 };

    let total = ((base_score as u16 * 3) + (accuracy_score as u16 * 4) + (uptime_score as u16 * 3)) / 10;

    Ok(ReputationScore {
        publisher: publisher.clone(),
        base_score, accuracy_score, uptime_score,
        total_score: total as u8,
    })
}

6.2 价格操纵防护

// src/security/price_guard.rs
use cosmwasm_std::{Decimal256, StdResult, Storage};
use crate::state::{Config, PRICES, CONFIG};

/// 检测价格操纵
pub fn detect_price_manipulation(
    storage: &dyn Storage,
    denom: &str,
    new_price: Decimal256,
    publisher_count: u32,
) -> StdResult<bool> {
    let config = CONFIG.load(storage)?;

    if publisher_count < config.min_publishers.max(3) {
        return Ok(true);
    }

    if let Some(aggregated) = PRICES.may_load(storage, denom)? {
        if !aggregated.twap_price.is_zero() {
            let deviation = if new_price > aggregated.twap_price {
                (new_price - aggregated.twap_price) / aggregated.twap_price
            } else {
                (aggregated.twap_price - new_price) / aggregated.twap_price
            };
            if deviation > config.deviation_threshold {
                return Ok(true);
            }
        }
    }
    Ok(false)
}

6.3 活性监控

// src/security/monitor.rs
use cosmwasm_std::{Addr, StdResult, Storage, Order};
use crate::state::{PublisherInfo, PUBLISHERS, CONFIG};

#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, Eq)]
pub struct LivenessReport {
    pub active_publishers: u32,
    pub total_publishers: u32,
    pub health_score: u8,
    pub stale_publishers: Vec<Addr>,
    pub warnings: Vec<String>,
}

pub fn check_system_liveness(storage: &dyn Storage, current_time: u64) -> StdResult<LivenessReport> {
    let config = CONFIG.load(storage)?;
    let mut active_count = 0u32;
    let mut total_count = 0u32;
    let mut stale = Vec::new();
    let mut warnings = Vec::new();

    let publishers: Vec<PublisherInfo> = PUBLISHERS
        .range(storage, None, None, Order::Ascending)
        .filter_map(|r| r.ok())
        .map(|(_, p)| p)
        .collect();

    total_count = publishers.len() as u32;

    for publisher in &publishers {
        if publisher.active { active_count += 1; }
        if current_time - publisher.joined_at > config.max_staleness * 10 {
            stale.push(publisher.address.clone());
        }
    }

    let health_score = if total_count == 0 { 0 }
    else { ((active_count * 100) / total_count) as u8 };

    if active_count < config.min_publishers {
        warnings.push(format!("Active publishers ({}) below minimum ({})", active_count, config.min_publishers));
    }
    if !stale.is_empty() {
        warnings.push(format!("Stale publishers: {}", stale.len()));
    }

    Ok(LivenessReport { active_publishers: active_count, total_publishers: total_count, health_score, stale_publishers: stale, warnings })
}

7. AI Agent 集成

7.1 Agent 价格查询

// src/ai_agent/price_query.rs
use cosmwasm_std::{Deps, StdResult, Decimal256, to_json_binary, WasmQuery, QueryRequest};

pub struct AgentPriceQuerier {
    pub oracle_addr: String,
}

impl AgentPriceQuerier {
    pub fn new(oracle_addr: &str) -> Self {
        Self { oracle_addr: oracle_addr.to_string() }
    }

    pub fn query_price(&self, deps: &Deps, denom: &str) -> StdResult<AgentPriceResponse> {
        let response: crate::msg::PriceResponse = deps.querier.query(
            &QueryRequest::Wasm(WasmQuery::Smart {
                contract_addr: self.oracle_addr.clone(),
                msg: to_json_binary(&crate::msg::QueryMsg::Price {
                    denom: denom.to_string(),
                })?,
            }),
        )?;

        Ok(AgentPriceResponse {
            denom: denom.to_string(),
            median_price: response.median_price.parse().unwrap_or(Decimal256::zero()),
            twap_price: response.twap_price.parse().unwrap_or(Decimal256::zero()),
            is_valid: response.is_valid,
            updated_at: response.updated_at,
        })
    }

    pub fn query_batch_prices(&self, deps: &Deps, denoms: &[&str]) -> StdResult<Vec<AgentPriceResponse>> {
        denoms.iter().filter_map(|d| {
            self.query_price(deps, d).ok()
        }).collect()
    }

    pub fn query_all_market_prices(&self, deps: &Deps) -> StdResult<Vec<AgentPriceResponse>> {
        let response: crate::msg::AllPricesResponse = deps.querier.query(
            &QueryRequest::Wasm(WasmQuery::Smart {
                contract_addr: self.oracle_addr.clone(),
                msg: to_json_binary(&crate::msg::QueryMsg::AllPrices {})?,
            }),
        )?;

        response.prices.into_iter().map(|p| {
            Ok(AgentPriceResponse {
                denom: p.denom,
                median_price: p.median_price.parse().unwrap_or(Decimal256::zero()),
                twap_price: p.twap_price.parse().unwrap_or(Decimal256::zero()),
                is_valid: p.is_valid,
                updated_at: p.updated_at,
            })
        }).collect()
    }
}

#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct AgentPriceResponse {
    pub denom: String,
    pub median_price: Decimal256,
    pub twap_price: Decimal256,
    pub is_valid: bool,
    pub updated_at: u64,
}

7.2 Agent 使用 VRF

// src/ai_agent/vrf_decision.rs
use cosmwasm_std::{Addr, Binary, Decimal256, DepsMut, Env, MessageInfo, Response, StdError, StdResult};
use sha2::{Digest, Sha256};
use crate::vrf_state::RandomnessSource;

#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, Eq, JsonSchema)]
pub enum DecisionMode { Deterministic, Probabilistic, Hybrid }

#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, Eq)]
pub struct AgentVrfConfig {
    pub agent_addr: Addr,
    pub decision_mode: DecisionMode,
    pub randomness_threshold: Decimal256,
    pub vrf_contract: Addr,
}

/// AI Agent 随机化决策
pub fn execute_agent_random_decision(
    deps: DepsMut,
    env: Env,
    info: MessageInfo,
    options: Vec<AgentDecision>,
) -> Result<Response, ContractError> {
    let request_id = env.block.height;
    let randomness_query = crate::vrf_msg::QueryMsg::Randomness { request_id };

    let response: crate::vrf_msg::RandomnessResponse = deps.querier.query(
        &QueryRequest::Wasm(WasmQuery::Smart {
            contract_addr: "msg1vrf".to_string(),
            msg: to_json_binary(&randomness_query)?,
        }),
    ).map_err(|_| StdError::generic_err("Failed to query randomness"))?;

    let random_bytes = Binary::from_base64(&response.randomness)
        .map_err(|_| StdError::generic_err("Invalid randomness format"))?;

    let decision_index = crate::vrf_crypto::random_number_in_range(
        &random_bytes, 0, (options.len() - 1) as u64,
    ) as usize;

    let selected = &options[decision_index];

    Ok(Response::new()
        .add_attribute("method", "agent_random_decision")
        .add_attribute("agent", info.sender.to_string())
        .add_attribute("decision_id", selected.id.to_string())
        .add_attribute("decision_action", &selected.action)
        .add_attribute("randomness_source", response.source))
}

#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct AgentDecision {
    pub id: u64,
    pub action: String,
    pub params: Binary,
    pub weight: u8,
}

7.3 Agent 发布预言机数据

// src/ai_agent/publisher.rs
use cosmwasm_std::{Addr, Binary, DepsMut, Env, MessageInfo, Response, StdError, StdResult};
use hmac::{Hmac, Mac};
use sha2::Sha256;

/// AI Agent 作为数据发布者
pub fn execute_agent_publish_analysis(
    deps: DepsMut,
    env: Env,
    info: MessageInfo,
    data_oracle_addr: String,
    analysis: AiAnalysis,
) -> Result<Response, ContractError> {
    let value = serde_json::to_vec(&analysis).map_err(|e| {
        StdError::generic_err(format!("Serialization error: {}", e))
    })?;

    type HmacSha256 = Hmac<Sha256>;
    let mut mac = HmacSha256::new_from_slice(info.sender.as_bytes()).unwrap();
    mac.update(&value);
    let signature = mac.finalize().into_bytes().to_vec();

    let signature_json = crate::data_oracle::msg::SignatureJson {
        signer: info.sender.to_string(),
        signature: Binary::from(signature).to_base64(),
        timestamp: env.block.time.seconds(),
    };

    let submit_msg = crate::data_oracle::msg::ExecuteMsg::SubmitData {
        data_id: format!("analysis_{}", env.block.height),
        source_id: "ai_agent_analysis".to_string(),
        data_type: crate::data_oracle::state::DataType::Json,
        value: Binary::from(value),
        decimals: 0,
        signatures: vec![signature_json],
    };

    let wasm_msg = cosmwasm_std::WasmMsg::Execute {
        contract_addr: data_oracle_addr,
        msg: to_json_binary(&submit_msg)?,
        funds: vec![],
    };

    Ok(Response::new()
        .add_message(wasm_msg)
        .add_attribute("method", "agent_publish_analysis")
        .add_attribute("agent", info.sender.to_string())
        .add_attribute("height", env.block.height.to_string()))
}

/// AI 分析结果
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct AiAnalysis {
    pub analysis_type: String,
    pub target: String,
    pub score: u8,
    pub confidence: Decimal256,
    pub recommendation: String,
    pub reasoning: String,
    pub model_version: String,
    pub analyzed_at: u64,
}

7.4 自动化策略

// src/ai_agent/strategy.rs
use cosmwasm_std::{Decimal256, Deps, StdResult, Addr};

/// 套利策略
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, Eq)]
pub struct ArbitrageStrategy {
    pub agent_addr: Addr,
    pub min_profit_bps: u16,
    pub max_slippage_bps: u16,
    pub max_position_size: Decimal256,
}

impl ArbitrageStrategy {
    pub fn evaluate_arbitrage(
        &self, _deps: &Deps,
        pool_a_price: Decimal256,
        pool_b_price: Decimal256,
    ) -> StdResult<Option<ArbitrageOpportunity>> {
        if pool_a_price.is_zero() || pool_b_price.is_zero() {
            return Ok(None);
        }
        let profit_ratio = if pool_a_price > pool_b_price {
            (pool_a_price - pool_b_price) / pool_b_price
        } else {
            (pool_b_price - pool_a_price) / pool_a_price
        };
        let profit_bps = (profit_ratio * Decimal256::from_ratio(10000u128, 1u128))
            .to_uint_floor().try_into().unwrap_or(0u16);

        if profit_bps >= self.min_profit_bps {
            Ok(Some(ArbitrageOpportunity {
                buy_pool: if pool_a_price < pool_b_price { "A" } else { "B" },
                sell_pool: if pool_a_price > pool_b_price { "A" } else { "B" },
                expected_profit_bps: profit_bps,
                estimated_gas_cost_bps: 5,
            }))
        } else { Ok(None) }
    }
}

#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, Eq)]
pub struct ArbitrageOpportunity {
    pub buy_pool: &'static str,
    pub sell_pool: &'static str,
    pub expected_profit_bps: u16,
    pub estimated_gas_cost_bps: u16,
}

/// 做市策略
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, Eq)]
pub struct MarketMakingStrategy {
    pub agent_addr: Addr,
    pub spread_bps: u16,
    pub order_size: Decimal256,
    pub max_inventory: Decimal256,
    pub rebalance_threshold: Decimal256,
}

impl MarketMakingStrategy {
    pub fn calculate_quotes(
        &self, mid_price: Decimal256, current_inventory: Decimal256,
    ) -> StdResult<(Decimal256, Decimal256)> {
        let spread_ratio = Decimal256::from_ratio(self.spread_bps as u128, 10000u128);
        let half_spread = spread_ratio / Decimal256::from_ratio(2u128, 1u128);
        let inventory_ratio = if self.max_inventory.is_zero() {
            Decimal256::zero()
        } else { current_inventory / self.max_inventory };
        let adjustment = inventory_ratio * half_spread;
        Ok((
            mid_price * (Decimal256::one() - half_spread - adjustment),
            mid_price * (Decimal256::one() + half_spread - adjustment),
        ))
    }
}

7.5 Agent 测试

#[cfg(test)]
mod ai_agent_tests {
    use super::*;
    use cosmwasm_std::testing::{mock_dependencies, mock_env, mock_info};
    use cosmwasm_std::{from_json, Uint128};

    #[test]
    fn test_ai_analysis_serialization() {
        let analysis = AiAnalysis {
            analysis_type: "sentiment".to_string(),
            target: "MSG".to_string(),
            score: 75,
            confidence: Decimal256::from_ratio(8u128, 10u128),
            recommendation: "buy".to_string(),
            reasoning: "Strong upward trend".to_string(),
            model_version: "v2.1.0".to_string(),
            analyzed_at: 1_700_000_000,
        };
        let serialized = serde_json::to_vec(&analysis).unwrap();
        let deserialized: AiAnalysis = serde_json::from_slice(&serialized).unwrap();
        assert_eq!(deserialized.recommendation, "buy");
    }

    #[test]
    fn test_arbitrage_evaluation() {
        let strategy = ArbitrageStrategy {
            agent_addr: Addr::unchecked("msg1agent"),
            min_profit_bps: 10,
            max_slippage_bps: 50,
            max_position_size: Decimal256::from_atomics(Uint128::new(1000), 6).unwrap(),
        };
        let deps = mock_dependencies();

        let pool_a = Decimal256::from_atomics(Uint128::new(100), 0).unwrap();
        let pool_b = Decimal256::from_atomics(Uint128::new(105), 0).unwrap();
        let opportunity = strategy.evaluate_arbitrage(&deps.as_ref(), pool_a, pool_b).unwrap();
        assert!(opportunity.is_some());

        let pool_c = Decimal256::from_atomics(Uint128::new(100), 0).unwrap();
        let pool_d = Decimal256::from_atomics(Uint128::new(100), 0).unwrap();
        let opportunity = strategy.evaluate_arbitrage(&deps.as_ref(), pool_c, pool_d).unwrap();
        assert!(opportunity.is_none());
    }

    #[test]
    fn test_market_making_quotes() {
        let strategy = MarketMakingStrategy {
            agent_addr: Addr::unchecked("msg1mmagent"),
            spread_bps: 50,
            order_size: Decimal256::from_atomics(Uint128::new(100), 6).unwrap(),
            max_inventory: Decimal256::from_atomics(Uint128::new(10000), 6).unwrap(),
            rebalance_threshold: Decimal256::from_ratio(5u128, 10u128),
        };
        let mid_price = Decimal256::from_atomics(Uint128::new(100), 0).unwrap();
        let low_inventory = Decimal256::from_atomics(Uint128::new(1000), 6).unwrap();
        let (bid, ask) = strategy.calculate_quotes(mid_price, low_inventory).unwrap();
        assert!(bid < mid_price);
        assert!(ask > mid_price);
        assert!(bid < ask);
    }

    #[test]
    fn test_agent_vrf_decision_selection() {
        let options = vec![
            AgentDecision { id: 1, action: "buy".to_string(), params: Binary::from(b"{\"amount\":100}"), weight: 1 },
            AgentDecision { id: 2, action: "sell".to_string(), params: Binary::from(b"{\"amount\":50}"), weight: 1 },
            AgentDecision { id: 3, action: "hold".to_string(), params: Binary::from(b"{}"), weight: 1 },
        ];
        let randomness = Binary::from(&[42u8; 32]);
        let index = crate::vrf_crypto::random_number_in_range(&randomness, 0, 2) as usize;
        assert!(index < options.len());
    }
}

8. 附录

8.1 Cargo.toml 依赖

[package]
name = "msgchain-oracle"
version = "2.0.0"
edition = "2021"

[lib]
crate-type = ["cdylib", "rlib"]

[dependencies]
cosmwasm-std = "2.1"
cosmwasm-schema = "2.1"
cw-storage-plus = "2.0"
cw2 = "2.0"
thiserror = "2"
serde = { version = "1", features = ["derive"] }
schemars = "0.8"
sha2 = "0.10"
hmac = "0.12"
serde_json = "1"

[dev-dependencies]
cosmwasm-vm = "2.1"

8.2 部署命令

# 编译合约
docker run --rm -v "$(pwd)":/code \
  --mount type=volume,source="$(basename "$(pwd)")_cache",target=/target \
  --mount type=volume,source=registry_cache,target=/usr/local/cargo/registry \
  cosmwasm/optimizer:0.16.0

# 部署价格预言机
msgd tx wasm store artifacts/msgchain_oracle.wasm \
  --from wallet --gas auto --gas-adjustment 1.3 -y

# 实例化
msgd tx wasm instantiate <CODE_ID> \
  '{"initial_publishers":["msg1publisher1","msg1publisher2"]}' \
  --from wallet --label "price-oracle" --admin <ADMIN_ADDR> -y

# 查询价格
msgd query wasm contract-state smart <CONTRACT_ADDR> \
  '{"price":{"denom":"MSG"}}'

8.3 架构概览

┌──────────────────────────────────────────────────────────┐
│                    MSG Chain 预言机系统                      │
├──────────────┬─────────────────┬──────────────────────────┤
│ 价格预言机     │   VRF 随机数     │   数据预言机              │
│ (PriceOracle) │ (Randomness)    │   (DataOracle)          │
├──────────────┼─────────────────┼──────────────────────────┤
│ · 多发布者聚合  │ · DAR 共识随机数  │ · 多数据源注册            │
│ · 中位数策略    │ · VRF 预言机回退  │ · 签名验证               │
│ · TWAP 计算   │ · 回调机制       │ · 信誉系统               │
│ · 过期检测     │ · 无偏范围生成    │ · 数据历史               │
└──────────────┴─────────────────┴──────────────────────────┘
                          │
                    ┌─────▼──────┐
                    │  消费者合约   │
                    │ (Consumer)  │
                    ├────────────┤
                    │ · 价格消费    │
                    │ · 随机数消费   │
                    │ · 交易决策    │
                    │ · 回退策略    │
                    └─────┬──────┘
                          │
                    ┌─────▼──────┐
                    │  AI Agent  │
                    ├────────────┤
                    │ · 价格查询    │
                    │ · VRF 决策   │
                    │ · 数据发布    │
                    │ · 策略执行    │
                    └────────────┘

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