MSG Chain 预言机集成与价格数据指南
主网状态: No-Go
目录
1. 概述
1.1 为什么预言机至关重要
区块链是确定性封闭系统,无法原生访问链外数据。预言机(Oracle)桥接了链上与链下世界,使得智能合约能够获取真实世界的信息。在 MSG Chain(msg 前缀地址)上,预言机的核心用途包括:
| 用例 | 描述 | 典型场景 |
|---|---|---|
| 价格数据 | 获取资产对美元的实时汇率 | DEX 定价、借贷清算、衍生品 |
| 链上随机数 | 生成不可预测的安全随机数 | 游戏抽奖、NFT 铸造、验证者选取 |
| 任意数据 | 任何链下数据源 | 天气数据、体育比分、AI 模型输出 |
1.2 预言机设计模型
推模型(Push Model)
预言机节点主动将数据提交到链上,数据按固定时间间隔或事件触发更新。
[数据源] → [预言机节点] → [MSG Chain 合约] → [消费者合约]
- 优点:数据即时可用,读取无延迟
- 缺点:Gas 成本高(每次推送都需要链上写入),数据可能过时
- 适合:DeFi 价格预言机(如 Band Protocol 标准模型)
拉模型(Pull Model)
消费者合约在需要时主动请求数据,预言机节点仅在被调用时提交数据。
[消费者合约] → [预言机合约(请求)] → [预言机节点(回应)] → [消费者合约(消费)]
- 优点:按需付费,Gas 成本更低
- 缺点:需要等待一个区块确认,延迟较高
- 适合:低频率查询、AI 决策
混合模型
结合推拉模型的优势:价格数据采用推模型高频更新,而随机数和自定义数据采用拉模型。
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/
