dApp Docs/AI Agent数据市场实现
Development reference. Not independently verified for production.

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. 概述
  2. 数据资产合约
  3. 数据交易流程
  4. 访问控制
  5. 数据验证与仲裁
  6. 前端实现
  7. AI Agent 集成
  8. 扩展功能
    附录

1. 概述

1.1 为什么 AI Agent 需要数据市场

AI Agent 的核心能力取决于其训练数据和推理数据的质量与多样性。在 MSG Chain 的生态中,Agent 已经能够通过 agent_registry_v1 注册发现、通过 agent_payment_v1 (AIPAY) 支付、通过 agent_a2a_v1 通信协作。然而,数据仍然是割裂的——每个 Agent 各自为战,无法高效共享和交易数据。

数据市场 (Data Marketplace) 解决了以下关键问题:

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 核心流程

  1. 数据提供方 注册数据资产 → 提交哈希承诺 → 设定价格和许可证
  2. 数据消费方 发现数据 → 发起购买(AIPAY 条件支付)
  3. 系统锁定资金至托管合约 → 提供方发送数据解密密钥/访问链接
  4. 消费方验证数据哈希 → 匹配则确认收货,资金释放
  5. 争议处理 → 哈希不匹配则进入仲裁流程

1.6 经济模型

数据市场使用 MSG 链原生代币 uMSG(18 位精度)进行交易:


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 哈希承诺-揭示机制

数据验证使用链下承诺+链上验证模式:

  1. 注册时: 数据提供方计算 SHA256(data),将哈希值写入链上
  2. 交付后: 买方获取数据,重新计算哈希,与链上承诺比对
  3. 匹配: 确认交付,资金释放给卖方
  4. 不匹配: 发起争议,进入仲裁流程
注册: 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 技术栈

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. 安全注意事项

  1. 数据哈希验证: 买方在确认交付前必须验证数据哈希。建议使用 Merkle 树验证大数据集
  2. 加密传输: 数据通过 IPFS 等去中心化存储传输时,应使用 AES-256-GCM 加密
  3. 访问控制: 对于私人数据,建议使用 Whitelist 或 TokenGated 模式
  4. 争议超时: 合同设置了最长仲裁时间,超时后自动触发默认操作
  5. Oracle 安全: 使用去中心化的 Oracle 网络进行数据验证,避免单点故障
  6. 重入攻击防护: 所有涉及资金转移的操作都遵循 Checks-Effects-Interactions 模式
  7. 版税上限: 转售版税率限制在 50% 以下,防止过度收费

D. 开发资源

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