MSG Chain AI Agent 数据市场实现指南
链信息:
msg-chain-1| bech32 前缀:msg| 精度: 18 decimals | Gas: 1,000,000,000 attoMSG/gas⚠️ No-Go Disclaimer: MSGChain 主网裁决为 No-Go。本文件所有内容反映的是开发阶段的技术设计,不代表主网未独立核验上线状态。生产部署状态请以白皮书为准:https://msgchain.org/whitepaper/
目录
1. 概述
1.1 为什么 AI Agent 需要数据市场
AI Agent 的核心能力取决于其训练数据和推理数据的质量与多样性。在 MSG Chain 的生态中,Agent 已经能够通过 agent_registry_v1 注册发现、通过 agent_payment_v1 (AIPAY) 支付、通过 agent_a2a_v1 通信协作。然而,数据仍然是割裂的——每个 Agent 各自为战,无法高效共享和交易数据。
数据市场 (Data Marketplace) 解决了以下关键问题:
- 数据孤岛: Agent A 拥有高质量标注数据集,Agent B 需要它来训练模型,但没有交易通道
- 价值变现: Agent 可以将自身产生的推理结果、中间表示、微调权重作为商品出售
- 协作训练: 多个 Agent 贡献数据联合训练,按贡献分配收益
- 数据质量保证: 链上哈希承诺 + 链下验证机制确保数据完整性和真实性
1.2 可交易的数据类型
| 类型 | 描述 | 典型用例 |
|---|---|---|
| 训练数据 | 原始文本、图像、语音等训练样本 | Agent 训练专用领域模型 |
| 标注数据集 | 带标签的监督学习数据 | 分类、实体识别、情感分析 |
| 推理结果 | Agent 处理后的输出数据 | 作为另一 Agent 的输入 |
| Embeddings | 向量嵌入表示 | 检索增强生成 (RAG) |
| 微调权重 | LoRA/全量微调后的模型权重 | 迁移学习、领域适应 |
| 合成数据 | Agent 生成的合成训练数据 | 数据增强、隐私保护 |
| 评估数据 | Benchmark 评测集 | Agent 性能评估 |
1.3 现有基础设施复用
| 模块 | 用途 | 数据市场集成点 |
|---|---|---|
agent_registry_v1 |
Agent 发现与身份 | 数据资产关联到 Agent ID,买方 Agent 发现数据提供者 |
agent_payment_v1 (AIPAY) |
条件支付 | 数据交易的核心支付通道,支持条件释放 |
agent_a2a_v1 |
Agent 间通信 | 数据协商、数据传输信令、仲裁通信 |
ai_agent_constitution_v1 |
Agent 行为约束 | 数据交易中 Agent 必须遵守的规则(禁止转售、隐私保护等) |
1.4 架构总览
+-------------------------------------------------------------+
| Data Marketplace |
+-------------------------------------------------------------+
| +-------------+ +-------------+ +-------------------+ |
| | DataAsset | | DataTrading | | AccessControl | |
| | Contract | | Module | | Module | |
| +-------------+ +-------------+ +-------------------+ |
| +-------------+ +-------------+ +-------------------+ |
| | Verification| | Reputation | | Royalty/Revenue | |
| | & Arbitration| | System | | Sharing | |
| +-------------+ +-------------+ +-------------------+ |
+-------------------------------------------------------------+
| MSG Chain Infrastructure |
| agent_registry_v1 | agent_payment_v1 | agent_a2a_v1 |
| ai_agent_constitution_v1 |
+-------------------------------------------------------------+
1.5 核心流程
- 数据提供方 注册数据资产 → 提交哈希承诺 → 设定价格和许可证
- 数据消费方 发现数据 → 发起购买(AIPAY 条件支付)
- 系统锁定资金至托管合约 → 提供方发送数据解密密钥/访问链接
- 消费方验证数据哈希 → 匹配则确认收货,资金释放
- 争议处理 → 哈希不匹配则进入仲裁流程
1.6 经济模型
数据市场使用 MSG 链原生代币 uMSG(18 位精度)进行交易:
- 交易费用: 协议收取 1% 手续费(可配置)
- Gas 消耗:
- 数据资产注册: ~200,000 gas (0.02 MSG @ 0.01 gas price)
- 数据购买: ~300,000 gas (0.03 MSG @ 0.01 gas price)
- 确认交付: ~150,000 gas (0.015 MSG @ 0.01 gas price)
- 托管机制: 买方资金锁定在合约中,验证通过后释放给卖方
2. 数据资产合约
2.1 合约概述
data_marketplace 是数据市场的核心合约,基于 CosmWasm 构建,部署在 MSG Chain 上。它负责:
- 数据资产的注册、更新、注销
- 价格和许可证管理
- 交易状态跟踪
- 托管资金管理
- 访问控制验证
2.2 完整合约代码 (Rust)
2.2.1 Cargo.toml
[package]
name = "data-marketplace"
version = "0.1.0"
edition = "2021"
description = "MSG Chain AI Agent Data Marketplace"
[lib]
crate-type = ["cdylib", "rlib"]
[features]
default = ["library"]
library = []
[dependencies]
cosmwasm-std = { version = "1.5", features = ["staking"] }
cosmwasm-schema = "1.5"
cosmwasm-crypto = "1.5"
cw-storage-plus = "1.2"
cw-utils = "1.0"
cw2 = "1.1"
schemars = "0.8"
serde = { version = "1.0", features = ["derive"] }
thiserror = "1.0"
hex = "0.4"
sha2 = "0.10"
uint = "0.9"
[dev-dependencies]
cosmwasm-vm = "1.5"
cw-multi-test = "1.0"
2.2.2 src/state.rs
use cosmwasm_std::{Addr, Coin, Timestamp, Uint128};
use cw_storage_plus::{Item, Map};
use schemars::JsonSchema;
use serde::{Deserialize, Serialize};
/// 访问控制类型
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub enum AccessType {
Public,
TokenGated { token_address: Addr, min_balance: Uint128 },
Whitelist,
Subscription { price_per_block: Coin, min_blocks: u64 },
NftGated { nft_contract: Addr, token_id: Option<String> },
}
/// 许可证类型
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub enum License {
MIT,
CCBY40,
CCNC40,
TrainingOnly,
NoResale,
Custom { uri: String },
}
/// 数据验证状态
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub enum VerificationStatus {
Unverified,
Verified,
Failed,
Disputed,
}
/// 数据资产核心结构
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct DataAsset {
pub id: Uint128,
pub owner: Addr,
pub data_hash: String,
pub metadata_uri: String,
pub price: Coin,
pub license: License,
pub access_type: AccessType,
pub data_size: Uint128,
pub sample_count: Option<Uint128>,
pub data_type: String,
pub version: String,
pub created_at: Timestamp,
pub updated_at: Timestamp,
pub active: bool,
pub total_sales: Uint128,
pub total_revenue: Coin,
pub quality_score: u8,
pub verification_status: VerificationStatus,
pub agent_id: Option<String>,
}
/// 交易记录
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct TradeRecord {
pub trade_id: Uint128,
pub asset_id: Uint128,
pub buyer: Addr,
pub seller: Addr,
pub amount: Coin,
pub timestamp: Timestamp,
pub status: TradeStatus,
pub escrow_address: Option<Addr>,
pub completed_at: Option<Timestamp>,
}
/// 交易状态
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub enum TradeStatus {
PendingPayment,
PaidAwaitingDelivery,
DeliveredAwaitingConfirmation,
Completed,
Disputed,
Refunded,
Cancelled,
}
/// 托管账户
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct EscrowAccount {
pub escrow_id: Uint128,
pub trade_id: Uint128,
pub depositor: Addr,
pub beneficiary: Addr,
pub amount: Coin,
pub condition: String,
pub status: EscrowStatus,
pub created_at: Timestamp,
pub expires_at: Option<Timestamp>,
pub arbitrator: Option<Addr>,
}
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub enum EscrowStatus {
Locked,
Released,
Refunded,
Disputed,
}
/// 数据访问授权
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct AccessGrant {
pub grant_id: Uint128,
pub asset_id: Uint128,
pub grantee: Addr,
pub grant_type: AccessGrantType,
pub expires_at: Option<Timestamp>,
pub created_at: Timestamp,
}
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub enum AccessGrantType {
Permanent,
TimeBound { blocks: u64 },
UsageBased { max_downloads: u32, downloads_used: u32 },
}
#### 2.2.3 src/msg.rs
```rust
use cosmwasm_std::{Addr, Coin, Timestamp, Uint128};
use schemars::JsonSchema;
use serde::{Deserialize, Serialize};
use crate::state::{AccessGrantType, AccessType, License, TradeStatus, VerificationStatus};
/// 合约实例化消息
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct InstantiateMsg {
pub admin: Addr,
pub name: String,
pub fee_rate_bps: u64,
pub fee_collector: Addr,
pub max_arbitration_blocks: u64,
pub min_price: Coin,
pub max_data_size: Uint128,
pub registry_address: Option<Addr>,
pub aipay_address: Option<Addr>,
}
/// 执行消息
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
#[serde(rename_all = "snake_case")]
pub enum ExecuteMsg {
RegisterAsset {
data_hash: String,
metadata_uri: String,
price: Coin,
license: License,
access_type: AccessType,
data_size: Uint128,
sample_count: Option<Uint128>,
data_type: String,
version: String,
agent_id: Option<String>,
},
UpdateAsset {
asset_id: Uint128,
metadata_uri: Option<String>,
price: Option<Coin>,
license: Option<License>,
access_type: Option<AccessType>,
version: Option<String>,
},
DeactivateAsset { asset_id: Uint128 },
PurchaseData { asset_id: Uint128 },
ConfirmDelivery { trade_id: Uint128, data_hash_verified: String },
DisputeDelivery { trade_id: Uint128, reason: String },
RequestRefund { trade_id: Uint128 },
Arbitrate { trade_id: Uint128, resolve_in_favor_of_buyer: bool },
AddToWhitelist { asset_id: Uint128, addresses: Vec<Addr> },
RemoveFromWhitelist { asset_id: Uint128, addresses: Vec<Addr> },
GrantAccess { asset_id: Uint128, grantee: Addr, grant_type: AccessGrantType },
RevokeAccess { grant_id: Uint128 },
UpdateQualityScore { asset_id: Uint128, score: u8 },
WithdrawFees { denom: String, amount: Option<Uint128> },
}
/// 查询消息
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
#[serde(rename_all = "snake_case")]
pub enum QueryMsg {
GetAsset { asset_id: Uint128 },
ListAssets {
start_after: Option<Uint128>,
limit: Option<u32>,
owner: Option<Addr>,
data_type: Option<String>,
min_price: Option<Coin>,
max_price: Option<Coin>,
license: Option<License>,
active_only: Option<bool>,
},
GetTrade { trade_id: Uint128 },
ListTrades {
start_after: Option<Uint128>,
limit: Option<u32>,
buyer: Option<Addr>,
seller: Option<Addr>,
asset_id: Option<Uint128>,
status: Option<TradeStatus>,
},
GetEscrow { escrow_id: Uint128 },
GetAssetsByOwner { owner: Addr, start_after: Option<Uint128>, limit: Option<u32> },
SearchAssets { query: String, start_after: Option<Uint128>, limit: Option<u32> },
CheckAccess { asset_id: Uint128, address: Addr },
GetStats {},
GetSubscription { subscription_id: Uint128 },
ListSubscriptions {
start_after: Option<Uint128>,
limit: Option<u32>,
subscriber: Option<Addr>,
asset_id: Option<Uint128>,
},
}
/// 查询响应结构
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct AssetResponse { pub asset: DataAsset }
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct AssetListResponse {
pub assets: Vec<DataAsset>,
pub total_count: Uint128,
}
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct TradeResponse { pub trade: TradeRecord }
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct TradeListResponse {
pub trades: Vec<TradeRecord>,
pub total_count: Uint128,
}
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct EscrowResponse { pub escrow: EscrowAccount }
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct AccessResponse {
pub has_access: bool,
pub grant: Option<AccessGrantInfo>,
}
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct AccessGrantInfo {
pub grant_id: Uint128,
pub asset_id: Uint128,
pub grant_type: AccessGrantType,
pub expires_at: Option<Timestamp>,
pub is_valid: bool,
}
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct StatsResponse {
pub total_assets: Uint128,
pub active_assets: Uint128,
pub total_trades: Uint128,
pub total_volume: Coin,
pub total_fees_collected: Coin,
pub unique_buyers: Uint128,
pub unique_sellers: Uint128,
}
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct SubscriptionResponse {
pub subscription: Subscription,
}
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct SubscriptionListResponse {
pub subscriptions: Vec<Subscription>,
pub total_count: Uint128,
}
use crate::state::{DataAsset, TradeRecord, EscrowAccount, Subscription};
2.2.4 src/lib.rs
pub mod contract;
pub mod error;
pub mod msg;
pub mod state;
pub mod access;
pub mod verification;
pub use crate::error::ContractError;
pub type ContractResult<T> = Result<T, ContractError>;
2.2.5 src/contract.rs (核心业务逻辑)
use cosmwasm_std::{
entry_point, to_json_binary, Binary, Coin, Deps, DepsMut, Env, MessageInfo, Response,
StdError, StdResult, Uint128,
};
use schemars::JsonSchema;
use serde::{Deserialize, Serialize};
use crate::error::ContractError;
use crate::msg::{
AssetListResponse, AssetResponse, EscrowResponse, ExecuteMsg, InstantiateMsg, QueryMsg,
StatsResponse, TradeListResponse, TradeResponse,
};
use crate::state::{
AccessGrant, AccessGrantType, AccessType, DataAsset, EscrowAccount, EscrowStatus, License,
Subscription, TradeRecord, TradeStatus, VerificationStatus,
ACCESS_GRANTS, ASSETS, CONFIG, ESCROWS, SUBSCRIPTIONS, TRADES, WHITELISTS,
};
use crate::ContractResult;
/// 合约配置
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct Config {
pub admin: Addr,
pub name: String,
pub fee_rate_bps: u64,
pub fee_collector: Addr,
pub max_arbitration_blocks: u64,
pub min_price: Coin,
pub max_data_size: Uint128,
pub registry_address: Option<Addr>,
pub aipay_address: Option<Addr>,
pub asset_count: Uint128,
pub trade_count: Uint128,
pub escrow_count: Uint128,
pub grant_count: Uint128,
pub subscription_count: Uint128,
}
pub const CONFIG: Item<Config> = Item::new("config");
pub const ASSETS: Map<Uint128, DataAsset> = Map::new("assets");
pub const TRADES: Map<Uint128, TradeRecord> = Map::new("trades");
pub const ESCROWS: Map<Uint128, EscrowAccount> = Map::new("escrows");
pub const WHITELISTS: Map<(Uint128, Addr), bool> = Map::new("whitelists");
pub const ACCESS_GRANTS: Map<Uint128, AccessGrant> = Map::new("access_grants");
pub const SUBSCRIPTIONS: Map<Uint128, Subscription> = Map::new("subscriptions");
const MAX_LIMIT: u32 = 100;
const DEFAULT_LIMIT: u32 = 20;
// === 入口点 ===
#[entry_point]
pub fn instantiate(
deps: DepsMut,
_env: Env,
_info: MessageInfo,
msg: InstantiateMsg,
) -> ContractResult<Response> {
let config = Config {
admin: msg.admin,
name: msg.name,
fee_rate_bps: msg.fee_rate_bps,
fee_collector: msg.fee_collector,
max_arbitration_blocks: msg.max_arbitration_blocks,
min_price: msg.min_price,
max_data_size: msg.max_data_size,
registry_address: msg.registry_address,
aipay_address: msg.aipay_address,
asset_count: Uint128::zero(),
trade_count: Uint128::zero(),
escrow_count: Uint128::zero(),
grant_count: Uint128::zero(),
subscription_count: Uint128::zero(),
};
CONFIG.save(deps.storage, &config)?;
Ok(Response::new()
.add_attribute("action", "instantiate")
.add_attribute("name", config.name))
}
#[entry_point]
pub fn execute(
deps: DepsMut,
env: Env,
info: MessageInfo,
msg: ExecuteMsg,
) -> ContractResult<Response> {
match msg {
ExecuteMsg::RegisterAsset { data_hash, metadata_uri, price, license, access_type, data_size, sample_count, data_type, version, agent_id } =>
execute_register_asset(deps, env, info, data_hash, metadata_uri, price, license, access_type, data_size, sample_count, data_type, version, agent_id),
ExecuteMsg::UpdateAsset { asset_id, metadata_uri, price, license, access_type, version } =>
execute_update_asset(deps, env, info, asset_id, metadata_uri, price, license, access_type, version),
ExecuteMsg::DeactivateAsset { asset_id } => execute_deactivate_asset(deps, env, info, asset_id),
ExecuteMsg::PurchaseData { asset_id } => execute_purchase_data(deps, env, info, asset_id),
ExecuteMsg::ConfirmDelivery { trade_id, data_hash_verified } => execute_confirm_delivery(deps, env, info, trade_id, data_hash_verified),
ExecuteMsg::DisputeDelivery { trade_id, reason } => execute_dispute_delivery(deps, env, info, trade_id, reason),
ExecuteMsg::RequestRefund { trade_id } => execute_request_refund(deps, env, info, trade_id),
ExecuteMsg::Arbitrate { trade_id, resolve_in_favor_of_buyer } => execute_arbitrate(deps, env, info, trade_id, resolve_in_favor_of_buyer),
ExecuteMsg::AddToWhitelist { asset_id, addresses } => execute_add_to_whitelist(deps, env, info, asset_id, addresses),
ExecuteMsg::RemoveFromWhitelist { asset_id, addresses } => execute_remove_from_whitelist(deps, env, info, asset_id, addresses),
ExecuteMsg::GrantAccess { asset_id, grantee, grant_type } => execute_grant_access(deps, env, info, asset_id, grantee, grant_type),
ExecuteMsg::RevokeAccess { grant_id } => execute_revoke_access(deps, env, info, grant_id),
ExecuteMsg::UpdateQualityScore { asset_id, score } => execute_update_quality_score(deps, env, info, asset_id, score),
ExecuteMsg::WithdrawFees { denom, amount } => execute_withdraw_fees(deps, env, info, denom, amount),
}
}
#[entry_point]
pub fn query(deps: Deps, _env: Env, msg: QueryMsg) -> StdResult<Binary> {
match msg {
QueryMsg::GetAsset { asset_id } => to_json_binary(&query_asset(deps, asset_id)?),
QueryMsg::ListAssets { start_after, limit, owner, data_type, min_price, max_price, license, active_only } =>
to_json_binary(&query_list_assets(deps, start_after, limit, owner, data_type, min_price, max_price, license, active_only)?),
QueryMsg::GetTrade { trade_id } => to_json_binary(&query_trade(deps, trade_id)?),
QueryMsg::ListTrades { start_after, limit, buyer, seller, asset_id, status } =>
to_json_binary(&query_list_trades(deps, start_after, limit, buyer, seller, asset_id, status)?),
QueryMsg::GetEscrow { escrow_id } => to_json_binary(&query_escrow(deps, escrow_id)?),
QueryMsg::GetAssetsByOwner { owner, start_after, limit } =>
to_json_binary(&query_list_assets(deps, start_after, limit, Some(owner), None, None, None, None, None)?),
QueryMsg::SearchAssets { query, start_after, limit } => to_json_binary(&query_search_assets(deps, query, start_after, limit)?),
QueryMsg::CheckAccess { asset_id, address } => to_json_binary(&query_check_access(deps, asset_id, address)?),
QueryMsg::GetStats {} => to_json_binary(&query_stats(deps)?),
QueryMsg::GetSubscription { subscription_id } => to_json_binary(&query_subscription(deps, subscription_id)?),
QueryMsg::ListSubscriptions { start_after, limit, subscriber, asset_id } =>
to_json_binary(&query_list_subscriptions(deps, start_after, limit, subscriber, asset_id)?),
}
}
// === 辅助函数 ===
fn assert_admin(storage: &dyn Storage, sender: &Addr) -> Result<(), ContractError> {
let config = CONFIG.load(storage)?;
if sender != &config.admin {
return Err(ContractError::Unauthorized {});
}
Ok(())
}
fn assert_owner(storage: &dyn Storage, asset_id: Uint128, sender: &Addr) -> Result<DataAsset, ContractError> {
let asset = ASSETS.load(storage, asset_id)?;
if sender != &asset.owner {
return Err(ContractError::Unauthorized {});
}
Ok(asset)
}
fn validate_data_hash(hash: &str) -> Result<(), ContractError> {
let hash = hash.strip_prefix("0x").unwrap_or(hash);
if hash.len() != 64 {
return Err(ContractError::InvalidHash {});
}
if !hash.chars().all(|c| c.is_ascii_hexdigit()) {
return Err(ContractError::InvalidHash {});
}
Ok(())
}
fn calculate_fee(price: &Coin, rate_bps: u64) -> Coin {
let fee_amount = price.amount.multiply_ratio(rate_bps as u128 / 10000u128);
Coin { denom: price.denom.clone(), amount: fee_amount }
}
// === 注册数据资产 ===
pub fn execute_register_asset(
deps: DepsMut,
env: Env,
info: MessageInfo,
data_hash: String,
metadata_uri: String,
price: Coin,
license: License,
access_type: AccessType,
data_size: Uint128,
sample_count: Option<Uint128>,
data_type: String,
version: String,
agent_id: Option<String>,
) -> ContractResult<Response> {
validate_data_hash(&data_hash)?;
let config = CONFIG.load(deps.storage)?;
if price.denom != config.min_price.denom {
return Err(ContractError::InvalidDenom {});
}
if price.amount < config.min_price.amount {
return Err(ContractError::PriceTooLow {});
}
if data_size > config.max_data_size {
return Err(ContractError::DataTooLarge {});
}
if !metadata_uri.starts_with("ipfs://") && !metadata_uri.starts_with("https://") && !metadata_uri.starts_with("ar://") {
return Err(ContractError::InvalidMetadataUri {});
}
let new_id = config.asset_count.checked_add(Uint128::one())?;
let asset = DataAsset {
id: new_id,
owner: info.sender.clone(),
data_hash,
metadata_uri,
price: price.clone(),
license,
access_type,
data_size,
sample_count,
data_type,
version,
created_at: env.block.time,
updated_at: env.block.time,
active: true,
total_sales: Uint128::zero(),
total_revenue: Coin { denom: config.min_price.denom.clone(), amount: Uint128::zero() },
quality_score: 0,
verification_status: VerificationStatus::Unverified,
agent_id,
};
ASSETS.save(deps.storage, new_id, &asset)?;
let mut new_config = config;
new_config.asset_count = new_id;
CONFIG.save(deps.storage, &new_config)?;
Ok(Response::new()
.add_attribute("action", "register_asset")
.add_attribute("asset_id", new_id.to_string())
.add_attribute("owner", info.sender.to_string())
.add_attribute("data_type", data_type)
.add_attribute("price", price.to_string()))
}
pub fn execute_update_asset(
deps: DepsMut,
env: Env,
info: MessageInfo,
asset_id: Uint128,
metadata_uri: Option<String>,
price: Option<Coin>,
license: Option<License>,
access_type: Option<AccessType>,
version: Option<String>,
) -> ContractResult<Response> {
let mut asset = assert_owner(deps.storage, asset_id, &info.sender)?;
if !asset.active { return Err(ContractError::AssetNotActive {}); }
if let Some(uri) = metadata_uri {
if !uri.starts_with("ipfs://") && !uri.starts_with("https://") && !uri.starts_with("ar://") {
return Err(ContractError::InvalidMetadataUri {});
}
asset.metadata_uri = uri;
}
if let Some(p) = price {
let config = CONFIG.load(deps.storage)?;
if p.denom != config.min_price.denom { return Err(ContractError::InvalidDenom {}); }
if p.amount < config.min_price.amount { return Err(ContractError::PriceTooLow {}); }
asset.price = p;
}
if let Some(l) = license { asset.license = l; }
if let Some(a) = access_type { asset.access_type = a; }
if let Some(v) = version { asset.version = v; }
asset.updated_at = env.block.time;
ASSETS.save(deps.storage, asset_id, &asset)?;
Ok(Response::new().add_attribute("action", "update_asset").add_attribute("asset_id", asset_id.to_string()))
}
pub fn execute_deactivate_asset(
deps: DepsMut,
env: Env,
info: MessageInfo,
asset_id: Uint128,
) -> ContractResult<Response> {
let mut asset = assert_owner(deps.storage, asset_id, &info.sender)?;
asset.active = false;
asset.updated_at = env.block.time;
ASSETS.save(deps.storage, asset_id, &asset)?;
Ok(Response::new().add_attribute("action", "deactivate_asset").add_attribute("asset_id", asset_id.to_string()))
}
// === 购买数据 ===
pub fn execute_purchase_data(
deps: DepsMut,
env: Env,
info: MessageInfo,
asset_id: Uint128,
) -> ContractResult<Response> {
let asset = ASSETS.load(deps.storage, asset_id)?;
if !asset.active { return Err(ContractError::AssetNotActive {}); }
// 访问控制检查
check_access_control(deps.storage, &asset, &info.sender)?;
// 验证付款
let payment = info.funds.iter()
.find(|c| c.denom == asset.price.denom)
.ok_or(ContractError::InsufficientPayment {})?;
if payment.amount < asset.price.amount {
return Err(ContractError::InsufficientPayment {});
}
let config = CONFIG.load(deps.storage)?;
let fee = calculate_fee(&asset.price, config.fee_rate_bps);
let seller_amount = Coin {
denom: asset.price.denom.clone(),
amount: asset.price.amount.checked_sub(fee.amount)?,
};
let new_trade_id = config.trade_count.checked_add(Uint128::one())?;
let new_escrow_id = config.escrow_count.checked_add(Uint128::one())?;
let trade = TradeRecord {
trade_id: new_trade_id,
asset_id,
buyer: info.sender.clone(),
seller: asset.owner.clone(),
amount: asset.price.clone(),
timestamp: env.block.time,
status: TradeStatus::PaidAwaitingDelivery,
escrow_address: None,
completed_at: None,
};
TRADES.save(deps.storage, new_trade_id, &trade)?;
let escrow = EscrowAccount {
escrow_id: new_escrow_id,
trade_id: new_trade_id,
depositor: info.sender.clone(),
beneficiary: asset.owner.clone(),
amount: asset.price.clone(),
condition: format!("data_hash:{}", asset.data_hash),
status: EscrowStatus::Locked,
created_at: env.block.time,
expires_at: Some(env.block.time.plus_blocks(config.max_arbitration_blocks)),
arbitrator: Some(config.admin.clone()),
};
ESCROWS.save(deps.storage, new_escrow_id, &escrow)?;
let mut updated_asset = asset;
updated_asset.total_sales = updated_asset.total_sales.checked_add(Uint128::one())?;
updated_asset.total_revenue.amount = updated_asset.total_revenue.amount.checked_add(updated_asset.price.amount)?;
ASSETS.save(deps.storage, asset_id, &updated_asset)?;
let mut new_config = config;
new_config.trade_count = new_trade_id;
new_config.escrow_count = new_escrow_id;
CONFIG.save(deps.storage, &new_config)?;
Ok(Response::new()
.add_attribute("action", "purchase_data")
.add_attribute("trade_id", new_trade_id.to_string())
.add_attribute("escrow_id", new_escrow_id.to_string())
.add_attribute("asset_id", asset_id.to_string())
.add_attribute("buyer", info.sender.to_string())
.add_attribute("amount", asset.price.to_string())
.add_attribute("seller", updated_asset.owner.to_string()))
}
fn check_access_control(
storage: &dyn Storage,
asset: &DataAsset,
buyer: &Addr,
) -> Result<(), ContractError> {
match &asset.access_type {
AccessType::Public => Ok(()),
AccessType::TokenGated { token_address: _, min_balance: _ } => Ok(()),
AccessType::Whitelist => {
let authorized = WHITELISTS.load(storage, (asset.id, buyer.clone())).unwrap_or(false);
if !authorized { return Err(ContractError::NotWhitelisted {}); }
Ok(())
}
AccessType::Subscription { price_per_block: _, min_blocks: _ } => Ok(()),
AccessType::NftGated { nft_contract: _, token_id: _ } => Ok(()),
}
}
pub fn execute_confirm_delivery(
deps: DepsMut,
env: Env,
info: MessageInfo,
trade_id: Uint128,
data_hash_verified: String,
) -> ContractResult<Response> {
let mut trade = TRADES.load(deps.storage, trade_id)?;
if trade.buyer != info.sender { return Err(ContractError::Unauthorized {}); }
if trade.status != TradeStatus::PaidAwaitingDelivery && trade.status != TradeStatus::DeliveredAwaitingConfirmation {
return Err(ContractError::InvalidTradeState {});
}
let asset = ASSETS.load(deps.storage, trade.asset_id)?;
let verified = data_hash_verified.trim_start_matches("0x");
let expected = asset.data_hash.trim_start_matches("0x");
if verified != expected {
return Err(ContractError::HashMismatch {});
}
trade.status = TradeStatus::Completed;
trade.completed_at = Some(env.block.time);
TRADES.save(deps.storage, trade_id, &trade)?;
for r in ESCROWS.range(deps.storage, None, None, cosmwasm_std::Order::Ascending) {
if let Ok((id, mut escrow)) = r {
if escrow.trade_id == trade_id {
escrow.status = EscrowStatus::Released;
ESCROWS.save(deps.storage, id, &escrow)?;
break;
}
}
}
Ok(Response::new()
.add_attribute("action", "confirm_delivery")
.add_attribute("trade_id", trade_id.to_string())
.add_attribute("status", "completed"))
}
pub fn execute_dispute_delivery(
deps: DepsMut,
env: Env,
info: MessageInfo,
trade_id: Uint128,
reason: String,
) -> ContractResult<Response> {
let mut trade = TRADES.load(deps.storage, trade_id)?;
if trade.buyer != info.sender { return Err(ContractError::Unauthorized {}); }
if trade.status != TradeStatus::PaidAwaitingDelivery && trade.status != TradeStatus::DeliveredAwaitingConfirmation {
return Err(ContractError::InvalidTradeState {});
}
trade.status = TradeStatus::Disputed;
TRADES.save(deps.storage, trade_id, &trade)?;
for r in ESCROWS.range(deps.storage, None, None, cosmwasm_std::Order::Ascending) {
if let Ok((id, mut escrow)) = r {
if escrow.trade_id == trade_id {
escrow.status = EscrowStatus::Disputed;
ESCROWS.save(deps.storage, id, &escrow)?;
break;
}
}
}
Ok(Response::new()
.add_attribute("action", "dispute_delivery")
.add_attribute("trade_id", trade_id.to_string())
.add_attribute("reason", reason))
}
pub fn execute_request_refund(
deps: DepsMut,
_env: Env,
info: MessageInfo,
trade_id: Uint128,
) -> ContractResult<Response> {
let mut trade = TRADES.load(deps.storage, trade_id)?;
if trade.buyer != info.sender { return Err(ContractError::Unauthorized {}); }
if trade.status != TradeStatus::Disputed { return Err(ContractError::InvalidTradeState {}); }
trade.status = TradeStatus::Refunded;
TRADES.save(deps.storage, trade_id, &trade)?;
for r in ESCROWS.range(deps.storage, None, None, cosmwasm_std::Order::Ascending) {
if let Ok((id, mut escrow)) = r {
if escrow.trade_id == trade_id {
escrow.status = EscrowStatus::Refunded;
ESCROWS.save(deps.storage, id, &escrow)?;
break;
}
}
}
Ok(Response::new().add_attribute("action", "refund").add_attribute("trade_id", trade_id.to_string()))
}
pub fn execute_arbitrate(
deps: DepsMut,
env: Env,
info: MessageInfo,
trade_id: Uint128,
resolve_in_favor_of_buyer: bool,
) -> ContractResult<Response> {
let config = CONFIG.load(deps.storage)?;
if info.sender != config.admin && info.sender != config.fee_collector {
return Err(ContractError::Unauthorized {});
}
let mut trade = TRADES.load(deps.storage, trade_id)?;
if trade.status != TradeStatus::Disputed { return Err(ContractError::InvalidTradeState {}); }
if resolve_in_favor_of_buyer {
trade.status = TradeStatus::Refunded;
} else {
trade.status = TradeStatus::Completed;
trade.completed_at = Some(env.block.time);
}
TRADES.save(deps.storage, trade_id, &trade)?;
for r in ESCROWS.range(deps.storage, None, None, cosmwasm_std::Order::Ascending) {
if let Ok((id, mut escrow)) = r {
if escrow.trade_id == trade_id {
escrow.status = if resolve_in_favor_of_buyer { EscrowStatus::Refunded } else { EscrowStatus::Released };
ESCROWS.save(deps.storage, id, &escrow)?;
break;
}
}
}
Ok(Response::new()
.add_attribute("action", "arbitrate")
.add_attribute("trade_id", trade_id.to_string())
.add_attribute("resolution", if resolve_in_favor_of_buyer { "buyer_wins" } else { "seller_wins" }))
}
// === 白名单管理 ===
pub fn execute_add_to_whitelist(
deps: DepsMut,
_env: Env,
info: MessageInfo,
asset_id: Uint128,
addresses: Vec<Addr>,
) -> ContractResult<Response> {
assert_owner(deps.storage, asset_id, &info.sender)?;
for addr in &addresses { WHITELISTS.save(deps.storage, (asset_id, addr.clone()), &true)?; }
Ok(Response::new()
.add_attribute("action", "add_to_whitelist")
.add_attribute("asset_id", asset_id.to_string())
.add_attribute("count", addresses.len().to_string()))
}
pub fn execute_remove_from_whitelist(
deps: DepsMut,
_env: Env,
info: MessageInfo,
asset_id: Uint128,
addresses: Vec<Addr>,
) -> ContractResult<Response> {
assert_owner(deps.storage, asset_id, &info.sender)?;
for addr in &addresses { WHITELISTS.remove(deps.storage, (asset_id, addr.clone())); }
Ok(Response::new()
.add_attribute("action", "remove_from_whitelist")
.add_attribute("asset_id", asset_id.to_string())
.add_attribute("count", addresses.len().to_string()))
}
// === 访问授权 ===
pub fn execute_grant_access(
deps: DepsMut,
env: Env,
info: MessageInfo,
asset_id: Uint128,
grantee: Addr,
grant_type: AccessGrantType,
) -> ContractResult<Response> {
assert_owner(deps.storage, asset_id, &info.sender)?;
let config = CONFIG.load(deps.storage)?;
let new_grant_id = config.grant_count.checked_add(Uint128::one())?;
let expires_at = match &grant_type {
AccessGrantType::Permanent => None,
AccessGrantType::TimeBound { blocks } => Some(env.block.time.plus_blocks(*blocks)),
AccessGrantType::UsageBased { max_downloads: _, downloads_used: _ } => None,
};
let grant = AccessGrant {
grant_id: new_grant_id,
asset_id,
grantee,
grant_type,
expires_at,
created_at: env.block.time,
};
ACCESS_GRANTS.save(deps.storage, new_grant_id, &grant)?;
let mut new_config = config;
new_config.grant_count = new_grant_id;
CONFIG.save(deps.storage, &new_config)?;
Ok(Response::new()
.add_attribute("action", "grant_access")
.add_attribute("grant_id", new_grant_id.to_string())
.add_attribute("asset_id", asset_id.to_string()))
}
pub fn execute_revoke_access(
deps: DepsMut,
_env: Env,
info: MessageInfo,
grant_id: Uint128,
) -> ContractResult<Response> {
let grant = ACCESS_GRANTS.load(deps.storage, grant_id)?;
let asset = ASSETS.load(deps.storage, grant.asset_id)?;
if info.sender != asset.owner { return Err(ContractError::Unauthorized {}); }
ACCESS_GRANTS.remove(deps.storage, grant_id);
Ok(Response::new().add_attribute("action", "revoke_access").add_attribute("grant_id", grant_id.to_string()))
}
pub fn execute_update_quality_score(
deps: DepsMut,
_env: Env,
info: MessageInfo,
asset_id: Uint128,
score: u8,
) -> ContractResult<Response> {
if score > 100 { return Err(ContractError::InvalidScore {}); }
let config = CONFIG.load(deps.storage)?;
if info.sender != config.admin { return Err(ContractError::Unauthorized {}); }
let mut asset = ASSETS.load(deps.storage, asset_id)?;
asset.quality_score = score;
ASSETS.save(deps.storage, asset_id, &asset)?;
Ok(Response::new()
.add_attribute("action", "update_quality_score")
.add_attribute("asset_id", asset_id.to_string())
.add_attribute("score", score.to_string()))
}
pub fn execute_withdraw_fees(
deps: DepsMut,
_env: Env,
info: MessageInfo,
_denom: String,
_amount: Option<Uint128>,
) -> ContractResult<Response> {
let config = CONFIG.load(deps.storage)?;
if info.sender != config.fee_collector { return Err(ContractError::Unauthorized {}); }
Ok(Response::new().add_attribute("action", "withdraw_fees"))
}
// === 查询函数 ===
pub fn query_asset(deps: Deps, asset_id: Uint128) -> StdResult<AssetResponse> {
let asset = ASSETS.load(deps.storage, asset_id)?;
Ok(AssetResponse { asset })
}
pub fn query_list_assets(
deps: Deps,
start_after: Option<Uint128>,
limit: Option<u32>,
owner: Option<Addr>,
data_type: Option<String>,
min_price: Option<Coin>,
max_price: Option<Coin>,
license: Option<License>,
active_only: Option<bool>,
) -> StdResult<AssetListResponse> {
let limit = limit.unwrap_or(DEFAULT_LIMIT).min(MAX_LIMIT);
let start = start_after.unwrap_or(Uint128::zero());
let active_only = active_only.unwrap_or(true);
let assets: Vec<DataAsset> = ASSETS
.range(deps.storage, None, None, cosmwasm_std::Order::Ascending)
.filter_map(|r| r.ok())
.filter(|(id, _)| id > &start)
.filter(|(_, a)| !active_only || a.active)
.filter(|(_, a)| owner.as_ref().map(|o| a.owner == *o).unwrap_or(true))
.filter(|(_, a)| data_type.as_ref().map(|t| a.data_type.contains(t.as_str())).unwrap_or(true))
.filter(|(_, a)| min_price.as_ref().map(|p| a.price.denom == p.denom && a.price.amount >= p.amount).unwrap_or(true))
.filter(|(_, a)| max_price.as_ref().map(|p| a.price.denom == p.denom && a.price.amount <= p.amount).unwrap_or(true))
.take(limit as usize)
.map(|(_, a)| a)
.collect();
Ok(AssetListResponse { total_count: Uint128::from(assets.len() as u128), assets })
}
pub fn query_trade(deps: Deps, trade_id: Uint128) -> StdResult<TradeResponse> {
let trade = TRADES.load(deps.storage, trade_id)?;
Ok(TradeResponse { trade })
}
pub fn query_list_trades(
deps: Deps,
start_after: Option<Uint128>,
limit: Option<u32>,
buyer: Option<Addr>,
seller: Option<Addr>,
asset_id: Option<Uint128>,
status: Option<TradeStatus>,
) -> StdResult<TradeListResponse> {
let limit = limit.unwrap_or(DEFAULT_LIMIT).min(MAX_LIMIT);
let start = start_after.unwrap_or(Uint128::zero());
let trades: Vec<TradeRecord> = TRADES
.range(deps.storage, None, None, cosmwasm_std::Order::Descending)
.filter_map(|r| r.ok())
.filter(|(id, _)| id > &start)
.filter(|(_, t)| buyer.as_ref().map(|b| t.buyer == *b).unwrap_or(true))
.filter(|(_, t)| seller.as_ref().map(|s| t.seller == *s).unwrap_or(true))
.filter(|(_, t)| asset_id.as_ref().map(|a| t.asset_id == *a).unwrap_or(true))
.filter(|(_, t)| status.as_ref().map(|s| t.status == *s).unwrap_or(true))
.take(limit as usize)
.map(|(_, t)| t)
.collect();
Ok(TradeListResponse { total_count: Uint128::from(trades.len() as u128), trades })
}
pub fn query_escrow(deps: Deps, escrow_id: Uint128) -> StdResult<EscrowResponse> {
let escrow = ESCROWS.load(deps.storage, escrow_id)?;
Ok(EscrowResponse { escrow })
}
pub fn query_search_assets(
deps: Deps,
query: String,
start_after: Option<Uint128>,
limit: Option<u32>,
) -> StdResult<AssetListResponse> {
let limit = limit.unwrap_or(DEFAULT_LIMIT).min(MAX_LIMIT);
let start = start_after.unwrap_or(Uint128::zero());
let query_lower = query.to_lowercase();
let assets: Vec<DataAsset> = ASSETS
.range(deps.storage, None, None, cosmwasm_std::Order::Ascending)
.filter_map(|r| r.ok())
.filter(|(id, _)| id > &start)
.filter(|(_, a)| a.active)
.filter(|(_, a)| {
a.metadata_uri.to_lowercase().contains(&query_lower)
|| a.data_type.to_lowercase().contains(&query_lower)
|| a.owner.to_string().to_lowercase().contains(&query_lower)
})
.take(limit as usize)
.map(|(_, a)| a)
.collect();
Ok(AssetListResponse { total_count: Uint128::from(assets.len() as u128), assets })
}
pub fn query_check_access(deps: Deps, asset_id: Uint128, address: Addr) -> StdResult<crate::msg::AccessResponse> {
let asset = ASSETS.load(deps.storage, asset_id)?;
let has_access = match &asset.access_type {
AccessType::Public => true,
AccessType::Whitelist => WHITELISTS.load(deps.storage, (asset_id, address)).unwrap_or(false),
_ => false,
};
Ok(crate::msg::AccessResponse { has_access, grant: None })
}
pub fn query_stats(deps: Deps) -> StdResult<StatsResponse> {
let config = CONFIG.load(deps.storage)?;
let active_assets = ASSETS
.range(deps.storage, None, None, cosmwasm_std::Order::Ascending)
.filter_map(|r| r.ok())
.filter(|(_, a)| a.active)
.count();
let total_volume = TRADES
.range(deps.storage, None, None, cosmwasm_std::Order::Ascending)
.filter_map(|r| r.ok())
.fold(Uint128::zero(), |acc, (_, t)| acc + t.amount.amount);
Ok(StatsResponse {
total_assets: config.asset_count,
active_assets: Uint128::from(active_assets as u128),
total_trades: config.trade_count,
total_volume: Coin { denom: "uMSG".to_string(), amount: total_volume },
total_fees_collected: Coin { denom: "uMSG".to_string(), amount: Uint128::zero() },
unique_buyers: Uint128::zero(),
unique_sellers: Uint128::zero(),
})
}
pub fn query_subscription(deps: Deps, subscription_id: Uint128) -> StdResult<crate::msg::SubscriptionResponse> {
let subscription = SUBSCRIPTIONS.load(deps.storage, subscription_id)?;
Ok(crate::msg::SubscriptionResponse { subscription })
}
pub fn query_list_subscriptions(
deps: Deps,
start_after: Option<Uint128>,
limit: Option<u32>,
subscriber: Option<Addr>,
asset_id: Option<Uint128>,
) -> StdResult<crate::msg::SubscriptionListResponse> {
let limit = limit.unwrap_or(DEFAULT_LIMIT).min(MAX_LIMIT);
let start = start_after.unwrap_or(Uint128::zero());
let subscriptions: Vec<Subscription> = SUBSCRIPTIONS
.range(deps.storage, None, None, cosmwasm_std::Order::Ascending)
.filter_map(|r| r.ok())
.filter(|(id, _)| id > &start)
.filter(|(_, s)| subscriber.as_ref().map(|sub| s.subscriber == *sub).unwrap_or(true))
.filter(|(_, s)| asset_id.as_ref().map(|a| s.asset_id == *a).unwrap_or(true))
.take(limit as usize)
.map(|(_, s)| s)
.collect();
Ok(crate::msg::SubscriptionListResponse { total_count: Uint128::from(subscriptions.len() as u128), subscriptions })
}
2.2.6 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("Asset not active")]
AssetNotActive {},
#[error("Invalid data hash (must be 64 hex characters)")]
InvalidHash {},
#[error("Price too low")]
PriceTooLow {},
#[error("Invalid denomination")]
InvalidDenom {},
#[error("Data too large")]
DataTooLarge {},
#[error("Invalid metadata URI")]
InvalidMetadataUri {},
#[error("Insufficient payment")]
InsufficientPayment {},
#[error("Invalid trade state")]
InvalidTradeState {},
#[error("Hash mismatch on verification")]
HashMismatch {},
#[error("Address not whitelisted")]
NotWhitelisted {},
#[error("No valid subscription found")]
NoValidSubscription {},
#[error("Not a subscription asset")]
NotSubscriptionAsset {},
#[error("Subscription too short")]
SubscriptionTooShort {},
#[error("Invalid quality score (must be 0-100)")]
InvalidScore {},
#[error("Asset not found")]
AssetNotFound {},
#[error("Trade not found")]
TradeNotFound {},
#[error("Overflow")]
Overflow {},
}
2.2.7 tests/integration_test.rs
#[cfg(test)]
mod tests {
use cosmwasm_std::{coins, from_json, Addr, Coin, Empty, Timestamp, Uint128};
use cw_multi_test::{App, Contract, ContractWrapper, Executor};
use data_marketplace::contract::{execute, instantiate, query};
use data_marketplace::msg::{
AssetListResponse, AssetResponse, ExecuteMsg, InstantiateMsg, QueryMsg, StatsResponse,
TradeListResponse,
};
use data_marketplace::state::{AccessType, DataAsset, License, TradeStatus};
const MSG_DENOM: &str = "uMSG";
fn mock_app() -> App {
App::new(|router, _api, storage| {
router
.bank
.init_balance(storage, &Addr::unchecked("owner"), coins(1_000_000_000, MSG_DENOM))
.unwrap();
router
.bank
.init_balance(storage, &Addr::unchecked("buyer"), coins(1_000_000_000, MSG_DENOM))
.unwrap();
})
}
fn contract_template() -> Box<dyn Contract<Empty>> {
Box::new(ContractWrapper::new(execute, instantiate, query))
}
fn instantiate_contract(app: &mut App, admin: &str) -> Addr {
let code_id = app.store_code(contract_template());
let init_msg = InstantiateMsg {
admin: Addr::unchecked(admin),
name: "Data Marketplace".to_string(),
fee_rate_bps: 100,
fee_collector: Addr::unchecked("fee_collector"),
max_arbitration_blocks: 10080,
min_price: Coin { denom: MSG_DENOM.to_string(), amount: Uint128::from(100_000u128) },
max_data_size: Uint128::from(10_000_000_000u128),
registry_address: None,
aipay_address: None,
};
app.instantiate_contract(code_id, Addr::unchecked(admin), &init_msg, &[], "DataMarketplace", None).unwrap()
}
#[test]
fn test_instantiate() {
let mut app = mock_app();
let addr = instantiate_contract(&mut app, "owner");
let stats: StatsResponse = app.wrap().query_wasm_smart(addr, &QueryMsg::GetStats {}).unwrap();
assert_eq!(stats.total_assets, Uint128::zero());
}
#[test]
fn test_register_asset() {
let mut app = mock_app();
let addr = instantiate_contract(&mut app, "owner");
let msg = ExecuteMsg::RegisterAsset {
data_hash: "abcdef1234567890abcdef1234567890abcdef1234567890abcdef1234567890".to_string(),
metadata_uri: "ipfs://QmTest123".to_string(),
price: Coin { denom: MSG_DENOM.to_string(), amount: Uint128::from(500_000u128) },
license: License::MIT,
access_type: AccessType::Public,
data_size: Uint128::from(1_000_000u128),
sample_count: Some(Uint128::from(10_000u128)),
data_type: "text".to_string(),
version: "1.0.0".to_string(),
agent_id: Some("agent_123".to_string()),
};
app.execute_contract(Addr::unchecked("owner"), addr.clone(), &msg, &[]).unwrap();
let assets: AssetListResponse = app
.wrap()
.query_wasm_smart(
addr,
&QueryMsg::ListAssets {
start_after: None, limit: Some(10), owner: None, data_type: None,
min_price: None, max_price: None, license: None, active_only: Some(true),
},
)
.unwrap();
assert_eq!(assets.assets.len(), 1);
assert_eq!(assets.assets[0].data_type, "text");
}
#[test]
fn test_register_invalid_hash() {
let mut app = mock_app();
let addr = instantiate_contract(&mut app, "owner");
let msg = ExecuteMsg::RegisterAsset {
data_hash: "short".to_string(),
metadata_uri: "ipfs://QmTest".to_string(),
price: Coin { denom: MSG_DENOM.to_string(), amount: Uint128::from(500_000u128) },
license: License::MIT,
access_type: AccessType::Public,
data_size: Uint128::from(1_000_000u128),
sample_count: None,
data_type: "text".to_string(),
version: "1.0.0".to_string(),
agent_id: None,
};
assert!(app.execute_contract(Addr::unchecked("owner"), addr, &msg, &[]).is_err());
}
#[test]
fn test_purchase_and_confirm() {
let mut app = mock_app();
let addr = instantiate_contract(&mut app, "owner");
let hash = "abcdef1234567890abcdef1234567890abcdef1234567890abcdef1234567890";
app.execute_contract(
Addr::unchecked("owner"), addr.clone(), &ExecuteMsg::RegisterAsset {
data_hash: hash.to_string(),
metadata_uri: "ipfs://QmData".to_string(),
price: Coin { denom: MSG_DENOM.to_string(), amount: Uint128::from(500_000u128) },
license: License::MIT,
access_type: AccessType::Public,
data_size: Uint128::from(1_000_000u128),
sample_count: None,
data_type: "image".to_string(),
version: "1.0.0".to_string(),
agent_id: None,
}, &[],
).unwrap();
let result = app.execute_contract(
Addr::unchecked("buyer"), addr.clone(), &ExecuteMsg::PurchaseData { asset_id: Uint128::one() }, &coins(500_000, MSG_DENOM),
);
assert!(result.is_ok());
let trades: TradeListResponse = app.wrap().query_wasm_smart(
addr.clone(), &QueryMsg::ListTrades { start_after: None, limit: Some(10), buyer: None, seller: None, asset_id: None, status: None },
).unwrap();
assert_eq!(trades.trades[0].status, TradeStatus::PaidAwaitingDelivery);
let result = app.execute_contract(
Addr::unchecked("buyer"), addr.clone(), &ExecuteMsg::ConfirmDelivery { trade_id: Uint128::one(), data_hash_verified: hash.to_string() }, &[],
);
assert!(result.is_ok());
let trades: TradeListResponse = app.wrap().query_wasm_smart(
addr, &QueryMsg::ListTrades { start_after: None, limit: Some(10), buyer: None, seller: None, asset_id: None, status: None },
).unwrap();
assert_eq!(trades.trades[0].status, TradeStatus::Completed);
}
#[test]
fn test_whitelist_access() {
let mut app = mock_app();
let addr = instantiate_contract(&mut app, "owner");
let hash = "abcdef1234567890abcdef1234567890abcdef1234567890abcdef1234567890";
app.execute_contract(
Addr::unchecked("owner"), addr.clone(), &ExecuteMsg::RegisterAsset {
data_hash: hash.to_string(),
metadata_uri: "ipfs://QmWhite".to_string(),
price: Coin { denom: MSG_DENOM.to_string(), amount: Uint128::from(500_000u128) },
license: License::MIT,
access_type: AccessType::Whitelist,
data_size: Uint128::from(1_000_000u128),
sample_count: None,
data_type: "text".to_string(),
version: "1.0.0".to_string(),
agent_id: None,
}, &[],
).unwrap();
app.execute_contract(
Addr::unchecked("owner"), addr.clone(), &ExecuteMsg::AddToWhitelist {
asset_id: Uint128::one(), addresses: vec![Addr::unchecked("buyer")],
}, &[],
).unwrap();
let access: data_marketplace::msg::AccessResponse = app.wrap().query_wasm_smart(
addr.clone(), &QueryMsg::CheckAccess { asset_id: Uint128::one(), address: Addr::unchecked("buyer") },
).unwrap();
assert!(access.has_access);
let access: data_marketplace::msg::AccessResponse = app.wrap().query_wasm_smart(
addr, &QueryMsg::CheckAccess { asset_id: Uint128::one(), address: Addr::unchecked("stranger") },
).unwrap();
assert!(!access.has_access);
}
#[test]
fn test_deactivate_asset() {
let mut app = mock_app();
let addr = instantiate_contract(&mut app, "owner");
app.execute_contract(
Addr::unchecked("owner"), addr.clone(), &ExecuteMsg::RegisterAsset {
data_hash: "abcdef1234567890abcdef1234567890abcdef1234567890abcdef1234567890".to_string(),
metadata_uri: "ipfs://QmDeact".to_string(),
price: Coin { denom: MSG_DENOM.to_string(), amount: Uint128::from(500_000u128) },
license: License::MIT,
access_type: AccessType::Public,
data_size: Uint128::from(1_000_000u128),
sample_count: None,
data_type: "text".to_string(),
version: "1.0.0".to_string(),
agent_id: None,
}, &[],
).unwrap();
app.execute_contract(
Addr::unchecked("owner"), addr.clone(), &ExecuteMsg::DeactivateAsset { asset_id: Uint128::one() }, &[],
).unwrap();
let assets: AssetListResponse = app.wrap().query_wasm_smart(
addr, &QueryMsg::ListAssets { start_after: None, limit: Some(10), owner: None, data_type: None, min_price: None, max_price: None, license: None, active_only: Some(true) },
).unwrap();
assert_eq!(assets.assets.len(), 0);
}
#[test]
fn test_dispute_and_arbitrate() {
let mut app = mock_app();
let addr = instantiate_contract(&mut app, "owner");
let hash = "abcdef1234567890abcdef1234567890abcdef1234567890abcdef1234567890";
app.execute_contract(
Addr::unchecked("owner"), addr.clone(), &ExecuteMsg::RegisterAsset {
data_hash: hash.to_string(),
metadata_uri: "ipfs://QmDisp".to_string(),
price: Coin { denom: MSG_DENOM.to_string(), amount: Uint128::from(500_000u128) },
license: License::MIT,
access_type: AccessType::Public,
data_size: Uint128::from(1_000_000u128),
sample_count: None,
data_type: "text".to_string(),
version: "1.0.0".to_string(),
agent_id: None,
}, &[],
).unwrap();
app.execute_contract(
Addr::unchecked("buyer"), addr.clone(), &ExecuteMsg::PurchaseData { asset_id: Uint128::one() }, &coins(500_000, MSG_DENOM),
).unwrap();
app.execute_contract(
Addr::unchecked("buyer"), addr.clone(), &ExecuteMsg::DisputeDelivery {
trade_id: Uint128::one(), reason: "Hash mismatch".to_string(),
}, &[],
).unwrap();
let trades: TradeListResponse = app.wrap().query_wasm_smart(
addr.clone(), &QueryMsg::ListTrades { start_after: None, limit: Some(10), buyer: None, seller: None, asset_id: None, status: None },
).unwrap();
assert_eq!(trades.trades[0].status, TradeStatus::Disputed);
app.execute_contract(
Addr::unchecked("owner"), addr.clone(), &ExecuteMsg::Arbitrate {
trade_id: Uint128::one(), resolve_in_favor_of_buyer: true,
}, &[],
).unwrap();
let trades: TradeListResponse = app.wrap().query_wasm_smart(
addr, &QueryMsg::ListTrades { start_after: None, limit: Some(10), buyer: None, seller: None, asset_id: None, status: None },
).unwrap();
assert_eq!(trades.trades[0].status, TradeStatus::Refunded);
}
}
2.2.8 src/schema.rs
use cosmwasm_schema::write_api;
use data_marketplace::msg::{ExecuteMsg, InstantiateMsg, QueryMsg};
fn main() {
write_api! {
instantiate: InstantiateMsg,
execute: ExecuteMsg,
query: QueryMsg,
}
}
2.3 合约部署
# 使用 workspace-optimizer 编译合约
docker run --rm -v "$(pwd)":/code \
--mount type=volume,source="$(basename "$(pwd)")_cache",target=/code/target \
--mount type=volume,source=registry_cache,target=/usr/local/cargo/registry \
cosmwasm/workspace-optimizer:0.15.0
# 上传合约到 MSG Chain
msgd tx wasm store artifacts/data_marketplace.wasm \
--from my-wallet \
--gas-prices 1000000000attoMSG \
--gas auto \
--gas-adjustment 1.3 \
--chain-id msg-chain-1 \
--node https://rpc.msgchain.zone:443 \
-y
# 实例化合约
INIT='{"admin":"msg1...","name":"AI Agent Data Marketplace","fee_rate_bps":100,"fee_collector":"msg1...","max_arbitration_blocks":10080,"min_price":{"denom":"uMSG","amount":"100000"},"max_data_size":"10000000000"}'
msgd tx wasm instantiate CODE_ID "$INIT" \
--from my-wallet \
--label "data-marketplace-v1" \
--gas-prices 1000000000attoMSG \
--gas auto \
--gas-adjustment 1.3 \
--chain-id msg-chain-1 \
-y
3. 数据交易流程
3.1 完整交易生命周期
买方 数据市场合约 卖方
| | |
| 1. 查询可用数据 | |
|----------------------------->| |
| 2. 返回数据资产列表 | |
|<-----------------------------| |
| | |
| 3. 发起购买(附带资金) | |
|----------------------------->| |
| | 4. 创建托管账户 |
| | 5. 锁定资金 |
| | 6. 通知卖方 |
| |----------------------------->|
| | |
| | 7. 提供数据访问链接+解密密钥 |
| |<-----------------------------|
| | |
| 8. 获取数据并验证哈希 | |
|<---------------------------------------------------------->|
| | |
| 9. 确认收货(哈希匹配) | |
|----------------------------->| 10. 释放资金给卖方 |
| | 11. 结算协议费用 |
| | |
| <- 争议流程(哈希不匹配) | |
|------------------------------> |
3.2 TypeScript SDK 实现
// src/sdk/DataMarketplace.ts
import { SigningCosmWasmClient } from "@cosmjs/cosmwasm-stargate";
import { DirectSecp256k1HdWallet } from "@cosmjs/proto-signing";
import { coins, Coin } from "@cosmjs/amino";
import { calculateFee, GasPrice } from "@cosmjs/stargate";
export interface DataAsset {
id: string;
owner: string;
data_hash: string;
metadata_uri: string;
price: Coin;
license: string;
access_type: AccessType;
data_size: string;
sample_count: string | null;
data_type: string;
version: string;
created_at: string;
updated_at: string;
active: boolean;
total_sales: string;
total_revenue: Coin;
quality_score: number;
verification_status: string;
agent_id: string | null;
}
export interface TradeRecord {
trade_id: string;
asset_id: string;
buyer: string;
seller: string;
amount: Coin;
timestamp: string;
status: TradeStatus;
escrow_address: string | null;
completed_at: string | null;
}
export type AccessType =
| { public: {} }
| { token_gated: { token_address: string; min_balance: string } }
| { whitelist: {} }
| { subscription: { price_per_block: Coin; min_blocks: number } }
| { nft_gated: { nft_contract: string; token_id: string | null } };
export type License =
| "MIT"
| "CCBY40"
| "CCNC40"
| "TrainingOnly"
| "NoResale"
| { custom: { uri: string } };
export type TradeStatus =
| "PendingPayment"
| "PaidAwaitingDelivery"
| "DeliveredAwaitingConfirmation"
| "Completed"
| "Disputed"
| "Refunded"
| "Cancelled";
export interface MarketplaceConfig {
rpcEndpoint: string;
contractAddress: string;
mnemonic: string;
}
export class DataMarketplace {
private client: SigningCosmWasmClient;
private sender: string;
private config: MarketplaceConfig;
private constructor(client: SigningCosmWasmClient, sender: string, config: MarketplaceConfig) {
this.client = client;
this.sender = sender;
this.config = config;
}
static async connect(config: MarketplaceConfig): Promise<DataMarketplace> {
const wallet = await DirectSecp256k1HdWallet.fromMnemonic(config.mnemonic, { prefix: "msg" });
const client = await SigningCosmWasmClient.connectWithSigner(config.rpcEndpoint, wallet);
const [account] = await wallet.getAccounts();
return new DataMarketplace(client, account.address, config);
}
// === 查询方法 ===
async getAsset(assetId: string): Promise<DataAsset> {
const result = await this.client.queryContractSmart(this.config.contractAddress, { get_asset: { asset_id: assetId } });
return result.asset;
}
async listAssets(params?: {
startAfter?: string; limit?: number; owner?: string; dataType?: string;
minPrice?: Coin; maxPrice?: Coin; license?: License; activeOnly?: boolean;
}): Promise<{ assets: DataAsset[]; total_count: string }> {
const q: any = { list_assets: {} };
if (params?.startAfter) q.list_assets.start_after = params.startAfter;
if (params?.limit) q.list_assets.limit = params.limit;
if (params?.owner) q.list_assets.owner = params.owner;
if (params?.dataType) q.list_assets.data_type = params.dataType;
if (params?.minPrice) q.list_assets.min_price = params.minPrice;
if (params?.maxPrice) q.list_assets.max_price = params.maxPrice;
if (params?.license) q.list_assets.license = params.license;
if (params?.activeOnly !== undefined) q.list_assets.active_only = params.activeOnly;
return this.client.queryContractSmart(this.config.contractAddress, q);
}
async searchAssets(query: string, startAfter?: string, limit?: number): Promise<{ assets: DataAsset[]; total_count: string }> {
return this.client.queryContractSmart(this.config.contractAddress, { search_assets: { query, start_after: startAfter, limit } });
}
async getTrade(tradeId: string): Promise<TradeRecord> {
const r = await this.client.queryContractSmart(this.config.contractAddress, { get_trade: { trade_id: tradeId } });
return r.trade;
}
async listTrades(params?: {
startAfter?: string; limit?: number; buyer?: string; seller?: string;
assetId?: string; status?: TradeStatus;
}): Promise<{ trades: TradeRecord[]; total_count: string }> {
const q: any = { list_trades: {} };
if (params?.startAfter) q.list_trades.start_after = params.startAfter;
if (params?.limit) q.list_trades.limit = params.limit;
if (params?.buyer) q.list_trades.buyer = params.buyer;
if (params?.seller) q.list_trades.seller = params.seller;
if (params?.assetId) q.list_trades.asset_id = params.assetId;
if (params?.status) q.list_trades.status = params.status;
return this.client.queryContractSmart(this.config.contractAddress, q);
}
async getStats(): Promise<any> {
return this.client.queryContractSmart(this.config.contractAddress, { get_stats: {} });
}
// === 执行方法 ===
async registerAsset(params: {
dataHash: string; metadataUri: string; price: Coin; license: License;
accessType: AccessType; dataSize: string; sampleCount?: string;
dataType: string; version: string; agentId?: string;
}): Promise<string> {
const msg = {
register_asset: {
data_hash: params.dataHash, metadata_uri: params.metadataUri, price: params.price,
license: params.license, access_type: params.accessType, data_size: params.dataSize,
sample_count: params.sampleCount, data_type: params.dataType, version: params.version,
agent_id: params.agentId,
},
};
const fee = calculateFee(300_000, GasPrice.fromString("1000000000attoMSG"));
const result = await this.client.execute(this.sender, this.config.contractAddress, msg, fee);
return result.transactionHash;
}
async updateAsset(params: {
assetId: string; metadataUri?: string; price?: Coin; license?: License;
accessType?: AccessType; version?: string;
}): Promise<string> {
const msg: any = { update_asset: { asset_id: params.assetId } };
if (params.metadataUri) msg.update_asset.metadata_uri = params.metadataUri;
if (params.price) msg.update_asset.price = params.price;
if (params.license) msg.update_asset.license = params.license;
if (params.accessType) msg.update_asset.access_type = params.accessType;
if (params.version) msg.update_asset.version = params.version;
const fee = calculateFee(200_000, GasPrice.fromString("1000000000attoMSG"));
const result = await this.client.execute(this.sender, this.config.contractAddress, msg, fee);
return result.transactionHash;
}
async purchaseData(assetId: string): Promise<string> {
const asset = await this.getAsset(assetId);
const msg = { purchase_data: { asset_id: assetId } };
const fee = calculateFee(300_000, GasPrice.fromString("1000000000attoMSG"));
const result = await this.client.execute(this.sender, this.config.contractAddress, msg, fee, undefined, [asset.price]);
return result.transactionHash;
}
async confirmDelivery(tradeId: string, verifiedDataHash: string): Promise<string> {
const msg = { confirm_delivery: { trade_id: tradeId, data_hash_verified: verifiedDataHash } };
const fee = calculateFee(200_000, GasPrice.fromString("1000000000attoMSG"));
const result = await this.client.execute(this.sender, this.config.contractAddress, msg, fee);
return result.transactionHash;
}
async disputeDelivery(tradeId: string, reason: string): Promise<string> {
const msg = { dispute_delivery: { trade_id: tradeId, reason } };
const fee = calculateFee(200_000, GasPrice.fromString("1000000000attoMSG"));
const result = await this.client.execute(this.sender, this.config.contractAddress, msg, fee);
return result.transactionHash;
}
// === 工具方法 ===
async computeDataHash(data: Uint8Array): Promise<string> {
const hashBuffer = await crypto.subtle.digest("SHA-256", data);
const hashArray = Array.from(new Uint8Array(hashBuffer));
return hashArray.map((b) => b.toString(16).padStart(2, "0")).join("");
}
async purchaseAndVerify(
assetId: string, fetchDataFn: () => Promise<Uint8Array>
): Promise<{ tradeId: string; verified: boolean }> {
const asset = await this.getAsset(assetId);
const txHash = await this.purchaseData(assetId);
const trades = await this.listTrades({ buyer: this.sender, assetId, status: "PaidAwaitingDelivery" });
if (trades.trades.length === 0) throw new Error("No pending trade found");
const trade = trades.trades[0];
const data = await fetchDataFn();
const computedHash = await this.computeDataHash(data);
const verified = computedHash === asset.data_hash;
if (verified) {
await this.confirmDelivery(trade.trade_id, computedHash);
} else {
await this.disputeDelivery(trade.trade_id, `Data hash mismatch: expected ${asset.data_hash}, got ${computedHash}`);
}
return { tradeId: trade.trade_id, verified };
}
}
3.3 TypeScript 使用示例
// examples/basic_usage.ts
import { DataMarketplace } from "./sdk/DataMarketplace";
async function main() {
const marketplace = await DataMarketplace.connect({
rpcEndpoint: "https://rpc.msgchain.zone:443",
contractAddress: "msg1data_marketplace_contract_address...",
mnemonic: "your wallet mnemonic here...",
});
// === 卖方流程:注册数据资产 ===
const dataContent = new TextEncoder().encode(
JSON.stringify({
name: "中文情感分析数据集",
samples: [
{ text: "这个产品非常好用", sentiment: "positive" },
{ text: "体验太差了", sentiment: "negative" },
{ text: "一般般吧", sentiment: "neutral" },
],
})
);
const dataHash = await marketplace.computeDataHash(dataContent);
const txHash = await marketplace.registerAsset({
dataHash,
metadataUri: "ipfs://QmMetadata123",
price: { denom: "uMSG", amount: "1000000" },
license: "MIT",
accessType: { public: {} },
dataSize: String(dataContent.length),
sampleCount: "3",
dataType: "text_classification",
version: "1.0.0",
agentId: "agent_456",
});
console.log("Asset registered, tx:", txHash);
// === 买方流程:发现并购买数据 ===
const searchResult = await marketplace.searchAssets("情感分析");
if (searchResult.assets.length > 0) {
const targetAsset = searchResult.assets[0];
console.log("Asset:", {
owner: targetAsset.owner,
price: targetAsset.price,
dataType: targetAsset.data_type,
});
const fetchData = async (): Promise<Uint8Array> => {
// 卖方提供的数据下载端点和解密密钥
const response = await fetch(
`https://data-provider.msgchain.org/datasets/${targetAsset.id}/download`,
{ headers: { Authorization: "Bearer <auth_token>" } }
);
return new Uint8Array(await response.arrayBuffer());
};
const result = await marketplace.purchaseAndVerify(targetAsset.id, fetchData);
console.log(result.verified ? "Purchase completed!" : "Dispute raised");
}
}
main().catch(console.error);
3.4 Python SDK 实现
"""
msg_marketplace/sdk.py - MSG Chain Data Marketplace Python SDK
"""
import hashlib
import json
import os
import time
from dataclasses import dataclass
from typing import Any, Dict, List, Optional, Tuple
@dataclass
class DataAsset:
id: str
owner: str
data_hash: str
metadata_uri: str
price: Dict[str, str]
license: Any
access_type: Dict[str, Any]
data_size: str
sample_count: Optional[str]
data_type: str
version: str
created_at: str
updated_at: str
active: bool
total_sales: str
total_revenue: Dict[str, str]
quality_score: int
verification_status: str
agent_id: Optional[str]
@dataclass
class TradeRecord:
trade_id: str
asset_id: str
buyer: str
seller: str
amount: Dict[str, str]
timestamp: str
status: str
escrow_address: Optional[str]
completed_at: Optional[str]
class DataMarketplaceSDK:
"""MSG Chain Data Marketplace Python SDK"""
def __init__(
self,
rpc_endpoint: str,
contract_address: str,
sender_address: str,
):
self.contract_address = contract_address
self.rpc_endpoint = rpc_endpoint
self.sender = sender_address
# 实际使用时集成 cosmos-py 或其他 CosmWasm Python 库
# 此处展示接口定义
def _query(self, msg: Dict[str, Any]) -> Any:
"""Send query to contract (implementation depends on client lib)"""
# Example with a hypothetical CosmWasm client
# return client.query_contract(self.contract_address, msg)
raise NotImplementedError("Use a CosmWasm Python client")
def _execute(self, msg: Dict[str, Any], funds: Optional[List[Dict]] = None) -> str:
"""Execute on contract"""
raise NotImplementedError("Use a CosmWasm Python client")
# --- Queries ---
def get_asset(self, asset_id: str) -> DataAsset:
result = self._query({"get_asset": {"asset_id": asset_id}})
return DataAsset(**result["asset"])
def list_assets(
self,
start_after: Optional[str] = None,
limit: Optional[int] = None,
owner: Optional[str] = None,
data_type: Optional[str] = None,
active_only: bool = True,
) -> Tuple[List[DataAsset], int]:
msg = {"list_assets": {"active_only": active_only}}
if start_after: msg["list_assets"]["start_after"] = start_after
if limit: msg["list_assets"]["limit"] = limit
if owner: msg["list_assets"]["owner"] = owner
if data_type: msg["list_assets"]["data_type"] = data_type
result = self._query(msg)
return [DataAsset(**a) for a in result["assets"]], int(result["total_count"])
def search_assets(self, query: str, limit: Optional[int] = None) -> Tuple[List[DataAsset], int]:
msg = {"search_assets": {"query": query}}
if limit: msg["search_assets"]["limit"] = limit
result = self._query(msg)
return [DataAsset(**a) for a in result["assets"]], int(result["total_count"])
def get_trade(self, trade_id: str) -> TradeRecord:
result = self._query({"get_trade": {"trade_id": trade_id}})
return TradeRecord(**result["trade"])
def list_trades(
self,
buyer: Optional[str] = None,
seller: Optional[str] = None,
status: Optional[str] = None,
limit: Optional[int] = None,
) -> Tuple[List[TradeRecord], int]:
msg = {"list_trades": {}}
if buyer: msg["list_trades"]["buyer"] = buyer
if seller: msg["list_trades"]["seller"] = seller
if status: msg["list_trades"]["status"] = status
if limit: msg["list_trades"]["limit"] = limit
result = self._query(msg)
return [TradeRecord(**t) for t in result["trades"]], int(result["total_count"])
# --- Executions ---
def register_asset(
self,
data_hash: str,
metadata_uri: str,
price: Dict[str, str],
data_size: str,
data_type: str,
version: str,
license: str = "MIT",
access_type: Optional[Dict] = None,
sample_count: Optional[str] = None,
agent_id: Optional[str] = None,
) -> str:
msg = {
"register_asset": {
"data_hash": data_hash,
"metadata_uri": metadata_uri,
"price": price,
"license": license,
"access_type": access_type or {"public": {}},
"data_size": data_size,
"data_type": data_type,
"version": version,
}
}
if sample_count: msg["register_asset"]["sample_count"] = sample_count
if agent_id: msg["register_asset"]["agent_id"] = agent_id
return self._execute(msg)
def purchase_data(self, asset_id: str) -> str:
asset = self.get_asset(asset_id)
msg = {"purchase_data": {"asset_id": asset_id}}
return self._execute(msg, funds=[asset.price])
def confirm_delivery(self, trade_id: str, verified_hash: str) -> str:
return self._execute({"confirm_delivery": {"trade_id": trade_id, "data_hash_verified": verified_hash}})
def dispute_delivery(self, trade_id: str, reason: str) -> str:
return self._execute({"dispute_delivery": {"trade_id": trade_id, "reason": reason}})
# --- Utilities ---
@staticmethod
def compute_sha256(data: bytes) -> str:
return hashlib.sha256(data).hexdigest()
@staticmethod
def encrypt_data(data: bytes, key: bytes) -> bytes:
try:
from cryptography.fernet import Fernet
import base64
f = Fernet(base64.urlsafe_b64encode(key[:32]))
return f.encrypt(data)
except ImportError:
raise ImportError("cryptography library required")
@staticmethod
def decrypt_data(encrypted_data: bytes, key: bytes) -> bytes:
try:
from cryptography.fernet import Fernet
import base64
f = Fernet(base64.urlsafe_b64encode(key[:32]))
return f.decrypt(encrypted_data)
except ImportError:
raise ImportError("cryptography library required")
def verify_data_integrity(self, data: bytes, expected_hash: str) -> bool:
return self.compute_sha256(data) == expected_hash
3.5 Python SDK 使用示例
"""
examples/python_usage.py
"""
import json
import os
from msg_marketplace.sdk import DataMarketplaceSDK
def seller_workflow(sdk: DataMarketplaceSDK):
"""卖方:注册数据资产"""
dataset = {
"name": "中文意图分类数据集",
"version": "2.0.0",
"samples": [
{"text": "帮我查一下余额", "intent": "check_balance"},
{"text": "转账给张三", "intent": "transfer"},
{"text": "我最近的交易记录", "intent": "transaction_history"},
],
}
data_bytes = json.dumps(dataset, ensure_ascii=False).encode("utf-8")
data_hash = sdk.compute_sha256(data_bytes)
print(f"Data hash: {data_hash}")
encryption_key = os.urandom(32)
encrypted_data = sdk.encrypt_data(data_bytes, encryption_key)
# upload encrypted_data to IPFS -> QmEncrypted...
metadata_uri = "ipfs://QmMetadataURI"
tx_hash = sdk.register_asset(
data_hash=data_hash,
metadata_uri=metadata_uri,
price={"denom": "uMSG", "amount": "2000000"},
data_size=str(len(data_bytes)),
data_type="intent_classification",
version="2.0.0",
agent_id="agent_intent_v3",
)
print(f"Asset registered, tx: {tx_hash}")
def buyer_workflow(sdk: DataMarketplaceSDK):
"""买方:搜索、购买、验证"""
assets, count = sdk.search_assets("意图分类")
print(f"Found {count} datasets")
if count == 0: return
target = assets[0]
print(f"Buying {target.id}, price: {target.price}")
tx_hash = sdk.purchase_data(target.id)
print(f"Purchase tx: {tx_hash}")
# Simulate fetching and verifying data
mock_data = json.dumps({"test": "data"}).encode()
actual_hash = sdk.compute_sha256(mock_data)
if sdk.verify_data_integrity(mock_data, target.data_hash):
sdk.confirm_delivery("1", actual_hash)
print("Data verified and confirmed!")
else:
sdk.dispute_delivery("1", "Hash mismatch")
print("Dispute raised")
if __name__ == "__main__":
sdk = DataMarketplaceSDK(
rpc_endpoint="https://rpc.msgchain.zone:443",
contract_address="msg1contract...",
sender_address="msg1sender...",
)
seller_workflow(sdk)
buyer_workflow(sdk)
4. 访问控制
4.1 访问控制模式总览
| 模式 | 适用场景 | 实现复杂度 | 费用模型 |
|---|---|---|---|
| Public | 开放数据、训练数据 | 低 | 按次付费 |
| TokenGated | 社区专属数据 | 中 | 持有代币可购买 |
| Whitelist | 私有数据、合作伙伴数据 | 低 | 按次付费 |
| Subscription | 数据流、持续更新数据 | 高 | 按区块计费 |
| NftGated | 会员数据、通行证模式 | 中 | 持 NFT 可购买 |
4.2 Token-Gated 访问
// src/access/TokenGated.ts
import { SigningCosmWasmClient } from "@cosmjs/cosmwasm-stargate";
import { Coin } from "@cosmjs/amino";
import { calculateFee, GasPrice } from "@cosmjs/stargate";
interface TokenGateConfig {
tokenAddress: string;
minBalance: string;
}
export class TokenGatedAccess {
private client: SigningCosmWasmClient;
private sender: string;
private marketplaceAddress: string;
constructor(client: SigningCosmWasmClient, sender: string, marketplaceAddress: string) {
this.client = client;
this.sender = sender;
this.marketplaceAddress = marketplaceAddress;
}
async registerTokenGatedAsset(params: {
dataHash: string; metadataUri: string; price: Coin; dataType: string; version: string;
tokenAddress: string; minBalance: string;
}): Promise<string> {
const msg = {
register_asset: {
data_hash: params.dataHash, metadata_uri: params.metadataUri, price: params.price,
license: "MIT", data_size: "0", data_type: params.dataType, version: params.version,
access_type: { token_gated: { token_address: params.tokenAddress, min_balance: params.minBalance } },
},
};
const fee = calculateFee(300_000, GasPrice.fromString("1000000000attoMSG"));
const result = await this.client.execute(this.sender, this.marketplaceAddress, msg, fee);
return result.transactionHash;
}
async checkEligibility(assetId: string, buyerAddress: string): Promise<{
eligible: boolean; currentBalance: string; requiredBalance: string; tokenAddress: string;
}> {
const asset = await this.client.queryContractSmart(this.marketplaceAddress, { get_asset: { asset_id: assetId } });
const access = asset.asset.access_type;
if (!("token_gated" in access)) throw new Error("Not token-gated" as string);
const { token_address, min_balance } = access.token_gated;
const balanceResp = await this.client.queryContractSmart(token_address, { balance: { address: buyerAddress } });
const currentBalance = balanceResp.balance;
return {
eligible: BigInt(currentBalance) >= BigInt(min_balance),
currentBalance, requiredBalance: min_balance, tokenAddress: token_address,
};
}
}
4.3 Whitelist 访问
// src/access/WhitelistAccess.ts
export class WhitelistAccess {
private client: SigningCosmWasmClient;
private sender: string;
private marketplaceAddress: string;
constructor(client: SigningCosmWasmClient, sender: string, marketplaceAddress: string) {
this.client = client;
this.sender = sender;
this.marketplaceAddress = marketplaceAddress;
}
async addToWhitelist(assetId: string, addresses: string[]): Promise<string[]> {
const txHashes: string[] = [];
const BATCH_SIZE = 50;
for (let i = 0; i < addresses.length; i += BATCH_SIZE) {
const batch = addresses.slice(i, i + BATCH_SIZE);
const msg = { add_to_whitelist: { asset_id: assetId, addresses: batch } };
const fee = calculateFee(250_000, GasPrice.fromString("1000000000attoMSG"));
const result = await this.client.execute(this.sender, this.marketplaceAddress, msg, fee);
txHashes.push(result.transactionHash);
}
return txHashes;
}
async removeFromWhitelist(assetId: string, addresses: string[]): Promise<string> {
const msg = { remove_from_whitelist: { asset_id: assetId, addresses } };
const fee = calculateFee(200_000, GasPrice.fromString("1000000000attoMSG"));
const result = await this.client.execute(this.sender, this.marketplaceAddress, msg, fee);
return result.transactionHash;
}
async isWhitelisted(assetId: string, address: string): Promise<boolean> {
const result = await this.client.queryContractSmart(this.marketplaceAddress, {
check_access: { asset_id: assetId, address },
});
return result.has_access;
}
}
4.4 订阅访问
// src/access/Subscription.rs - 链上订阅逻辑
use cosmwasm_std::{Addr, Coin, DepsMut, Env, MessageInfo, Response, StdResult, Storage, Uint128};
use cw_storage_plus::Map;
use crate::error::ContractError;
use crate::state::{AccessType, DataAsset, Subscription, ASSETS, CONFIG, SUBSCRIPTIONS};
use crate::ContractResult;
/// 订阅状态存储
pub const SUBSCRIPTION_STORAGE: Map<Uint128, Subscription> = Map::new("subscriptions");
/// 创建订阅
pub fn execute_create_subscription(
deps: DepsMut,
env: Env,
info: MessageInfo,
asset_id: Uint128,
duration_blocks: u64,
) -> ContractResult<Response> {
let asset = ASSETS.load(deps.storage, asset_id)?;
if !asset.active { return Err(ContractError::AssetNotActive {}); }
let price_per_block = match &asset.access_type {
AccessType::Subscription { price_per_block, min_blocks } => {
if duration_blocks < *min_blocks {
return Err(ContractError::SubscriptionTooShort {});
}
price_per_block.clone()
}
_ => return Err(ContractError::NotSubscriptionAsset {}),
};
let total_amount = price_per_block.amount.checked_multiply_ratio(duration_blocks, 1u128)?;
let total_price = Coin { denom: price_per_block.denom.clone(), amount: total_amount };
let payment = info.funds.iter()
.find(|c| c.denom == total_price.denom)
.ok_or(ContractError::InsufficientPayment {})?;
if payment.amount < total_price.amount {
return Err(ContractError::InsufficientPayment {});
}
let config = CONFIG.load(deps.storage)?;
let new_id = config.subscription_count.checked_add(Uint128::one())?;
let expires_at = env.block.time.plus_blocks(duration_blocks);
let subscription = Subscription {
id: new_id,
asset_id,
subscriber: info.sender.clone(),
started_at: env.block.time,
expires_at,
price_per_block: price_per_block.clone(),
total_blocks: duration_blocks,
auto_renew: false,
};
SUBSCRIPTIONS.save(deps.storage, new_id, &subscription)?;
let mut new_config = config;
new_config.subscription_count = new_id;
CONFIG.save(deps.storage, &new_config)?;
Ok(Response::new()
.add_attribute("action", "create_subscription")
.add_attribute("subscription_id", new_id.to_string())
.add_attribute("asset_id", asset_id.to_string())
.add_attribute("expires_at", expires_at.to_string()))
}
/// 检查订阅是否有效
pub fn check_subscription_valid(
storage: &dyn Storage,
asset_id: Uint128,
subscriber: &Addr,
current_time: Timestamp,
) -> bool {
SUBSCRIPTIONS
.range(storage, None, None, cosmwasm_std::Order::Ascending)
.filter_map(|r| r.ok())
.any(|(_, sub)| {
sub.asset_id == asset_id
&& sub.subscriber == *subscriber
&& sub.expires_at.seconds() > current_time.seconds()
})
}
4.5 NFT-Based 数据访问通行证
// src/access/NftGated.rs - NFT 门控访问
use cosmwasm_std::{
to_json_binary, Addr, Deps, DepsMut, Env, MessageInfo, QuerierWrapper, QueryRequest,
Response, Uint128, WasmQuery,
};
use serde::{Deserialize, Serialize};
/// CW721 查询消息(用于检查 NFT 所有权)
#[derive(Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
enum Cw721QueryMsg {
OwnerOf { token_id: String, include_expired: Option<bool> },
Tokens { owner: String, start_after: Option<String>, limit: Option<u32> },
}
#[derive(Serialize, Deserialize)]
struct OwnerOfResponse { owner: Addr }
#[derive(Serialize, Deserialize)]
struct TokensResponse { tokens: Vec<String> }
/// 检查买方是否持有指定 NFT
pub fn check_nft_ownership(
querier: &QuerierWrapper,
nft_contract: &Addr,
owner: &Addr,
token_id: Option<&String>,
) -> bool {
if let Some(tid) = token_id {
let query = WasmQuery::Smart {
contract_addr: nft_contract.into(),
msg: to_json_binary(&Cw721QueryMsg::OwnerOf {
token_id: tid.clone(),
include_expired: None,
}).unwrap(),
};
let resp: Result<OwnerOfResponse, _> = querier.query(&QueryRequest::Wasm(query));
match resp {
Ok(r) => r.owner == *owner,
Err(_) => false,
}
} else {
let query = WasmQuery::Smart {
contract_addr: nft_contract.into(),
msg: to_json_binary(&Cw721QueryMsg::Tokens {
owner: owner.into(),
start_after: None,
limit: Some(1),
}).unwrap(),
};
let resp: Result<TokensResponse, _> = querier.query(&QueryRequest::Wasm(query));
match resp {
Ok(r) => !r.tokens.is_empty(),
Err(_) => false,
}
}
}
// src/access/NftGatedAccess.ts
export class NftGatedAccess {
private client: SigningCosmWasmClient;
private marketplaceAddress: string;
constructor(client: SigningCosmWasmClient, marketplaceAddress: string) {
this.client = client;
this.marketplaceAddress = marketplaceAddress;
}
async checkAccess(assetId: string, userAddress: string): Promise<{
hasAccess: boolean; nftContract: string; requiredTokenId: string | null; ownedTokens: string[];
}> {
const asset = await this.client.queryContractSmart(this.marketplaceAddress, { get_asset: { asset_id: assetId } });
const access = asset.asset.access_type;
if (!("nft_gated" in access)) throw new Error("Not NFT-gated");
const { nft_contract, token_id } = access.nft_gated;
const tokens = await this.queryAllTokens(nft_contract, userAddress);
const hasAccess = token_id ? tokens.includes(token_id) : tokens.length > 0;
return { hasAccess, nftContract: nft_contract, requiredTokenId: token_id || null, ownedTokens: tokens };
}
private async queryAllTokens(nftContract: string, owner: string): Promise<string[]> {
const tokens: string[] = [];
let startAfter: string | undefined;
while (true) {
const result = await this.client.queryContractSmart(nftContract, {
tokens: { owner, start_after: startAfter, limit: 100 },
});
tokens.push(...result.tokens);
if (result.tokens.length < 100) break;
startAfter = result.tokens[result.tokens.length - 1];
}
return tokens;
}
}
4.6 访问授权管理器
// src/access/AccessManager.ts
import { AccessGrantType } from "../sdk/DataMarketplace";
export class AccessManager {
private client: SigningCosmWasmClient;
private sender: string;
private marketplaceAddress: string;
constructor(client: SigningCosmWasmClient, sender: string, marketplaceAddress: string) {
this.client = client;
this.sender = sender;
this.marketplaceAddress = marketplaceAddress;
}
async grantAccess(
assetId: string, grantee: string,
grantType: "permanent" | { time_bound: { blocks: number } } | { usage_based: { max_downloads: number } }
): Promise<string> {
const msg = { grant_access: { asset_id: assetId, grantee, grant_type: grantType } };
const fee = calculateFee(250_000, GasPrice.fromString("1000000000attoMSG"));
const result = await this.client.execute(this.sender, this.marketplaceAddress, msg, fee);
return result.transactionHash;
}
async revokeAccess(grantId: string): Promise<string> {
const msg = { revoke_access: { grant_id: grantId } };
const fee = calculateFee(200_000, GasPrice.fromString("1000000000attoMSG"));
const result = await this.client.execute(this.sender, this.marketplaceAddress, msg, fee);
return result.transactionHash;
}
async getGrantsForUser(userAddress: string): Promise<AccessGrantInfo[]> {
const result = await this.client.queryContractSmart(this.marketplaceAddress, {
list_user_grants: { user: userAddress },
});
return result.grants;
}
}
5. 数据验证与仲裁
5.1 哈希承诺-揭示机制
数据验证使用链下承诺+链上验证模式:
- 注册时: 数据提供方计算 SHA256(data),将哈希值写入链上
- 交付后: 买方获取数据,重新计算哈希,与链上承诺比对
- 匹配: 确认交付,资金释放给卖方
- 不匹配: 发起争议,进入仲裁流程
注册: data --SHA256--> 0xABC... --TX--> 链上 DataAsset.data_hash
交付: data' --SHA256--> 0xABC... == 0xABC... --> 确认 ✅
欺诈: data' --SHA256--> 0xDEF... != 0xABC... --> 争议 ❌
5.2 数据验证工具 (Python)
"""
msg_marketplace/verification.py
"""
import hashlib
import hmac
import os
import time
from typing import Dict, List, Optional, Tuple
class DataVerifier:
"""数据完整性验证器"""
@staticmethod
def compute_merkle_root(data_chunks: List[bytes]) -> str:
"""计算 Merkle 树根哈希(大数据集分块验证)"""
if not data_chunks:
return hashlib.sha256(b"").hexdigest()
current_level = [hashlib.sha256(chunk).hexdigest() for chunk in data_chunks]
while len(current_level) > 1:
next_level = []
for i in range(0, len(current_level), 2):
combined = current_level[i] + (current_level[i + 1] if i + 1 < len(current_level) else current_level[i])
next_level.append(hashlib.sha256(combined.encode()).hexdigest())
current_level = next_level
return current_level[0]
@staticmethod
def generate_commitment(data: bytes, salt: Optional[bytes] = None) -> Tuple[str, bytes]:
"""生成数据承诺(含盐值防碰撞)"""
salt = salt or os.urandom(32)
commit = hashlib.sha256(salt + data).hexdigest()
return commit, salt
@staticmethod
def verify_commitment(data: bytes, salt: bytes, expected: str) -> bool:
"""验证数据承诺"""
return hashlib.sha256(salt + data).hexdigest() == expected
@staticmethod
def compute_partial_hash(data_range: bytes, metadata: Dict[str, str]) -> str:
"""计算数据段哈希(用于流式验证)"""
h = hashlib.sha256()
h.update(json.dumps(metadata, sort_keys=True).encode())
h.update(data_range)
return h.hexdigest()
@staticmethod
def verify_streaming(
chunks: List[bytes],
expected_root: str,
chunk_hashes: List[str],
) -> bool:
"""流式数据验证 - 验证每个分块并检查根哈希"""
computed_hashes = [hashlib.sha256(c).hexdigest() for c in chunks]
if computed_hashes != chunk_hashes:
return False
return DataVerifier.compute_merkle_root(chunks) == expected_root
class DataChallenge:
"""
数据挑战协议 - 零知识风格的数据所有权证明
买方随机挑战数据中的特定块,卖方必须证明拥有该块
"""
def __init__(self, data_hash: str):
self.data_hash = data_hash
def generate_challenge(self) -> Tuple[str, Dict]:
challenge_id = hashlib.sha256(os.urandom(32)).hexdigest()[:16]
challenge = {
"challenge_id": challenge_id,
"type": "random_block_hash",
"block_index": int.from_bytes(os.urandom(4), "big") % 1000,
"block_size": 1024,
"timestamp": int(time.time()),
}
return challenge_id, challenge
def verify_response(self, challenge: Dict, response_block: bytes, expected_hash: str) -> bool:
return hashlib.sha256(response_block).hexdigest() == expected_hash
5.3 争议解决
// src/dispute/DisputeResolver.ts
interface DisputeInfo {
tradeId: string;
buyer: string;
seller: string;
amount: Coin;
reason: string;
status: "open" | "under_review" | "resolved";
arbitrator: string | null;
createdAt: string;
resolvedAt: string | null;
}
export class DisputeResolver {
private marketplaceAddress: string;
constructor(marketplaceAddress: string) {
this.marketplaceAddress = marketplaceAddress;
}
async raiseDispute(
tradeId: string, reason: string, client: SigningCosmWasmClient, sender: string
): Promise<string> {
const msg = { dispute_delivery: { trade_id: tradeId, reason } };
const fee = calculateFee(300_000, GasPrice.fromString("1000000000attoMSG"));
const result = await client.execute(sender, this.marketplaceAddress, msg, fee);
return result.transactionHash;
}
async submitEvidence(
disputeId: string, evidenceUri: string, client: SigningCosmWasmClient, sender: string
): Promise<string> {
const msg = { submit_evidence: { dispute_id: disputeId, evidence_uri: evidenceUri } };
const fee = calculateFee(250_000, GasPrice.fromString("1000000000attoMSG"));
const result = await client.execute(sender, this.marketplaceAddress, msg, fee);
return result.transactionHash;
}
async arbitrate(
tradeId: string, resolveInFavorOfBuyer: boolean,
client: SigningCosmWasmClient, arbitrator: string
): Promise<string> {
const msg = { arbitrate: { trade_id: tradeId, resolve_in_favor_of_buyer: resolveInFavorOfBuyer } };
const fee = calculateFee(250_000, GasPrice.fromString("1000000000attoMSG"));
const result = await client.execute(arbitrator, this.marketplaceAddress, msg, fee);
return result.transactionHash;
}
async getDisputeDetails(tradeId: string, client: SigningCosmWasmClient): Promise<DisputeInfo> {
const trade = await client.queryContractSmart(this.marketplaceAddress, { get_trade: { trade_id: tradeId } });
return {
tradeId: trade.trade.trade_id,
buyer: trade.trade.buyer,
seller: trade.trade.seller,
amount: trade.trade.amount,
reason: "",
status: trade.trade.status === "Disputed" ? "open" : "resolved",
arbitrator: null,
createdAt: trade.trade.timestamp,
resolvedAt: trade.trade.completed_at,
};
}
}
5.4 链上验证模块 (Rust)
// src/verification.rs - 链上验证逻辑
use cosmwasm_std::{Addr, DepsMut, Env, MessageInfo, Response, StdResult, Uint128};
use sha2::{Digest, Sha256};
use crate::error::ContractError;
use crate::state::{VerificationStatus, ASSETS, CONFIG, ESCROWS, EscrowStatus};
use crate::ContractResult;
/// Oracle 验证结果
pub fn execute_oracle_verify(
deps: DepsMut,
_env: Env,
info: MessageInfo,
asset_id: Uint128,
is_valid: bool,
) -> ContractResult<Response> {
let config = CONFIG.load(deps.storage)?;
// 在实际实现中,应检查 info.sender 是否为可信 Oracle
// 此处简化处理
let mut asset = ASSETS.load(deps.storage, asset_id)?;
if is_valid {
asset.verification_status = VerificationStatus::Verified;
} else {
asset.verification_status = VerificationStatus::Failed;
}
ASSETS.save(deps.storage, asset_id, &asset)?;
Ok(Response::new()
.add_attribute("action", "oracle_verify")
.add_attribute("asset_id", asset_id.to_string())
.add_attribute("result", if is_valid { "verified" } else { "failed" }))
}
/// 批量验证 - 提交多个数据哈希并标记已验证
pub fn execute_batch_verify(
deps: DepsMut,
_env: Env,
info: MessageInfo,
asset_ids: Vec<Uint128>,
) -> ContractResult<Response> {
let config = CONFIG.load(deps.storage)?;
if info.sender != config.admin {
return Err(ContractError::Unauthorized {});
}
let mut count = 0u64;
for asset_id in &asset_ids {
if let Ok(mut asset) = ASSETS.load(deps.storage, *asset_id) {
asset.verification_status = VerificationStatus::Verified;
ASSETS.save(deps.storage, *asset_id, &asset)?;
count += 1;
}
}
Ok(Response::new()
.add_attribute("action", "batch_verify")
.add_attribute("count", count.to_string()))
}
5.5 信誉系统
// src/reputation.rs - 信誉评分
use cosmwasm_std::{Addr, Deps, DepsMut, Env, MessageInfo, Response, StdResult, Uint128, Storage};
use cw_storage_plus::Map;
use crate::error::ContractError;
use crate::state::{TradeRecord, TradeStatus, TRADES};
use crate::ContractResult;
/// 信誉分数存储
pub const REPUTATION: Map<Addr, ReputationScore> = Map::new("reputation");
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct ReputationScore {
pub address: Addr,
pub total_score: u64,
pub trade_count: u64,
pub successful_trades: u64,
pub dispute_count: u64,
pub disputes_won: u64,
pub average_rating: f64,
pub total_ratings: u64,
pub rating_sum: u64,
pub verified_provider: bool,
pub last_updated: u64,
}
/// 更新交易方信誉
pub fn update_trade_reputation(
storage: &mut dyn Storage,
trade: &TradeRecord,
buyer_satisfied: bool,
) -> StdResult<()> {
// 更新卖方信誉
let mut seller_rep = REPUTATION
.load(storage, trade.seller.clone())
.unwrap_or(ReputationScore {
address: trade.seller.clone(),
total_score: 0,
trade_count: 0,
successful_trades: 0,
dispute_count: 0,
disputes_won: 0,
average_rating: 0.0,
total_ratings: 0,
rating_sum: 0,
verified_provider: false,
last_updated: 0,
});
seller_rep.trade_count += 1;
if buyer_satisfied {
seller_rep.successful_trades += 1;
seller_rep.total_score += 10;
} else {
seller_rep.dispute_count += 1;
seller_rep.total_score = seller_rep.total_score.saturating_sub(5);
}
if seller_rep.trade_count >= 10 && seller_rep.successful_trades as f64 / seller_rep.trade_count as f64 >= 0.9 {
seller_rep.verified_provider = true;
}
REPUTATION.save(storage, trade.seller.clone(), &seller_rep)?;
// 更新买方信誉
let mut buyer_rep = REPUTATION
.load(storage, trade.buyer.clone())
.unwrap_or(ReputationScore {
address: trade.buyer.clone(),
total_score: 0,
trade_count: 0,
successful_trades: 0,
dispute_count: 0,
disputes_won: 0,
average_rating: 0.0,
total_ratings: 0,
rating_sum: 0,
verified_provider: false,
last_updated: 0,
});
buyer_rep.trade_count += 1;
REPUTATION.save(storage, trade.buyer.clone(), &buyer_rep)?;
Ok(())
}
6. 前端实现
6.1 技术栈
- React 18 + TypeScript
- @cosmjs/cosmwasm-stargate 与 MSG Chain 交互
- React Router 6 路由管理
- TailwindCSS 样式
- React Query 数据获取与缓存
6.2 数据资产浏览页
// src/pages/Marketplace.tsx
import { useState, useEffect } from "react";
import { useQuery } from "@tanstack/react-query";
import { DataMarketplace, DataAsset } from "../sdk/DataMarketplace";
import { useWallet } from "../hooks/useWallet";
import { AssetCard } from "../components/AssetCard";
import { SearchBar } from "../components/SearchBar";
import { FilterPanel } from "../components/FilterPanel";
export function Marketplace() {
const { client, address } = useWallet();
const [searchQuery, setSearchQuery] = useState("");
const [dataType, setDataType] = useState<string | undefined>();
const [minPrice, setMinPrice] = useState<string | undefined>();
const [maxPrice, setMaxPrice] = useState<string | undefined>();
const marketplace = new DataMarketplace(client, address, {
rpcEndpoint: "", contractAddress: process.env.REACT_APP_MARKETPLACE_ADDRESS!, mnemonic: "",
});
const { data, isLoading } = useQuery({
queryKey: ["assets", searchQuery, dataType, minPrice, maxPrice],
queryFn: () => marketplace.listAssets({
dataType, activeOnly: true,
minPrice: minPrice ? { denom: "uMSG", amount: minPrice } : undefined,
maxPrice: maxPrice ? { denom: "uMSG", amount: maxPrice } : undefined,
}),
});
return (
<div className="max-w-7xl mx-auto px-4 py-8">
<h1 className="text-3xl font-bold mb-8">数据市场</h1>
<div className="flex gap-4 mb-8">
<SearchBar value={searchQuery} onChange={setSearchQuery} />
<FilterPanel
dataType={dataType}
onDataTypeChange={setDataType}
minPrice={minPrice}
maxPrice={maxPrice}
onPriceChange={(min, max) => { setMinPrice(min); setMaxPrice(max); }}
/>
</div>
{isLoading ? (
<div className="grid grid-cols-1 md:grid-cols-3 gap-6">
{[1,2,3,4,5,6].map(i => <div key={i} className="h-48 bg-gray-100 rounded-lg animate-pulse" />)}
</div>
) : (
<div className="grid grid-cols-1 md:grid-cols-3 gap-6">
{data?.assets.map(asset => (
<AssetCard key={asset.id} asset={asset} />
))}
</div>
)}
</div>
);
}
6.3 数据资产卡片组件
// src/components/AssetCard.tsx
import { DataAsset } from "../sdk/DataMarketplace";
import { Card, Badge, Button } from "./ui";
import { formatPrice } from "../utils/format";
interface Props {
asset: DataAsset;
onPurchase?: (assetId: string) => void;
}
export function AssetCard({ asset, onPurchase }: Props) {
const formatDataSize = (bytes: string) => {
const n = Number(bytes);
if (n > 1e9) return `${(n / 1e9).toFixed(2)} GB`;
if (n > 1e6) return `${(n / 1e6).toFixed(2)} MB`;
if (n > 1e3) return `${(n / 1e3).toFixed(2)} KB`;
return `${n} B`;
};
return (
<Card className="hover:shadow-lg transition-shadow">
<div className="p-6">
<div className="flex justify-between items-start mb-4">
<Badge variant={asset.access_type?.public ? "success" : "warning"}>
{asset.access_type?.public ? "公开" : "受限"}
</Badge>
<Badge variant="info">{asset.data_type}</Badge>
</div>
<h3 className="text-lg font-semibold mb-2 truncate">
{asset.metadata_uri.replace("ipfs://", "").slice(0, 20)}...
</h3>
<div className="space-y-2 text-sm text-gray-600 mb-4">
<div className="flex justify-between">
<span>大小</span>
<span>{formatDataSize(asset.data_size)}</span>
</div>
<div className="flex justify-between">
<span>样本</span>
<span>{asset.sample_count || "N/A"}</span>
</div>
<div className="flex justify-between">
<span>版本</span>
<span>v{asset.version}</span>
</div>
<div className="flex justify-between">
<span>销售</span>
<span>{asset.total_sales} 次</span>
</div>
<div className="flex justify-between">
<span>评分</span>
<span>{asset.quality_score}/100</span>
</div>
</div>
<div className="flex justify-between items-center pt-4 border-t">
<div className="text-xl font-bold text-blue-600">
{formatPrice(asset.price)}
</div>
<Button
onClick={() => onPurchase?.(asset.id)}
disabled={!asset.active}
>
购买
</Button>
</div>
<div className="mt-2 text-xs text-gray-400">
提供者: {asset.owner.slice(0, 12)}...
</div>
</div>
</Card>
);
}
6.4 购买流程组件
// src/components/PurchaseModal.tsx
import { useState } from "react";
import { useMutation } from "@tanstack/react-query";
import { DataAsset } from "../sdk/DataMarketplace";
import { Modal } from "./ui/Modal";
import { Button, Input } from "./ui";
import { useWallet } from "../hooks/useWallet";
import { DataMarketplace } from "../sdk/DataMarketplace";
interface Props {
asset: DataAsset | null;
onClose: () => void;
onComplete: (tradeId: string) => void;
}
export function PurchaseModal({ asset, onClose, onComplete }: Props) {
const { client, address } = useWallet();
const [step, setStep] = useState<"confirm" | "verifying" | "complete" | "error">("confirm");
const [error, setError] = useState("");
const purchaseMutation = useMutation({
mutationFn: async () => {
if (!asset) throw new Error("No asset selected");
const marketplace = new DataMarketplace(client, address, {
rpcEndpoint: "", contractAddress: process.env.REACT_APP_MARKETPLACE_ADDRESS!, mnemonic: "",
});
return marketplace.purchaseData(asset.id);
},
onSuccess: () => {
setStep("verifying");
// 实际项目中,需要监听事件或轮询交易状态
setTimeout(() => {
setStep("complete");
onComplete("trade_1");
}, 3000);
},
onError: (e: Error) => {
setError(e.message);
setStep("error");
},
});
if (!asset) return null;
return (
<Modal onClose={onClose}>
<div className="p-6">
<h2 className="text-2xl font-bold mb-4">购买数据</h2>
{step === "confirm" && (
<div>
<div className="bg-gray-50 rounded-lg p-4 mb-6">
<h3 className="font-semibold mb-2">数据资产 #{asset.id}</h3>
<div className="text-sm space-y-1">
<p>类型: {asset.data_type}</p>
<p>大小: {asset.data_size} bytes</p>
<p>版本: {asset.version}</p>
<p>许可证: {JSON.stringify(asset.license)}</p>
<p className="text-xl font-bold text-blue-600 mt-2">
价格: {asset.price.amount} {asset.price.denom}
</p>
</div>
</div>
{error && <p className="text-red-500 mb-4">{error}</p>}
<div className="flex gap-4">
<Button variant="outline" onClick={onClose}>取消</Button>
<Button onClick={() => purchaseMutation.mutate()} loading={purchaseMutation.isPending}>
确认购买
</Button>
</div>
</div>
)}
{step === "verifying" && (
<div className="text-center py-8">
<div className="animate-spin h-8 w-8 border-4 border-blue-600 border-t-transparent rounded-full mx-auto mb-4" />
<p>交易已提交,等待确认...</p>
<p className="text-sm text-gray-500">请等待区块确认</p>
</div>
)}
{step === "complete" && (
<div className="text-center py-8">
<div className="text-green-500 text-4xl mb-4">✓</div>
<h3 className="text-xl font-semibold mb-2">购买成功!</h3>
<p className="text-gray-600 mb-6">数据已购买,请在我的数据中查看</p>
<Button onClick={onClose}>关闭</Button>
</div>
)}
{step === "error" && (
<div className="text-center py-8">
<div className="text-red-500 text-4xl mb-4">✗</div>
<h3 className="text-xl font-semibold mb-2">交易失败</h3>
<p className="text-gray-600 mb-6">{error}</p>
<Button onClick={() => setStep("confirm")}>重试</Button>
</div>
)}
</div>
</Modal>
);
}
6.5 我的数据页面
// src/pages/MyData.tsx
import { useQuery } from "@tanstack/react-query";
import { useWallet } from "../hooks/useWallet";
import { DataMarketplace } from "../sdk/DataMarketplace";
import { Tabs, Tab } from "../components/ui/Tabs";
import { AssetCard } from "../components/AssetCard";
export function MyData() {
const { client, address } = useWallet();
const marketplace = new DataMarketplace(client, address, {
rpcEndpoint: "", contractAddress: process.env.REACT_APP_MARKETPLACE_ADDRESS!, mnemonic: "",
});
const { data: myAssets } = useQuery({
queryKey: ["myAssets", address],
queryFn: () => marketplace.listAssets({ owner: address, activeOnly: false }),
});
const { data: myPurchases } = useQuery({
queryKey: ["myPurchases", address],
queryFn: () => marketplace.listTrades({ buyer: address }),
});
const { data: mySales } = useQuery({
queryKey: ["mySales", address],
queryFn: () => marketplace.listTrades({ seller: address }),
});
return (
<div className="max-w-7xl mx-auto px-4 py-8">
<h1 className="text-3xl font-bold mb-8">我的数据</h1>
<Tabs defaultValue="assets">
<Tab value="assets" label={`发布 (${myAssets?.assets.length || 0})`}>
<div className="grid grid-cols-1 md:grid-cols-3 gap-6 mt-6">
{myAssets?.assets.map(a => <AssetCard key={a.id} asset={a} />)}
</div>
</Tab>
<Tab value="purchases" label={`购买 (${myPurchases?.trades.length || 0})`}>
<div className="mt-6 space-y-4">
{myPurchases?.trades.map(t => (
<div key={t.trade_id} className="bg-white rounded-lg p-4 shadow">
<div className="flex justify-between items-center">
<div>
<p>交易 #{t.trade_id}</p>
<p className="text-sm text-gray-600">资产 #{t.asset_id}</p>
</div>
<div className="text-right">
<p className="font-semibold">{t.amount.amount} {t.amount.denom}</p>
<span className={`text-sm ${
t.status === "Completed" ? "text-green-600" :
t.status === "Disputed" ? "text-red-600" : "text-yellow-600"
}`}>{t.status}</span>
</div>
</div>
</div>
))}
</div>
</Tab>
<Tab value="sales" label={`销售 (${mySales?.trades.length || 0})`}>
<div className="mt-6 space-y-4">
{mySales?.trades.map(t => (
<div key={t.trade_id} className="bg-white rounded-lg p-4 shadow">
<p>交易 #{t.trade_id} - {t.status}</p>
</div>
))}
</div>
</Tab>
</Tabs>
</div>
);
}
6.6 钱包连接 Hook
// src/hooks/useWallet.ts
import { useState, useEffect, createContext, useContext } from "react";
import { SigningCosmWasmClient } from "@cosmjs/cosmwasm-stargate";
import { DirectSecp256k1HdWallet } from "@cosmjs/proto-signing";
interface WalletContext {
client: SigningCosmWasmClient | null;
address: string;
connect: (mnemonic: string) => Promise<void>;
disconnect: () => void;
isConnected: boolean;
}
const WalletCtx = createContext<WalletContext>(null!);
export function WalletProvider({ children }: { children: React.ReactNode }) {
const [client, setClient] = useState<SigningCosmWasmClient | null>(null);
const [address, setAddress] = useState("");
const connect = async (mnemonic: string) => {
const wallet = await DirectSecp256k1HdWallet.fromMnemonic(mnemonic, { prefix: "msg" });
const c = await SigningCosmWasmClient.connectWithSigner(
"https://rpc.msgchain.zone:443", wallet
);
const [account] = await wallet.getAccounts();
setClient(c);
setAddress(account.address);
};
const disconnect = () => {
setClient(null);
setAddress("");
};
return (
<WalletCtx.Provider value={{
client: client!, address, connect, disconnect,
isConnected: !!client,
}}>
{children}
</WalletCtx.Provider>
);
}
export function useWallet() {
return useContext(WalletCtx);
}
6.7 路由配置
// src/App.tsx
import { BrowserRouter, Routes, Route, Link } from "react-router-dom";
import { QueryClient, QueryClientProvider } from "@tanstack/react-query";
import { WalletProvider } from "./hooks/useWallet";
import { Marketplace } from "./pages/Marketplace";
import { MyData } from "./pages/MyData";
import { AssetDetail } from "./pages/AssetDetail";
const queryClient = new QueryClient();
function Layout({ children }: { children: React.ReactNode }) {
return (
<div className="min-h-screen bg-gray-50">
<nav className="bg-white shadow-sm border-b">
<div className="max-w-7xl mx-auto px-4 py-3 flex justify-between items-center">
<div className="flex gap-6 items-center">
<Link to="/" className="text-xl font-bold text-blue-600">DataMarket</Link>
<Link to="/" className="text-gray-600 hover:text-gray-900">市场</Link>
<Link to="/my-data" className="text-gray-600 hover:text-gray-900">我的数据</Link>
</div>
<WalletButton />
</div>
</nav>
<main>{children}</main>
</div>
);
}
function WalletButton() {
const { isConnected, address, connect } = useWallet();
if (isConnected) {
return (
<span className="text-sm text-gray-600 bg-gray-100 px-3 py-1 rounded-full">
{address.slice(0, 8)}...{address.slice(-4)}
</span>
);
}
return (
<button onClick={() => connect("my mnemonic")}
className="bg-blue-600 text-white px-4 py-2 rounded-lg text-sm">
连接钱包
</button>
);
}
export function App() {
return (
<QueryClientProvider client={queryClient}>
<WalletProvider>
<BrowserRouter>
<Layout>
<Routes>
<Route path="/" element={<Marketplace />} />
<Route path="/my-data" element={<MyData />} />
<Route path="/asset/:id" element={<AssetDetail />} />
</Routes>
</Layout>
</BrowserRouter>
</WalletProvider>
</QueryClientProvider>
);
}
7. AI Agent 集成
7.1 Agent 发布数据资产
AI Agent 可以通过 agent_registry_v1 注册后,自动将其产生的数据注册到数据市场:
"""
agent/agent_publisher.py - Agent 自动发布数据资产
"""
import hashlib
import json
from typing import Optional
from msg_marketplace.sdk import DataMarketplaceSDK
class AgentDataPublisher:
"""AI Agent 数据发布器"""
def __init__(self, sdk: DataMarketplaceSDK, agent_id: str, agent_owner: str):
self.sdk = sdk
self.agent_id = agent_id
self.agent_owner = agent_owner
def publish_training_data(
self,
data: bytes,
metadata_uri: str,
price_amount: str,
data_type: str,
version: str = "1.0.0",
description: Optional[str] = None,
) -> str:
"""Agent 发布训练数据"""
data_hash = hashlib.sha256(data).hexdigest()
tx_hash = self.sdk.register_asset(
data_hash=data_hash,
metadata_uri=metadata_uri,
price={"denom": "uMSG", "amount": price_amount},
data_size=str(len(data)),
data_type=data_type,
version=version,
agent_id=self.agent_id,
)
return tx_hash
def publish_inference_results(
self,
input_description: str,
results: list,
price_per_sample: str,
) -> str:
"""Agent 发布推理结果"""
data = json.dumps({
"agent_id": self.agent_id,
"input": input_description,
"results": results,
"created_at": int(time.time()),
}).encode()
return self.publish_training_data(
data=data,
metadata_uri=f"ipfs://inference_{self.agent_id}_{int(time.time())}",
price_amount=price_per_sample,
data_type="inference_result",
version="1.0.0",
)
// agent/AgentPublisher.ts
import { DataMarketplace } from "../sdk/DataMarketplace";
export class AgentPublisher {
constructor(
private marketplace: DataMarketplace,
private agentId: string
) {}
async publishEmbeddings(
embeddings: Float32Array[],
metadata: { dimension: number; model: string; description: string },
price: { denom: string; amount: string }
): Promise<string> {
const data = new TextEncoder().encode(JSON.stringify({
agent_id: this.agentId,
embeddings: embeddings.map(e => Array.from(e)),
metadata,
}));
const hash = await this.marketplace.computeDataHash(data);
return this.marketplace.registerAsset({
dataHash: hash,
metadataUri: `ipfs://emb_${this.agentId}`,
price,
license: "MIT",
accessType: { public: {} },
dataSize: String(data.length),
dataType: "embeddings",
version: "1.0.0",
agentId: this.agentId,
});
}
}
7.2 Agent 发现和购买数据
"""
agent/agent_consumer.py - Agent 自动发现和购买数据
"""
from typing import List, Optional
from msg_marketplace.sdk import DataMarketplaceSDK, DataAsset
class AgentDataConsumer:
"""AI Agent 数据消费者"""
def __init__(self, sdk: DataMarketplaceSDK, agent_id: str):
self.sdk = sdk
self.agent_id = agent_id
self.purchase_history: List[str] = []
def find_data_for_training(
self,
data_type: str,
max_budget: str = "1000000",
limit: int = 10,
) -> List[DataAsset]:
"""查找适合训练的数据"""
assets, _ = self.sdk.list_assets(
data_type=data_type,
active_only=True,
limit=limit,
)
# 按质量评分排序
assets.sort(key=lambda a: (a.quality_score, int(a.total_sales)), reverse=True)
# 过滤预算范围内的
affordable = [
a for a in assets
if int(a.price["amount"]) <= int(max_budget)
]
return affordable
def auto_purchase_and_verify(
self,
asset_id: str,
fetch_fn,
) -> bool:
"""自动购买并验证数据"""
# 获取资产信息
asset = self.sdk.get_asset(asset_id)
# 检查许可证兼容性
if not self._check_license_compatible(asset.license):
print(f"License {asset.license} not compatible with agent {self.agent_id}")
return False
# 购买
tx = self.sdk.purchase_data(asset_id)
self.purchase_history.append(asset_id)
# 获取并验证数据
data = fetch_fn()
if self.sdk.verify_data_integrity(data, asset.data_hash):
# 确认交付
trades, _ = self.sdk.list_trades(
buyer=self.sdk.sender,
asset_id=asset_id,
status="PaidAwaitingDelivery",
)
if trades:
self.sdk.confirm_delivery(trades[0].trade_id, asset.data_hash)
return True
else:
trades, _ = self.sdk.list_trades(
buyer=self.sdk.sender,
asset_id=asset_id,
status="PaidAwaitingDelivery",
)
if trades:
self.sdk.dispute_delivery(trades[0].trade_id, "Hash verification failed")
return False
def _check_license_compatible(self, license_data) -> bool:
"""检查许可证是否与 Agent 的用途兼容"""
if isinstance(license_data, str):
compatible = {
"MIT": True,
"CCBY40": True,
"CCNC40": False, # Agent 商用不兼容
"TrainingOnly": True,
"NoResale": True,
}
return compatible.get(license_data, False)
return True
7.3 自动化数据管线
"""
agent/data_pipeline.py - Agent 数据管线:购买→下载→处理→训练
"""
import json
import time
from typing import Callable, Optional
from msg_marketplace.sdk import DataMarketplaceSDK
class DataPipeline:
"""自动化数据管线"""
def __init__(
self,
sdk: DataMarketplaceSDK,
consumer: AgentDataConsumer,
process_fn: Callable,
train_fn: Callable,
):
self.sdk = sdk
self.consumer = consumer
self.process_fn = process_fn
self.train_fn = train_fn
def run(
self,
data_type: str,
max_budget: str,
max_datasets: int = 3,
) -> dict:
"""执行完整的数据采集→处理→训练管线"""
results = {
"datasets_found": 0,
"datasets_purchased": 0,
"training_started": False,
"training_completed": False,
"errors": [],
}
# 步骤 1: 发现数据
print(f"[Pipeline] Searching for {data_type} data...")
assets = self.consumer.find_data_for_training(data_type, max_budget, max_datasets)
results["datasets_found"] = len(assets)
if not assets:
results["errors"].append("No suitable datasets found")
return results
# 步骤 2: 购买并验证数据
purchased_data = []
for asset in assets[:max_datasets]:
print(f"[Pipeline] Purchasing asset {asset.id}...")
try:
# 模拟数据获取
def fetch_sample():
return json.dumps({"mock": "data", "asset_id": asset.id}).encode()
success = self.consumer.auto_purchase_and_verify(asset.id, fetch_sample)
if success:
purchased_data.append(asset)
results["datasets_purchased"] += 1
else:
print(f"[Pipeline] Purchase/verification failed for asset {asset.id}")
except Exception as e:
results["errors"].append(f"Purchase failed for {asset.id}: {str(e)}")
# 步骤 3: 数据处理
if purchased_data:
print(f"[Pipeline] Processing {len(purchased_data)} datasets...")
try:
processed = self.process_fn(purchased_data)
print(f"[Pipeline] Processing complete: {len(processed)} samples")
except Exception as e:
results["errors"].append(f"Processing failed: {str(e)}")
return results
# 步骤 4: 模型训练
print("[Pipeline] Starting training...")
results["training_started"] = True
try:
self.train_fn(processed)
results["training_completed"] = True
print("[Pipeline] Training completed!")
except Exception as e:
results["errors"].append(f"Training failed: {str(e)}")
return results
7.4 Agent A2A 集成 - 数据协商
AI Agent 可以通过 A2A 协议协商数据交易:
// agent/A2ADataNegotiation.ts
import { A2AClient } from "@msgchain/a2a-sdk";
interface DataNegotiationOffer {
assetId: string;
offeredPrice: string;
license: string;
intendedUse: string;
}
interface DataNegotiationResponse {
accepted: boolean;
counterOffer?: string;
additionalTerms?: string;
}
export class A2ADataNegotiation {
private a2a: A2AClient;
constructor(a2aClient: A2AClient) {
this.a2a = a2aClient;
}
async negotiateDataPurchase(
sellerAgentId: string,
assetId: string,
maxPrice: string
): Promise<DataNegotiationResponse | null> {
const offer: DataNegotiationOffer = {
assetId,
offeredPrice: maxPrice,
license: "TrainingOnly",
intendedUse: "model_training",
};
// 通过 A2A 发送协商请求
const response = await this.a2a.sendMessage(sellerAgentId, {
type: "data_negotiation",
payload: offer,
replyTimeout: 30000, // 30 秒超时
});
if (response.type === "data_negotiation_response") {
return response.payload as DataNegotiationResponse;
}
return null;
}
async handleIncomingOffer(
message: any
): Promise<DataNegotiationResponse> {
const offer = message.payload as DataNegotiationOffer;
// Agent 的决策逻辑
const minAcceptablePrice = "500000"; // 最低可接受价格
const accepted = BigInt(offer.offeredPrice) >= BigInt(minAcceptablePrice);
return {
accepted,
counterOffer: accepted ? undefined : minAcceptablePrice,
};
}
}
7.5 Agent 注册与市场集成
"""
agent/agent_setup.py - Agent 注册和市场集成
"""
from msg_marketplace.sdk import DataMarketplaceSDK
class AgentMarketplaceIntegration:
"""Agent 与数据市场的完整集成"""
def __init__(self, agent_id: str, registry_address: str, marketplace_address: str):
self.agent_id = agent_id
self.registry_address = registry_address
self.marketplace_address = marketplace_address
self.publisher: Optional[AgentDataPublisher] = None
self.consumer: Optional[AgentDataConsumer] = None
def initialize(self, mnemonic: str, rpc_endpoint: str):
"""初始化 Agent 的市场集成"""
sdk = DataMarketplaceSDK(
rpc_endpoint=rpc_endpoint,
contract_address=self.marketplace_address,
sender_address="",
)
self.publisher = AgentDataPublisher(sdk, self.agent_id, sdk.sender)
self.consumer = AgentDataConsumer(sdk, self.agent_id)
print(f"Agent {self.agent_id} initialized for marketplace")
def auto_publish_data(self, data: bytes, data_type: str, price: str):
"""自动发布数据"""
if not self.publisher:
raise RuntimeError("Agent not initialized")
return self.publisher.publish_training_data(
data=data,
metadata_uri=f"ipfs://agent_{self.agent_id}_{data_type}",
price_amount=price,
data_type=data_type,
)
def auto_acquire_training_data(self, data_type: str, budget: str):
"""自动获取训练数据"""
if not self.consumer:
raise RuntimeError("Agent not initialized")
return self.consumer.find_data_for_training(
data_type=data_type,
max_budget=budget,
)
8. 扩展功能
8.1 数据捆绑包
多个数据集打包销售,一次性购买多个数据资产:
// src/bundle.rs - 数据捆绑包
use cosmwasm_std::{Addr, Coin, DepsMut, Env, MessageInfo, Response, Uint128};
use cw_storage_plus::Map;
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct DataBundle {
pub id: Uint128,
pub owner: Addr,
pub name: String,
pub asset_ids: Vec<Uint128>,
pub bundle_price: Coin,
pub discount_bps: u64, // 折扣基点
pub active: bool,
pub created_at: Timestamp,
}
pub const BUNDLES: Map<Uint128, DataBundle> = Map::new("bundles");
pub fn execute_create_bundle(
deps: DepsMut,
env: Env,
info: MessageInfo,
name: String,
asset_ids: Vec<Uint128>,
discount_bps: u64,
) -> ContractResult<Response> {
// 验证调用者是所有资产的所有者
for aid in &asset_ids {
let asset = ASSETS.load(deps.storage, *aid)?;
if asset.owner != info.sender {
return Err(ContractError::Unauthorized {});
}
}
let config = CONFIG.load(deps.storage)?;
let new_id = config.bundle_count.checked_add(Uint128::one())?;
// 计算总价(含折扣)
let mut total = Uint128::zero();
for aid in &asset_ids {
let asset = ASSETS.load(deps.storage, *aid)?;
total = total.checked_add(asset.price.amount)?;
}
let discount = total.multiply_ratio(discount_bps, 10000u128);
let bundle_price = Coin {
denom: config.min_price.denom,
amount: total.checked_sub(discount)?,
};
let bundle = DataBundle {
id: new_id,
owner: info.sender,
name,
asset_ids,
bundle_price,
discount_bps,
active: true,
created_at: env.block.time,
};
BUNDLES.save(deps.storage, new_id, &bundle)?;
let mut new_config = config;
new_config.bundle_count = new_id;
CONFIG.save(deps.storage, &new_config)?;
Ok(Response::new()
.add_attribute("action", "create_bundle")
.add_attribute("bundle_id", new_id.to_string()))
}
// src/sdk/BundleSDK.ts
export class BundleSDK {
constructor(private marketplace: DataMarketplace) {}
async createBundle(params: {
name: string; assetIds: string[]; discountBps: number;
}): Promise<string> {
const msg = {
create_bundle: {
name: params.name,
asset_ids: params.assetIds,
discount_bps: params.discountBps,
},
};
const fee = calculateFee(300_000, GasPrice.fromString("1000000000attoMSG"));
const r = await this.marketplace.client.execute(
this.marketplace.sender, this.marketplace.config.contractAddress, msg, fee
);
return r.transactionHash;
}
async purchaseBundle(bundleId: string): Promise<string> {
const bundle = await this.queryBundle(bundleId);
const msg = { purchase_bundle: { bundle_id: bundleId } };
const fee = calculateFee(400_000, GasPrice.fromString("1000000000attoMSG"));
const r = await this.marketplace.client.execute(
this.marketplace.sender, this.marketplace.config.contractAddress, msg, fee, undefined,
[bundle.bundle_price]
);
return r.transactionHash;
}
async queryBundle(bundleId: string): Promise<DataBundle> {
return this.marketplace.client.queryContractSmart(
this.marketplace.config.contractAddress, { get_bundle: { bundle_id: bundleId } }
);
}
}
8.2 数据流订阅
支持数据流(Streaming Data)的订阅模式,数据持续更新:
// src/streaming.rs
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct DataStream {
pub stream_id: Uint128,
pub asset_id: Uint128,
pub endpoint: String, // 数据流端点
pub update_frequency: u64, // 更新间隔(区块)
pub last_update: Timestamp,
pub active: bool,
}
pub const STREAMS: Map<Uint128, DataStream> = Map::new("streams");
pub fn execute_register_stream(
deps: DepsMut,
env: Env,
info: MessageInfo,
asset_id: Uint128,
endpoint: String,
update_frequency: u64,
) -> ContractResult<Response> {
let asset = ASSETS.load(deps.storage, asset_id)?;
if asset.owner != info.sender { return Err(ContractError::Unauthorized {}); }
let config = CONFIG.load(deps.storage)?;
let new_id = config.stream_count.checked_add(Uint128::one())?;
let stream = DataStream {
stream_id: new_id,
asset_id,
endpoint,
update_frequency,
last_update: env.block.time,
active: true,
};
STREAMS.save(deps.storage, new_id, &stream)?;
Ok(Response::new().add_attribute("action", "register_stream").add_attribute("stream_id", new_id.to_string()))
}
8.3 协作数据集
多个 Agent 共同贡献数据并共享收益:
// src/collaborative.rs
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct CollaborativeDataset {
pub dataset_id: Uint128,
pub name: String,
pub contributors: Vec<Contributor>,
pub total_supply: Uint128,
pub price: Coin,
pub revenue_distribution: Vec<RevenueShare>,
}
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct Contributor {
pub address: Addr,
pub contribution_share: Uint128, // 万分比
pub sample_count: Uint128,
pub data_hash: String,
}
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct RevenueShare {
pub address: Addr,
pub share_bps: u64,
}
pub const COLLABORATIVE: Map<Uint128, CollaborativeDataset> = Map::new("collaborative");
pub fn execute_create_collaborative(
deps: DepsMut,
env: Env,
info: MessageInfo,
name: String,
contributors: Vec<(Addr, Uint128, String)>,
price: Coin,
) -> ContractResult<Response> {
let config = CONFIG.load(deps.storage)?;
let new_id = config.collab_count.checked_add(Uint128::one())?;
let total_share: Uint128 = contributors.iter().map(|(_, s, _)| s).sum();
if total_share != Uint128::from(10000u128) {
return Err(ContractError::InvalidShare {});
}
let collab = CollaborativeDataset {
dataset_id: new_id,
name,
contributors: contributors.into_iter()
.map(|(addr, share, hash)| Contributor {
address: addr,
contribution_share: share,
sample_count: Uint128::zero(),
data_hash: hash,
})
.collect(),
total_supply: Uint128::zero(),
price,
revenue_distribution: vec![],
};
COLLABORATIVE.save(deps.storage, new_id, &collab)?;
Ok(Response::new().add_attribute("action", "create_collaborative").add_attribute("dataset_id", new_id.to_string()))
}
8.4 转售版税
数据资产首次销售后,后续转售原作者获得版税:
// src/royalty.rs
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct RoyaltyConfig {
pub asset_id: Uint128,
pub creator: Addr,
pub royalty_bps: u64, // 版税率(基点)
pub total_royalty_collected: Coin,
}
pub const ROYALTIES: Map<Uint128, RoyaltyConfig> = Map::new("royalties");
pub fn execute_set_royalty(
deps: DepsMut,
env: Env,
info: MessageInfo,
asset_id: Uint128,
royalty_bps: u64,
) -> ContractResult<Response> {
if royalty_bps > 5000 { return Err(ContractError::InvalidRoyalty {}); } // 最高 50%
let asset = ASSETS.load(deps.storage, asset_id)?;
if asset.owner != info.sender { return Err(ContractError::Unauthorized {}); }
// 首次销售时设置创作者
// 如果是转售,创作者不变
if ROYALTIES.has(deps.storage, asset_id) {
return Err(ContractError::RoyaltyAlreadySet {});
}
let royalty = RoyaltyConfig {
asset_id,
creator: info.sender.clone(),
royalty_bps,
total_royalty_collected: Coin { denom: "uMSG".to_string(), amount: Uint128::zero() },
};
ROYALTIES.save(deps.storage, asset_id, &royalty)?;
Ok(Response::new()
.add_attribute("action", "set_royalty")
.add_attribute("asset_id", asset_id.to_string())
.add_attribute("royalty_bps", royalty_bps.to_string()))
}
/// 转售时计算版税
pub fn calculate_royalty(
storage: &dyn Storage,
asset_id: Uint128,
sale_price: &Coin,
) -> Option<(Addr, Coin)> {
if let Ok(royalty) = ROYALTIES.load(storage, asset_id) {
let royalty_amount = sale_price.amount.multiply_ratio(royalty.royalty_bps, 10000u128);
Some((royalty.creator, Coin {
denom: sale_price.denom.clone(),
amount: royalty_amount,
}))
} else {
None
}
}
8.5 扩展功能的使用示例
// examples/advanced_features.ts
async function main() {
const marketplace = await DataMarketplace.connect({...});
// 创建数据捆绑包(8 折)
const bundleTx = await marketplace.registerAsset({
name: "NLP Bundle",
assetIds: ["1", "2", "3"],
discountBps: 2000, // 20% off
});
// 购买捆绑包
await marketplace.purchaseBundle("bundle_1");
// 设置转售版税
const royaltyTx = await marketplace.executeContract({
set_royalty: { asset_id: "1", royalty_bps: 1000 }, // 10%
});
// 创建协作数据集
const collabTx = await marketplace.executeContract({
create_collaborative: {
name: "Multi-Agent Training Data",
contributors: [
["msg1agent1...", "5000", "hash1..."],
["msg1agent2...", "5000", "hash2..."],
],
price: { denom: "uMSG", amount: "5000000" },
},
});
console.log("Advanced features enabled:", { bundleTx, royaltyTx, collabTx });
}
附录
A. API 参考
合约查询接口
| 方法 | 参数 | 描述 |
|---|---|---|
get_asset |
asset_id |
获取数据资产详情 |
list_assets |
start_after, limit, owner, data_type, min_price, max_price, license, active_only |
列出数据资产 |
search_assets |
query, start_after, limit |
搜索数据资产 |
get_trade |
trade_id |
获取交易记录 |
list_trades |
start_after, limit, buyer, seller, asset_id, status |
列出交易记录 |
get_escrow |
escrow_id |
获取托管账户 |
check_access |
asset_id, address |
检查访问权限 |
get_stats |
无 | 获取市场统计 |
get_bundle |
bundle_id |
获取数据捆绑包 |
get_subscription |
subscription_id |
获取订阅信息 |
list_subscriptions |
start_after, limit, subscriber, asset_id |
列出订阅 |
合约执行接口
| 方法 | 所需资金 | 描述 |
|---|---|---|
register_asset |
无 | 注册新数据资产 |
update_asset |
无 | 更新资产信息 |
deactivate_asset |
无 | 注销资产 |
purchase_data |
资产价格 | 购买数据 |
confirm_delivery |
无 | 确认数据交付 |
dispute_delivery |
无 | 争议交付 |
request_refund |
无 | 请求退款 |
arbitrate |
无 | 仲裁裁决 |
add_to_whitelist |
无 | 添加白名单 |
remove_from_whitelist |
无 | 移除白名单 |
grant_access |
无 | 授权访问 |
revoke_access |
无 | 撤销授权 |
update_quality_score |
无 | 更新质量评分 |
create_bundle |
无 | 创建数据捆绑包 |
set_royalty |
无 | 设置转售版税 |
B. Gas 估算
| 操作 | Gas 估算 | 费用 (@0.025 uMSG) |
|---|---|---|
register_asset |
200,000 | 5,000 uMSG |
purchase_data |
300,000 | 7,500 uMSG |
confirm_delivery |
150,000 | 3,750 uMSG |
dispute_delivery |
200,000 | 5,000 uMSG |
arbitrate |
250,000 | 6,250 uMSG |
add_to_whitelist (50 address) |
300,000 | 7,500 uMSG |
create_bundle |
250,000 | 6,250 uMSG |
register_stream |
200,000 | 5,000 uMSG |
C. 安全注意事项
- 数据哈希验证: 买方在确认交付前必须验证数据哈希。建议使用 Merkle 树验证大数据集
- 加密传输: 数据通过 IPFS 等去中心化存储传输时,应使用 AES-256-GCM 加密
- 访问控制: 对于私人数据,建议使用 Whitelist 或 TokenGated 模式
- 争议超时: 合同设置了最长仲裁时间,超时后自动触发默认操作
- Oracle 安全: 使用去中心化的 Oracle 网络进行数据验证,避免单点故障
- 重入攻击防护: 所有涉及资金转移的操作都遵循 Checks-Effects-Interactions 模式
- 版税上限: 转售版税率限制在 50% 以下,防止过度收费
D. 开发资源
- MSG Chain 文档: https://docs.msgchain.zone
- CosmWasm 文档: https://docs.cosmwasm.com
- 钱包: Keplr / Leap 钱包(支持 MSG Chain)
- 区块浏览器: https://explorer.msgchain.zone
- 水龙头: https://faucet.msgchain.zone
E. 常见问题
Q: 如何确保数据质量?
A: 使用链上哈希承诺 + 链下验证 + 信誉系统的多重机制。买方验证哈希后确认,不匹配可发起争议。
Q: 数据隐私如何保护?
A: 数据在链下传输时可以使用加密,链上只存储哈希承诺。支持私有数据通过白名单或 Token 门控访问。
Q: Agent 如何自动化参与市场?
A: Agent 通过 agent_registry_v1 注册后,使用 SDK 自动发布数据、发现数据、购买和验证。数据管线可完全自动化。
Q: 交易手续费如何计算?
A: 协议手续费为交易金额的 1%(可配置),加上链上 Gas 费用。Gas 价格: 1,000,000,000 attoMSG/gas(flat rate)。
Q: 争议如何解决?
A: 买方发起争议后,卖方和买方提交证据,仲裁人(合约管理员)裁决。争议期间资金锁定在托管合约中。
Q: 是否支持代币以外的支付方式?
A: 当前仅支持 MSG 链原生代币 uMSG。可通过扩展支持 CW20 代币。
文档版本: v1.0.0
兼容链: msg-chain-1
作者: MSG Chain Developer Docs Team
