去中心化AI算力市场实现指南
链:
msg-chain-1| Bech32:msg| 精度: 18 decimals
主网状态: No-Go
核心协议: AIPAY (支付) + Compute Registry (注册发现) + Match Engine (匹配)
目标: 构建类似 Akash Network 的去中心化算力市场,让 GPU/CPU 提供者与 AI Agent 消费者直接匹配
目录
1. 概述
1.1 Why Decentralized Compute for AI
AI Agent 的推理 (inference)、训练 (training)、微调 (fine-tuning) 需要大量 GPU 算力。传统云服务(AWS、GCP、Azure)存在以下问题:
| 问题 | 传统云 | 去中心化市场 |
|---|---|---|
| 成本 | 高溢价,锁定合同 | 自由定价,市场竞争 |
| 可用性 | 区域受限,配额制 | 全球节点,无许可 |
| 隐私 | 数据必须上传 | 可执行 WASM/TEE,隐私保护 |
| 抗审查 | 平台可拒绝服务 | 无许可接入 |
去中心化算力市场让算力提供商(Provider)和 AI Agent 消费者(Client)直接在 MSG Chain 上匹配,通过智能合约托管资金,确保双方公平。
1.2 市场模型
+-----------------------------+
| MSG Chain |
| msg-chain-1 |
+----------+----------------+
|
+-----------+-----------+
| |
v v
+--------------------+ +--------------------+
| Compute Registry | | AIPAY Payment |
| (Provider DID) | | (Escrow/Release) |
+--------------------+ +--------------------+
| |
v v
+--------------------+ +--------------------+
| Task Manager |<->| Match Engine |
| (Order Book) | |(Bid/Ask/R.Auction) |
+--------------------+ +--------------------+
- Provider: 注册 GPU/CPU 资源,设定价格,运行容器
- Client: 发布计算任务,锁定资金,获取结果
- Matcher: 链下/链上匹配引擎,撮合供需
- AIPAY: 托管支付,任务完成释放,争议仲裁
1.3 MSG Chain 集成点
| 组件 | 合约/模块 | 作用 |
|---|---|---|
| Provider Registry | compute_registry.wasm |
算力提供者 DID 注册、资源声明、状态管理 |
| Task Manager | compute_task.wasm |
任务发布、状态机、生命周期 |
| Match Engine | 链下 Python/Go | 订单簿匹配、反向拍卖 |
| Payment | AIPAY 合约 | 托管、释放、退款、争议 |
| Execution | CosmWasm + WASM | 容器化执行与结果验证 |
1.4 经济模型
- 结算币种: AIPAY (18 decimals)
- Provider 质押: 注册需锁定
1_000_000 AIPAY防止女巫攻击 - 市场费率: 每笔交易收取
0.5%协议费 - 惩罚: 未完成任务扣除质押的
10%
1.5 安全假设
- Provider 需证明拥有所声明的资源(远程证明 / PoC)
- 任务结果通过哈希提交,争议时由仲裁者重新执行
- 所有通信通过 MSG Chain 签名消息验证身份
2. 算力提供者注册
2.1 数据模型
// contracts/compute-registry/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 struct ComputeProvider {
/// 所有者地址
pub owner: Addr,
/// 提供者 DID(分布式身份)
pub provider_did: String,
/// 可用资源
pub resources: Resources,
/// 定价策略
pub pricing: Pricing,
/// 地理区域
pub region: String,
/// 当前状态
pub status: ProviderStatus,
/// 累计收入
pub total_earned: Uint128,
/// 注册时间
pub registered_at: Timestamp,
/// 上次更新
pub updated_at: Timestamp,
/// 质押金额
pub bond_amount: Coin,
/// 完成任务数
pub completed_tasks: u64,
/// 失败任务数
pub failed_tasks: u64,
/// 平均评分 (0-1000, 除10得实际分数)
pub avg_rating: u32,
}
/// 资源声明
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct Resources {
/// GPU 型号,例如 "A100", "H100", "RTX4090", "L40S"
pub gpu_model: String,
/// GPU 数量
pub gpu_count: u32,
/// vCPU 核心数
pub vcpu: u32,
/// 内存 GB
pub ram_gb: u32,
/// 存储 GB
pub storage_gb: u32,
/// 网络带宽 Mbps
pub bandwidth_mbps: u32,
/// 是否支持 TEE
pub supports_tee: bool,
/// TEE 类型 (如 "Intel SGX", "AMD SEV", "Nitro")
pub tee_type: Option<String>,
}
/// 定价策略
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct Pricing {
/// 每小时 GPU 价格 (AIPAY, 18 decimals)
pub price_per_gpu_hour: Uint128,
/// 每小时 vCPU 价格
pub price_per_vcpu_hour: Uint128,
/// 每小时内存价格 (per GB)
pub price_per_ram_gb_hour: Uint128,
/// 每小时存储价格 (per GB)
pub price_per_storage_gb_hour: Uint128,
/// 最低计费时间(秒)
pub min_billing_duration: u64,
/// 是否接受议价
pub negotiable: bool,
}
/// 提供者状态
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub enum ProviderStatus {
/// 在线,可接受任务
Active,
/// 离线,不接受新任务
Inactive,
/// 暂停(维护中)
Paused,
/// 已被惩罚冻结
Slashed,
/// 已注销
Revoked,
}
/// 提供者声誉记录
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct ReputationRecord {
pub provider: Addr,
pub client: Addr,
pub task_id: String,
pub rating: u8,
pub comment: String,
pub rated_at: Timestamp,
}
/// 提供者硬件证明
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct HardwareAttestation {
pub provider: Addr,
pub gpu_uuids: Vec<String>,
pub cpu_info: String,
pub ram_total_gb: u32,
pub attestation_doc: String,
pub verified: bool,
pub verified_at: Option<Timestamp>,
}
// 存储
pub const PROVIDERS: Map<&Addr, ComputeProvider> = Map::new("providers");
pub const PROVIDER_BY_DID: Map<&str, Addr> = Map::new("provider_did");
pub const REPUTATIONS: Map<&str, ReputationRecord> = Map::new("reputations");
pub const ATTESTATIONS: Map<&Addr, HardwareAttestation> = Map::new("attestations");
pub const PROVIDER_COUNT: Item<u64> = Item::new("provider_count");
pub const TOTAL_STAKED: Item<Uint128> = Item::new("total_staked");
2.2 合约消息定义
// contracts/compute-registry/src/msg.rs
use cosmwasm_std::Coin;
use schemars::JsonSchema;
use serde::{Deserialize, Serialize};
use crate::state::{Pricing, ProviderStatus, Resources};
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct InstantiateMsg {
pub admin: String,
pub min_bond: Coin,
pub protocol_fee_bps: u64,
pub dispute_duration: u64,
}
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
#[serde(rename_all = "snake_case")]
pub enum ExecuteMsg {
Register {
provider_did: String,
resources: Resources,
pricing: Pricing,
region: String,
},
UpdateResources { resources: Resources },
UpdatePricing { pricing: Pricing },
UpdateRegion { region: String },
UpdateStatus { status: ProviderStatus },
SubmitAttestation {
gpu_uuids: Vec<String>,
cpu_info: String,
ram_total_gb: u32,
attestation_doc: String,
},
WithdrawEarnings { amount: Option<Coin> },
Revoke {},
RateProvider {
provider: String,
task_id: String,
rating: u8,
comment: String,
},
}
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
#[serde(rename_all = "snake_case")]
pub enum QueryMsg {
GetProvider { address: String },
GetProviderByDid { did: String },
ListProviders {
start_after: Option<String>,
limit: Option<u32>,
status_filter: Option<ProviderStatus>,
},
QueryAvailableProviders {
min_gpu_count: Option<u32>,
min_vcpu: Option<u32>,
min_ram_gb: Option<u32>,
min_storage_gb: Option<u32>,
region: Option<String>,
gpu_model: Option<String>,
max_price_per_hour: Option<Uint128>,
supports_tee: Option<bool>,
},
GetReputation { provider: String },
GetStats {},
GetMinBond {},
}
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct ProviderListResponse {
pub providers: Vec<ComputeProvider>,
pub total_count: u64,
}
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct AvailableProvidersResponse {
pub providers: Vec<ComputeProvider>,
pub count: u32,
}
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct StatsResponse {
pub total_providers: u64,
pub active_providers: u64,
pub total_staked: Uint128,
pub total_earned: Uint128,
pub total_tasks_completed: u64,
}
2.3 合约实现
// contracts/compute-registry/src/contract.rs
#[entry_point]
pub fn instantiate(
deps: DepsMut,
_env: Env,
_info: MessageInfo,
msg: InstantiateMsg,
) -> StdResult<Response> {
set_contract_version(deps.storage, "compute-registry", "1.0.0")?;
ADMIN.save(deps.storage, &deps.api.addr_validate(&msg.admin)?)?;
MIN_BOND.save(deps.storage, &msg.min_bond)?;
PROTOCOL_FEE_BPS.save(deps.storage, &msg.protocol_fee_bps)?;
DISPUTE_DURATION.save(deps.storage, &msg.dispute_duration)?;
PROVIDER_COUNT.save(deps.storage, &0u64)?;
TOTAL_STAKED.save(deps.storage, &Uint128::zero())?;
Ok(Response::new()
.add_attribute("action", "instantiate")
.add_attribute("admin", msg.admin)
.add_attribute("min_bond", msg.min_bond.to_string()))
}
#[entry_point]
pub fn execute(deps: DepsMut, env: Env, info: MessageInfo, msg: ExecuteMsg)
-> Result<Response, ContractError> {
match msg {
ExecuteMsg::Register { provider_did, resources, pricing, region } =>
execute_register(deps, env, info, provider_did, resources, pricing, region),
ExecuteMsg::UpdateResources { resources } => execute_update_resources(deps, env, info, resources),
ExecuteMsg::UpdatePricing { pricing } => execute_update_pricing(deps, env, info, pricing),
ExecuteMsg::UpdateRegion { region } => execute_update_region(deps, env, info, region),
ExecuteMsg::UpdateStatus { status } => execute_update_status(deps, env, info, status),
ExecuteMsg::SubmitAttestation { gpu_uuids, cpu_info, ram_total_gb, attestation_doc } =>
execute_submit_attestation(deps, env, info, gpu_uuids, cpu_info, ram_total_gb, attestation_doc),
ExecuteMsg::WithdrawEarnings { amount } => execute_withdraw_earnings(deps, env, info, amount),
ExecuteMsg::Revoke {} => execute_revoke(deps, env, info),
ExecuteMsg::RateProvider { provider, task_id, rating, comment } =>
execute_rate_provider(deps, env, info, provider, task_id, rating, comment),
}
}
2.4 Register Execution
pub fn execute_register(
deps: DepsMut,
env: Env,
info: MessageInfo,
provider_did: String,
resources: Resources,
pricing: Pricing,
region: String,
) -> Result<Response, ContractError> {
if PROVIDER_BY_DID.may_load(deps.storage, &provider_did)?.is_some() {
return Err(ContractError::DidAlreadyExists {});
}
if PROVIDERS.may_load(deps.storage, &info.sender)?.is_some() {
return Err(ContractError::AlreadyRegistered {});
}
let min_bond = MIN_BOND.load(deps.storage)?;
let sent_bond = info.funds.iter()
.find(|c| c.denom == min_bond.denom)
.ok_or(ContractError::MissingBond {})?;
if sent_bond.amount < min_bond.amount {
return Err(ContractError::InsufficientBond {
required: min_bond.amount, sent: sent_bond.amount,
});
}
if resources.gpu_model.is_empty() {
return Err(ContractError::InvalidResource {
msg: "gpu_model cannot be empty".to_string() });
}
if resources.gpu_count == 0 && resources.vcpu == 0 {
return Err(ContractError::InvalidResource {
msg: "must provide at least gpu or vcpu".to_string() });
}
if pricing.price_per_gpu_hour.is_zero()
&& pricing.price_per_vcpu_hour.is_zero()
&& pricing.price_per_ram_gb_hour.is_zero()
&& pricing.price_per_storage_gb_hour.is_zero()
{
return Err(ContractError::InvalidPricing {
msg: "at least one price must be non-zero".to_string() });
}
let provider = ComputeProvider {
owner: info.sender.clone(),
provider_did: provider_did.clone(),
resources, pricing, region,
status: ProviderStatus::Active,
total_earned: Uint128::zero(),
registered_at: env.block.time,
updated_at: env.block.time,
bond_amount: sent_bond.clone(),
completed_tasks: 0, failed_tasks: 0, avg_rating: 0,
};
PROVIDERS.save(deps.storage, &info.sender, &provider)?;
PROVIDER_BY_DID.save(deps.storage, &provider_did, &info.sender)?;
PROVIDER_COUNT.update(deps.storage, |c| Ok(c + 1))?;
TOTAL_STAKED.update(deps.storage, |t| Ok(t + sent_bond.amount))?;
Ok(Response::new()
.add_attribute("action", "register_provider")
.add_attribute("owner", info.sender.to_string())
.add_attribute("did", provider_did)
.add_attribute("gpu_model", provider.resources.gpu_model)
.add_attribute("gpu_count", provider.resources.gpu_count.to_string())
.add_attribute("bond", sent_bond.to_string()))
}
2.5 Update and Attestation Handlers
pub fn execute_update_resources(deps: DepsMut, env: Env, info: MessageInfo, resources: Resources)
-> Result<Response, ContractError> {
let mut provider = PROVIDERS.load(deps.storage, &info.sender).map_err(|_| ContractError::NotRegistered {})?;
if provider.status == ProviderStatus::Revoked || provider.status == ProviderStatus::Slashed {
return Err(ContractError::ProviderNotActive {});
}
provider.resources = resources;
provider.updated_at = env.block.time;
PROVIDERS.save(deps.storage, &info.sender, &provider)?;
Ok(Response::new().add_attribute("action", "update_resources").add_attribute("owner", info.sender.to_string()))
}
pub fn execute_update_pricing(deps: DepsMut, env: Env, info: MessageInfo, pricing: Pricing)
-> Result<Response, ContractError> {
let mut provider = PROVIDERS.load(deps.storage, &info.sender).map_err(|_| ContractError::NotRegistered {})?;
if provider.status == ProviderStatus::Revoked || provider.status == ProviderStatus::Slashed {
return Err(ContractError::ProviderNotActive {});
}
if pricing.price_per_gpu_hour.is_zero() && pricing.price_per_vcpu_hour.is_zero()
&& pricing.price_per_ram_gb_hour.is_zero() && pricing.price_per_storage_gb_hour.is_zero()
{ return Err(ContractError::InvalidPricing { msg: "at least one price must be non-zero".to_string() }); }
provider.pricing = pricing;
provider.updated_at = env.block.time;
PROVIDERS.save(deps.storage, &info.sender, &provider)?;
Ok(Response::new().add_attribute("action", "update_pricing").add_attribute("owner", info.sender.to_string()))
}
pub fn execute_update_status(deps: DepsMut, env: Env, info: MessageInfo, status: ProviderStatus)
-> Result<Response, ContractError> {
let mut provider = PROVIDERS.load(deps.storage, &info.sender).map_err(|_| ContractError::NotRegistered {})?;
if provider.status == ProviderStatus::Revoked { return Err(ContractError::ProviderNotActive {}); }
if status == ProviderStatus::Slashed || status == ProviderStatus::Revoked {
return Err(ContractError::Unauthorized {});
}
provider.status = status;
provider.updated_at = env.block.time;
PROVIDERS.save(deps.storage, &info.sender, &provider)?;
Ok(Response::new().add_attribute("action", "update_status").add_attribute("status", format!("{:?}", status)))
}
pub fn execute_submit_attestation(deps: DepsMut, env: Env, info: MessageInfo,
gpu_uuids: Vec<String>, cpu_info: String, ram_total_gb: u32, attestation_doc: String)
-> Result<Response, ContractError> {
PROVIDERS.load(deps.storage, &info.sender).map_err(|_| ContractError::NotRegistered {})?;
ATTESTATIONS.save(deps.storage, &info.sender, &HardwareAttestation {
provider: info.sender.clone(), gpu_uuids, cpu_info, ram_total_gb, attestation_doc,
verified: false, verified_at: None,
})?;
Ok(Response::new().add_attribute("action", "submit_attestation").add_attribute("owner", info.sender.to_string()))
}
pub fn execute_withdraw_earnings(deps: DepsMut, _env: Env, info: MessageInfo, amount: Option<Coin>)
-> Result<Response, ContractError> {
let provider = PROVIDERS.load(deps.storage, &info.sender).map_err(|_| ContractError::NotRegistered {})?;
if provider.total_earned.is_zero() { return Err(ContractError::NoEarnings {}); }
let withdraw_amount = match amount {
Some(coin) => { if coin.amount > provider.total_earned {
return Err(ContractError::InsufficientEarnings { available: provider.total_earned, requested: coin.amount });
} coin }
None => Coin { denom: "uaipay".to_string(), amount: provider.total_earned },
};
PROVIDERS.update(deps.storage, &info.sender, |p| -> StdResult<_> {
let mut p = p.unwrap(); p.total_earned = p.total_earned.checked_sub(withdraw_amount.amount)?; Ok(p)
})?;
Ok(Response::new()
.add_message(BankMsg::Send { to_address: info.sender.to_string(), amount: vec![withdraw_amount.clone()] })
.add_attribute("action", "withdraw_earnings")
.add_attribute("amount", withdraw_amount.to_string()))
}
2.6 Revoke and Rate
pub fn execute_revoke(deps: DepsMut, env: Env, info: MessageInfo) -> Result<Response, ContractError> {
let provider = PROVIDERS.load(deps.storage, &info.sender).map_err(|_| ContractError::NotRegistered {})?;
if provider.status == ProviderStatus::Slashed { return Err(ContractError::AlreadySlashed {}); }
let return_bond = if provider.status == ProviderStatus::Slashed {
provider.bond_amount.amount * Uint128::from(90u128) / Uint128::from(100u128)
} else { provider.bond_amount.amount };
PROVIDERS.update(deps.storage, &info.sender, |p| -> StdResult<_> {
let mut p = p.unwrap(); p.status = ProviderStatus::Revoked; p.updated_at = env.block.time; Ok(p)
})?;
PROVIDER_COUNT.update(deps.storage, |c| Ok(c.saturating_sub(1)))?;
TOTAL_STAKED.update(deps.storage, |t| Ok(t.checked_sub(provider.bond_amount.amount)?))?;
Ok(Response::new()
.add_message(BankMsg::Send {
to_address: info.sender.to_string(),
amount: vec![Coin { denom: provider.bond_amount.denom, amount: return_bond }],
})
.add_attribute("action", "revoke").add_attribute("owner", info.sender.to_string()))
}
pub fn execute_rate_provider(deps: DepsMut, env: Env, info: MessageInfo,
provider_addr: String, task_id: String, rating: u8, comment: String)
-> Result<Response, ContractError> {
if rating < 1 || rating > 5 { return Err(ContractError::InvalidRating {}); }
let provider = deps.api.addr_validate(&provider_addr)?;
PROVIDERS.load(deps.storage, &provider).map_err(|_| ContractError::ProviderNotFound {})?;
let record_key = format!("{}/{}", provider_addr, task_id);
REPUTATIONS.save(deps.storage, &record_key, &ReputationRecord {
provider: provider.clone(), client: info.sender, task_id: task_id.clone(),
rating, comment, rated_at: env.block.time,
})?;
let all_ratings: Vec<u32> = REPUTATIONS.range(deps.storage, None, None, Order::Ascending)
.filter_map(|r| { let (_, rec) = r.ok()?; if rec.provider == provider { Some(rec.rating as u32) } else { None } })
.collect();
let avg = if all_ratings.is_empty() { 0 } else {
all_ratings.iter().sum::<u32>() * 100 / all_ratings.len() as u32 * 10
};
PROVIDERS.update(deps.storage, &provider, |p| -> StdResult<_> {
let mut p = p.unwrap(); p.avg_rating = avg; Ok(p)
})?;
Ok(Response::new().add_attribute("action", "rate_provider").add_attribute("rating", rating.to_string()))
}
2.7 Query Implementation
#[entry_point]
pub fn query(deps: Deps, _env: Env, msg: QueryMsg) -> StdResult<Binary> {
match msg {
QueryMsg::GetProvider { address } => {
let addr = deps.api.addr_validate(&address)?;
to_json_binary(&PROVIDERS.load(deps.storage, &addr)?)
}
QueryMsg::GetProviderByDid { did } => {
let addr = PROVIDER_BY_DID.load(deps.storage, &did)?;
to_json_binary(&PROVIDERS.load(deps.storage, &addr)?)
}
QueryMsg::ListProviders { start_after, limit, status_filter } => {
let limit = limit.unwrap_or(30).min(100) as usize;
let start = start_after.as_ref().map(|s| deps.api.addr_validate(s).unwrap());
let providers: Vec<ComputeProvider> = PROVIDERS
.range(deps.storage, start.as_ref(), None, Order::Ascending)
.filter_map(|r| { let (_, p) = r.ok()?;
match &status_filter { Some(f) if p.status != *f => None, _ => Some(p) }
}).take(limit).collect();
to_json_binary(&ProviderListResponse { providers, total_count: PROVIDER_COUNT.load(deps.storage)? })
}
QueryMsg::QueryAvailableProviders { min_gpu_count, min_vcpu, min_ram_gb, min_storage_gb,
region, gpu_model, max_price_per_hour, supports_tee } => {
let providers: Vec<ComputeProvider> = PROVIDERS
.range(deps.storage, None, None, Order::Ascending)
.filter_map(|r| {
let (_, p) = r.ok()?;
if p.status != ProviderStatus::Active { return None; }
if let Some(mg) = min_gpu_count { if p.resources.gpu_count < mg { return None; } }
if let Some(mv) = min_vcpu { if p.resources.vcpu < mv { return None; } }
if let Some(mr) = min_ram_gb { if p.resources.ram_gb < mr { return None; } }
if let Some(ms) = min_storage_gb { if p.resources.storage_gb < ms { return None; } }
if let Some(ref r) = region { if p.region != *r { return None; } }
if let Some(ref m) = gpu_model { if p.resources.gpu_model != *m { return None; } }
if let Some(mp) = max_price_per_hour { if p.pricing.price_per_gpu_hour > mp { return None; } }
if let Some(tee) = supports_tee { if p.resources.supports_tee != tee { return None; } }
Some(p)
}).collect();
to_json_binary(&AvailableProvidersResponse { count: providers.len() as u32, providers })
}
QueryMsg::GetReputation { provider } => {
let addr = deps.api.addr_validate(&provider)?;
let reps: Vec<ReputationRecord> = REPUTATIONS
.range(deps.storage, None, None, Order::Ascending)
.filter_map(|r| { let (_, rec) = r.ok()?;
if rec.provider == addr { Some(rec) } else { None }
}).collect();
to_json_binary(&reps)
}
QueryMsg::GetStats {} => {
let total_providers = PROVIDER_COUNT.load(deps.storage)?;
let total_staked = TOTAL_STAKED.load(deps.storage)?;
let active = PROVIDERS.range(deps.storage, None, None, Order::Ascending)
.filter_map(|r| { let (_, p) = r.ok()?;
if p.status == ProviderStatus::Active { Some(()) } else { None }
}).count() as u64;
let earned = PROVIDERS.range(deps.storage, None, None, Order::Ascending)
.fold(Uint128::zero(), |acc, r| { let (_, p) = r.unwrap(); acc + p.total_earned });
let completed = PROVIDERS.range(deps.storage, None, None, Order::Ascending)
.fold(0u64, |acc, r| { let (_, p) = r.unwrap(); acc + p.completed_tasks });
to_json_binary(&StatsResponse {
total_providers, active_providers: active, total_staked,
total_earned: earned, total_tasks_completed: completed,
})
}
QueryMsg::GetMinBond {} => to_json_binary(&MIN_BOND.load(deps.storage)?),
}
}
2.8 Error Types
// contracts/compute-registry/src/error.rs
#[derive(Error, Debug, PartialEq)]
pub enum ContractError {
#[error("{0}")] Std(#[from] StdError),
#[error("Unauthorized")] Unauthorized {},
#[error("Already registered")] AlreadyRegistered {},
#[error("Not registered")] NotRegistered {},
#[error("DID already exists")] DidAlreadyExists {},
#[error("Provider not found")] ProviderNotFound {},
#[error("Provider not active")] ProviderNotActive {},
#[error("Already slashed")] AlreadySlashed {},
#[error("Missing bond")] MissingBond {},
#[error("Insufficient bond: required {required}, sent {sent}")]
InsufficientBond { required: Uint128, sent: Uint128 },
#[error("Invalid resource: {msg}")] InvalidResource { msg: String },
#[error("Invalid pricing: {msg}")] InvalidPricing { msg: String },
#[error("No earnings")] NoEarnings {},
#[error("Insufficient earnings")]
InsufficientEarnings { available: Uint128, requested: Uint128 },
#[error("Invalid rating")] InvalidRating {},
}
2.9 TypeScript Provider Registration CLI
// scripts/register_provider.ts
import { SigningCosmWasmClient } from "@cosmjs/cosmwasm-stargate";
import { DirectSecp256k1HdWallet } from "@cosmjs/proto-signing";
import { GasPrice, calculateFee } from "@cosmjs/stargate";
const MSG_RPC = "https://rpc.msg-chain-1.msgchain.zone";
interface RegisterParams {
mnemonic: string;
providerDid: string;
gpuModel: string; gpuCount: number; vcpu: number;
ramGb: number; storageGb: number;
pricePerGpuHour: string; pricePerVcpuHour: string;
pricePerRamGbHour: string; pricePerStorageGbHour: string;
region: string; minBillingDuration: number;
registryAddr: string;
}
async function registerProvider(params: RegisterParams) {
const wallet = await DirectSecp256k1HdWallet.fromMnemonic(params.mnemonic, { prefix: "msg" });
const [account] = await wallet.getAccounts();
const client = await SigningCosmWasmClient.connectWithSigner(
MSG_RPC, wallet,
{ gasPrice: GasPrice.fromString("1000000000attoMSG") }
);
const msg = {
register: {
provider_did: params.providerDid,
resources: {
gpu_model: params.gpuModel, gpu_count: params.gpuCount,
vcpu: params.vcpu, ram_gb: params.ramGb,
storage_gb: params.storageGb, bandwidth_mbps: 1000,
supports_tee: false, tee_type: null,
},
pricing: {
price_per_gpu_hour: params.pricePerGpuHour,
price_per_vcpu_hour: params.pricePerVcpuHour,
price_per_ram_gb_hour: params.pricePerRamGbHour,
price_per_storage_gb_hour: params.pricePerStorageGbHour,
min_billing_duration: params.minBillingDuration,
negotiable: false,
},
region: params.region,
},
};
const bondAmount = [{ denom: "uaipay", amount: "1000000000000000000000000" }];
const fee = calculateFee(200000, GasPrice.fromString("1000000000attoMSG"));
const result = await client.execute(account.address, params.registryAddr, msg, fee,
"Register GPU Provider", bondAmount);
console.log(`Provider registered! TX: ${result.transactionHash}`);
console.log(`DID: ${params.providerDid}, GPU: ${params.gpuModel} x${params.gpuCount}`);
}
async function queryProviders(registryAddr: string, gpuModel?: string) {
const wallet = await DirectSecp256k1HdWallet.fromMnemonic(
"dummy mnemonic", { prefix: "msg" }
);
const client = await SigningCosmWasmClient.connectWithSigner(
MSG_RPC, wallet, { gasPrice: GasPrice.fromString("1000000000attoMSG") }
);
const res = await client.queryContractSmart(registryAddr, {
query_available_providers: { gpu_model: gpuModel || null },
});
console.log(`Found ${res.count} providers:`);
res.providers.forEach((p: any) => console.log(`- ${p.provider_did}: ${p.resources.gpu_model} x${p.resources.gpu_count}`));
}
async function updateStatus(mnemonic: string, registryAddr: string, status: string) {
const wallet = await DirectSecp256k1HdWallet.fromMnemonic(mnemonic, { prefix: "msg" });
const [account] = await wallet.getAccounts();
const client = await SigningCosmWasmClient.connectWithSigner(
MSG_RPC, wallet, { gasPrice: GasPrice.fromString("1000000000attoMSG") }
);
const msg = { update_status: { status: status.toUpperCase() } };
const fee = calculateFee(150000, GasPrice.fromString("1000000000attoMSG"));
await client.execute(account.address, registryAddr, msg, fee);
console.log(`Status updated to ${status}`);
}
async function withdrawEarnings(mnemonic: string, registryAddr: string) {
const wallet = await DirectSecp256k1HdWallet.fromMnemonic(mnemonic, { prefix: "msg" });
const [account] = await wallet.getAccounts();
const client = await SigningCosmWasmClient.connectWithSigner(
MSG_RPC, wallet, { gasPrice: GasPrice.fromString("1000000000attoMSG") }
);
const msg = { withdraw_earnings: { amount: null } };
const fee = calculateFee(150000, GasPrice.fromString("1000000000attoMSG"));
const res = await client.execute(account.address, registryAddr, msg, fee);
console.log(`Earnings withdrawn: ${res.transactionHash}`);
}
3. 计算任务管理
3.1 数据模型
// contracts/compute-task/src/state.rs
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub enum TaskType {
Inference, Training, FineTuning,
BatchProcessing, DataPreprocessing, GeneralCompute,
}
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub enum TaskStatus {
Created, Matched, Running, Completed, Settled, Cancelled, Disputed, Failed,
}
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct ComputeTask {
pub id: String,
pub client: Addr,
pub provider: Option<Addr>,
pub task_type: TaskType,
pub container_image: String,
pub resource_requirements: Resources,
pub max_price: Coin,
pub settled_price: Option<Coin>,
pub status: TaskStatus,
pub input_data_cid: String,
pub output_data_cid: Option<String>,
pub expected_duration: u64,
pub actual_duration: Option<u64>,
pub created_at: Timestamp,
pub matched_at: Option<Timestamp>,
pub started_at: Option<Timestamp>,
pub completed_at: Option<Timestamp>,
pub dispute_deadline: Option<Timestamp>,
pub command: Vec<String>,
pub env_vars: Vec<EnvVar>,
pub use_tee: bool,
pub metadata: Option<String>,
}
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct EnvVar { pub key: String, pub value: String, pub secret: bool }
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct TaskEvent {
pub task_id: String, pub from_status: TaskStatus, pub to_status: TaskStatus,
pub actor: Addr, pub timestamp: Timestamp, pub reason: Option<String>,
}
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct Dispute {
pub task_id: String, pub initiator: Addr, pub reason: String,
pub evidence: Vec<String>, pub opened_at: Timestamp,
pub resolved_at: Option<Timestamp>, pub resolution: Option<DisputeResolution>,
}
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub enum DisputeResolution { InFavorOfClient, InFavorOfProvider, SplitPayment { client_pct: u8, provider_pct: u8 } }
pub const TASKS: Map<&str, ComputeTask> = Map::new("tasks");
pub const TASK_EVENTS: Map<&str, Vec<TaskEvent>> = Map::new("task_events");
pub const DISPUTES: Map<&str, Dispute> = Map::new("disputes");
pub const TASK_COUNT: Item<u64> = Item::new("task_count");
pub const CLIENT_TASKS: Map<&Addr, Vec<String>> = Map::new("client_tasks");
pub const PROVIDER_TASKS: Map<&Addr, Vec<String>> = Map::new("provider_tasks");
3.2 Task Messages
// contracts/compute-task/src/msg.rs
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct InstantiateMsg {
pub registry_addr: String, pub aipay_addr: String, pub admin: String,
pub dispute_duration: u64, pub protocol_fee_bps: u64, pub min_task_deposit: Coin,
}
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
#[serde(rename_all = "snake_case")]
pub enum ExecuteMsg {
CreateTask { task_type: TaskType, container_image: String,
resource_requirements: Resources, max_price: Coin, input_data_cid: String,
expected_duration: u64, command: Vec<String>, env_vars: Vec<EnvVar>,
use_tee: bool, metadata: Option<String> },
BidTask { task_id: String, bid_price: Coin },
SelectProvider { task_id: String, provider: String },
StartTask { task_id: String },
CompleteTask { task_id: String, output_data_cid: String, actual_duration: u64 },
ConfirmCompletion { task_id: String },
DisputeTask { task_id: String, reason: String, evidence: Vec<String> },
ResolveDispute { task_id: String, resolution: DisputeResolution },
CancelTask { task_id: String },
ReportFailure { task_id: String, reason: String },
}
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
#[serde(rename_all = "snake_case")]
pub enum QueryMsg {
GetTask { task_id: String },
ListClientTasks { client: String, status_filter: Option<TaskStatus>,
start_after: Option<String>, limit: Option<u32> },
ListProviderTasks { provider: String, status_filter: Option<TaskStatus>,
start_after: Option<String>, limit: Option<u32> },
ListOpenTasks { start_after: Option<String>, limit: Option<u32> },
GetTaskEvents { task_id: String }, GetDispute { task_id: String }, GetStats {},
}
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct TaskListResponse { pub tasks: Vec<ComputeTask>, pub total: u64 }
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct OpenTasksResponse { pub tasks: Vec<ComputeTask>, pub count: u32 }
3.3 Task Contract Instantiate and Execute
// contracts/compute-task/src/contract.rs
const REGISTRY_ADDR: Item<Addr> = Item::new("registry_addr");
const AIPAY_ADDR: Item<Addr> = Item::new("aipay_addr");
const ADMIN: Item<Addr> = Item::new("admin");
const DISPUTE_DURATION: Item<u64> = Item::new("dispute_duration");
const PROTOCOL_FEE_BPS: Item<u64> = Item::new("protocol_fee_bps");
const MIN_TASK_DEPOSIT: Item<Coin> = Item::new("min_task_deposit");
#[entry_point]
pub fn instantiate(deps: DepsMut, _env: Env, _info: MessageInfo, msg: InstantiateMsg) -> StdResult<Response> {
set_contract_version(deps.storage, "compute-task", "1.0.0")?;
REGISTRY_ADDR.save(deps.storage, &deps.api.addr_validate(&msg.registry_addr)?)?;
AIPAY_ADDR.save(deps.storage, &deps.api.addr_validate(&msg.aipay_addr)?)?;
ADMIN.save(deps.storage, &deps.api.addr_validate(&msg.admin)?)?;
DISPUTE_DURATION.save(deps.storage, &msg.dispute_duration)?;
PROTOCOL_FEE_BPS.save(deps.storage, &msg.protocol_fee_bps)?;
MIN_TASK_DEPOSIT.save(deps.storage, &msg.min_task_deposit)?;
TASK_COUNT.save(deps.storage, &0u64)?;
Ok(Response::new().add_attribute("action", "instantiate"))
}
#[entry_point]
pub fn execute(deps: DepsMut, env: Env, info: MessageInfo, msg: ExecuteMsg)
-> Result<Response, ContractError> {
match msg {
ExecuteMsg::CreateTask { task_type, container_image, resource_requirements,
max_price, input_data_cid, expected_duration, command, env_vars, use_tee, metadata } =>
execute_create_task(deps, env, info, task_type, container_image, resource_requirements,
max_price, input_data_cid, expected_duration, command, env_vars, use_tee, metadata),
ExecuteMsg::BidTask { task_id, bid_price } => execute_bid_task(deps, env, info, task_id, bid_price),
ExecuteMsg::SelectProvider { task_id, provider } => execute_select_provider(deps, env, info, task_id, provider),
ExecuteMsg::StartTask { task_id } => execute_start_task(deps, env, info, task_id),
ExecuteMsg::CompleteTask { task_id, output_data_cid, actual_duration } =>
execute_complete_task(deps, env, info, task_id, output_data_cid, actual_duration),
ExecuteMsg::ConfirmCompletion { task_id } => execute_confirm_completion(deps, env, info, task_id),
ExecuteMsg::DisputeTask { task_id, reason, evidence } =>
execute_dispute_task(deps, env, info, task_id, reason, evidence),
ExecuteMsg::ResolveDispute { task_id, resolution } =>
execute_resolve_dispute(deps, env, info, task_id, resolution),
ExecuteMsg::CancelTask { task_id } => execute_cancel_task(deps, env, info, task_id),
ExecuteMsg::ReportFailure { task_id, reason } => execute_report_failure(deps, env, info, task_id, reason),
}
}
fn record_event(deps: &mut DepsMut, task_id: &str, from: TaskStatus, to: TaskStatus,
actor: &Addr, reason: Option<String>) -> StdResult<()> {
let event = TaskEvent { task_id: task_id.to_string(), from_status: from, to_status: to,
actor: actor.clone(), timestamp: deps.api.block_info()?.time, reason };
let mut events = TASK_EVENTS.load(deps.storage, task_id).unwrap_or_default();
events.push(event); TASK_EVENTS.save(deps.storage, task_id, &events)
}
3.4 Create Task
pub fn execute_create_task(deps: DepsMut, env: Env, info: MessageInfo,
task_type: TaskType, container_image: String, resource_requirements: Resources,
max_price: Coin, input_data_cid: String, expected_duration: u64,
command: Vec<String>, env_vars: Vec<EnvVar>, use_tee: bool, metadata: Option<String>)
-> Result<Response, ContractError> {
let min_deposit = MIN_TASK_DEPOSIT.load(deps.storage)?;
let deposit = info.funds.iter().find(|c| c.denom == min_deposit.denom)
.ok_or(ContractError::InsufficientDeposit {})?;
if deposit.amount < min_deposit.amount { return Err(ContractError::InsufficientDeposit {}); }
if container_image.is_empty() { return Err(ContractError::InvalidTaskParam {
msg: "container_image cannot be empty".to_string() }); }
if max_price.amount.is_zero() { return Err(ContractError::InvalidTaskParam {
msg: "max_price must be > 0".to_string() }); }
if expected_duration == 0 { return Err(ContractError::InvalidTaskParam {
msg: "expected_duration must be > 0".to_string() }); }
let task_id = format!("task-{}", Uuid::new_v4());
let task = ComputeTask {
id: task_id.clone(), client: info.sender.clone(), provider: None,
task_type, container_image, resource_requirements, max_price: max_price.clone(),
settled_price: None, status: TaskStatus::Created, input_data_cid,
output_data_cid: None, expected_duration, actual_duration: None,
created_at: env.block.time, matched_at: None, started_at: None,
completed_at: None, dispute_deadline: None, command, env_vars, use_tee, metadata,
};
TASKS.save(deps.storage, &task_id, &task)?;
TASK_COUNT.update(deps.storage, |c| Ok(c + 1))?;
let mut ct = CLIENT_TASKS.load(deps.storage, &info.sender).unwrap_or_default();
ct.push(task_id.clone()); CLIENT_TASKS.save(deps.storage, &info.sender, &ct)?;
record_event(&mut deps, &task_id, TaskStatus::Created, TaskStatus::Created, &info.sender, None)?;
Ok(Response::new().add_attribute("action", "create_task").add_attribute("task_id", task_id)
.add_attribute("type", format!("{:?}", task_type)).add_attribute("max_price", max_price.to_string()))
}
3.5 Bid, Select, Start, Complete
pub fn execute_bid_task(deps: DepsMut, env: Env, info: MessageInfo,
task_id: String, bid_price: Coin) -> Result<Response, ContractError> {
let mut task = TASKS.load(deps.storage, &task_id).map_err(|_| ContractError::TaskNotFound {})?;
if task.status != TaskStatus::Created { return Err(ContractError::InvalidTaskStatus {
expected: "Created".to_string(), actual: format!("{:?}", task.status) }); }
if bid_price.amount > task.max_price.amount { return Err(ContractError::BidExceedsMaxPrice {}); }
if bid_price.amount.is_zero() { return Err(ContractError::InvalidBid {}); }
task.provider = Some(info.sender.clone()); task.settled_price = Some(bid_price);
task.status = TaskStatus::Matched; task.matched_at = Some(env.block.time);
TASKS.save(deps.storage, &task_id, &task)?;
let mut pt = PROVIDER_TASKS.load(deps.storage, &info.sender).unwrap_or_default();
pt.push(task_id.clone()); PROVIDER_TASKS.save(deps.storage, &info.sender, &pt)?;
record_event(&mut deps, &task_id, TaskStatus::Created, TaskStatus::Matched,
&info.sender, Some("Provider bid accepted".to_string()))?;
Ok(Response::new().add_attribute("action", "bid_task").add_attribute("task_id", task_id))
}
pub fn execute_start_task(deps: DepsMut, env: Env, info: MessageInfo, task_id: String)
-> Result<Response, ContractError> {
let mut task = TASKS.load(deps.storage, &task_id).map_err(|_| ContractError::TaskNotFound {})?;
if task.provider.as_ref() != Some(&info.sender) { return Err(ContractError::Unauthorized {}); }
if task.status != TaskStatus::Matched { return Err(ContractError::InvalidTaskStatus {
expected: "Matched".to_string(), actual: format!("{:?}", task.status) }); }
task.status = TaskStatus::Running; task.started_at = Some(env.block.time);
TASKS.save(deps.storage, &task_id, &task)?;
record_event(&mut deps, &task_id, TaskStatus::Matched, TaskStatus::Running, &info.sender, None)?;
Ok(Response::new().add_attribute("action", "start_task").add_attribute("task_id", task_id))
}
pub fn execute_complete_task(deps: DepsMut, env: Env, info: MessageInfo,
task_id: String, output_data_cid: String, actual_duration: u64)
-> Result<Response, ContractError> {
let mut task = TASKS.load(deps.storage, &task_id).map_err(|_| ContractError::TaskNotFound {})?;
if task.provider.as_ref() != Some(&info.sender) { return Err(ContractError::Unauthorized {}); }
if task.status != TaskStatus::Running { return Err(ContractError::InvalidTaskStatus {
expected: "Running".to_string(), actual: format!("{:?}", task.status) }); }
task.output_data_cid = Some(output_data_cid.clone());
task.actual_duration = Some(actual_duration);
task.status = TaskStatus::Completed; task.completed_at = Some(env.block.time);
let dur = DISPUTE_DURATION.load(deps.storage)?;
task.dispute_deadline = Some(env.block.time.plus_seconds(dur));
TASKS.save(deps.storage, &task_id, &task)?;
record_event(&mut deps, &task_id, TaskStatus::Running, TaskStatus::Completed, &info.sender, None)?;
Ok(Response::new().add_attribute("action", "complete_task")
.add_attribute("task_id", task_id).add_attribute("duration", actual_duration.to_string()))
}
3.6 Confirm, Cancel, Failure
pub fn execute_confirm_completion(deps: DepsMut, env: Env, info: MessageInfo, task_id: String)
-> Result<Response, ContractError> {
let mut task = TASKS.load(deps.storage, &task_id).map_err(|_| ContractError::TaskNotFound {})?;
if task.client != info.sender { return Err(ContractError::Unauthorized {}); }
if task.status != TaskStatus::Completed { return Err(ContractError::InvalidTaskStatus {
expected: "Completed".to_string(), actual: format!("{:?}", task.status) }); }
task.status = TaskStatus::Settled;
let settled = task.settled_price.clone().unwrap();
let fee_bps = PROTOCOL_FEE_BPS.load(deps.storage)?;
let protocol_fee = settled.amount * Uint128::from(fee_bps as u128) / Uint128::from(10000u128);
TASKS.save(deps.storage, &task_id, &task)?;
record_event(&mut deps, &task_id, TaskStatus::Completed, TaskStatus::Settled, &info.sender, None)?;
Ok(Response::new().add_attribute("action", "confirm_completion")
.add_attribute("task_id", task_id).add_attribute("fee", protocol_fee.to_string()))
}
pub fn execute_cancel_task(deps: DepsMut, env: Env, info: MessageInfo, task_id: String)
-> Result<Response, ContractError> {
let mut task = TASKS.load(deps.storage, &task_id).map_err(|_| ContractError::TaskNotFound {})?;
if task.client != info.sender { return Err(ContractError::Unauthorized {}); }
if task.status != TaskStatus::Created { return Err(ContractError::InvalidTaskStatus {
expected: "Created".to_string(), actual: format!("{:?}", task.status) }); }
task.status = TaskStatus::Cancelled; TASKS.save(deps.storage, &task_id, &task)?;
record_event(&mut deps, &task_id, TaskStatus::Created, TaskStatus::Cancelled, &info.sender, None)?;
Ok(Response::new().add_attribute("action", "cancel_task").add_attribute("task_id", task_id))
}
pub fn execute_report_failure(deps: DepsMut, env: Env, info: MessageInfo, task_id: String, reason: String)
-> Result<Response, ContractError> {
let mut task = TASKS.load(deps.storage, &task_id).map_err(|_| ContractError::TaskNotFound {})?;
if task.provider.as_ref() != Some(&info.sender) { return Err(ContractError::Unauthorized {}); }
if task.status != TaskStatus::Running { return Err(ContractError::InvalidTaskStatus {
expected: "Running".to_string(), actual: format!("{:?}", task.status) }); }
task.status = TaskStatus::Failed; task.completed_at = Some(env.block.time);
TASKS.save(deps.storage, &task_id, &task)?;
record_event(&mut deps, &task_id, TaskStatus::Running, TaskStatus::Failed, &info.sender, Some(reason))?;
Ok(Response::new().add_attribute("action", "report_failure").add_attribute("task_id", task_id))
}
3.7 Dispute Handlers
pub fn execute_dispute_task(deps: DepsMut, env: Env, info: MessageInfo,
task_id: String, reason: String, evidence: Vec<String>) -> Result<Response, ContractError> {
let task = TASKS.load(deps.storage, &task_id).map_err(|_| ContractError::TaskNotFound {})?;
if task.client != info.sender && task.provider.as_ref() != Some(&info.sender) {
return Err(ContractError::Unauthorized {});
}
if task.status != TaskStatus::Completed && task.status != TaskStatus::Running {
return Err(ContractError::InvalidTaskStatus { expected: "Running/Completed".to_string(),
actual: format!("{:?}", task.status) });
}
if let Some(deadline) = task.dispute_deadline {
if env.block.time > deadline { return Err(ContractError::DisputePeriodExpired {}); }
}
if DISPUTES.may_load(deps.storage, &task_id)?.is_some() {
return Err(ContractError::DisputeAlreadyExists {});
}
let dispute = Dispute { task_id: task_id.clone(), initiator: info.sender.clone(),
reason, evidence, opened_at: env.block.time, resolved_at: None, resolution: None };
DISPUTES.save(deps.storage, &task_id, &dispute)?;
TASKS.update(deps.storage, &task_id, |t| -> StdResult<_> {
let mut t = t.unwrap(); t.status = TaskStatus::Disputed; Ok(t)
})?;
record_event(&mut deps, &task_id, task.status, TaskStatus::Disputed,
&info.sender, Some("Dispute initiated".to_string()))?;
Ok(Response::new().add_attribute("action", "dispute_task").add_attribute("task_id", task_id))
}
pub fn execute_resolve_dispute(deps: DepsMut, env: Env, info: MessageInfo,
task_id: String, resolution: DisputeResolution) -> Result<Response, ContractError> {
let admin = ADMIN.load(deps.storage)?;
if info.sender != admin { return Err(ContractError::Unauthorized {}); }
let mut dispute = DISPUTES.load(deps.storage, &task_id).map_err(|_| ContractError::DisputeNotFound {})?;
dispute.resolution = Some(resolution.clone());
dispute.resolved_at = Some(env.block.time);
DISPUTES.save(deps.storage, &task_id, &dispute)?;
let current_status = TASKS.load(deps.storage, &task_id)?.status;
match resolution {
DisputeResolution::InFavorOfClient => {
TASKS.update(deps.storage, &task_id, |t| -> StdResult<_> {
let mut t = t.unwrap(); t.status = TaskStatus::Failed; Ok(t)
})?;
}
DisputeResolution::InFavorOfProvider => {
TASKS.update(deps.storage, &task_id, |t| -> StdResult<_> {
let mut t = t.unwrap(); t.status = TaskStatus::Settled; Ok(t)
})?;
}
DisputeResolution::SplitPayment { .. } => {
TASKS.update(deps.storage, &task_id, |t| -> StdResult<_> {
let mut t = t.unwrap(); t.status = TaskStatus::Settled; Ok(t)
})?;
}
}
let new_status = TASKS.load(deps.storage, &task_id)?.status;
record_event(&mut deps, &task_id, current_status, new_status, &info.sender, None)?;
Ok(Response::new().add_attribute("action", "resolve_dispute")
.add_attribute("task_id", task_id))
}
3.8 Task Query Implementation
#[entry_point]
pub fn query(deps: Deps, _env: Env, msg: QueryMsg) -> StdResult<Binary> {
match msg {
QueryMsg::GetTask { task_id } => to_json_binary(&TASKS.load(deps.storage, &task_id)?),
QueryMsg::ListClientTasks { client, status_filter, start_after, limit } => {
let addr = deps.api.addr_validate(&client)?;
let ids = CLIENT_TASKS.load(deps.storage, &addr).unwrap_or_default();
let lim = limit.unwrap_or(30).min(100) as usize;
let start = start_after.and_then(|s| ids.iter().position(|id| id == &s)).map(|i| i + 1).unwrap_or(0);
let tasks: Vec<ComputeTask> = ids.iter().skip(start).filter_map(|id| {
let t = TASKS.load(deps.storage, id).ok()?;
match &status_filter { Some(f) if t.status != *f => None, _ => Some(t) }
}).take(lim).collect();
to_json_binary(&TaskListResponse { total: ids.len() as u64, tasks })
}
QueryMsg::ListProviderTasks { provider, status_filter, start_after, limit } => {
let addr = deps.api.addr_validate(&provider)?;
let ids = PROVIDER_TASKS.load(deps.storage, &addr).unwrap_or_default();
let lim = limit.unwrap_or(30).min(100) as usize;
let start = start_after.and_then(|s| ids.iter().position(|id| id == &s)).map(|i| i + 1).unwrap_or(0);
let tasks: Vec<ComputeTask> = ids.iter().skip(start).filter_map(|id| {
let t = TASKS.load(deps.storage, id).ok()?;
match &status_filter { Some(f) if t.status != *f => None, _ => Some(t) }
}).take(lim).collect();
to_json_binary(&TaskListResponse { total: ids.len() as u64, tasks })
}
QueryMsg::ListOpenTasks { start_after: _, limit } => {
let lim = limit.unwrap_or(30).min(100) as usize;
let tasks: Vec<ComputeTask> = TASKS.range(deps.storage, None, None, Order::Ascending)
.filter_map(|r| { let (_, t) = r.ok()?; if t.status == TaskStatus::Created { Some(t) } else { None } })
.take(lim).collect();
to_json_binary(&OpenTasksResponse { count: tasks.len() as u32, tasks })
}
QueryMsg::GetTaskEvents { task_id } => {
to_json_binary(&TASK_EVENTS.load(deps.storage, &task_id).unwrap_or_default())
}
QueryMsg::GetDispute { task_id } => to_json_binary(&DISPUTES.load(deps.storage, &task_id)?),
QueryMsg::GetStats {} => {
let total = TASK_COUNT.load(deps.storage)?;
let open = TASKS.range(deps.storage, None, None, Order::Ascending)
.filter_map(|r| { let (_, t) = r.ok()?; if t.status == TaskStatus::Created { Some(()) } else { None } }).count() as u64;
let running = TASKS.range(deps.storage, None, None, Order::Ascending)
.filter_map(|r| { let (_, t) = r.ok()?; if t.status == TaskStatus::Running { Some(()) } else { None } }).count() as u64;
let completed = TASKS.range(deps.storage, None, None, Order::Ascending)
.filter_map(|r| { let (_, t) = r.ok()?;
if t.status == TaskStatus::Completed || t.status == TaskStatus::Settled { Some(()) } else { None }
}).count() as u64;
to_json_binary(&TaskStatsResponse { total_tasks: total, open_tasks: open, running_tasks: running, completed_tasks: completed, total_volume: Uint128::zero(), avg_completion_time: 0 })
}
}
}
3.9 Task Error Types
// contracts/compute-task/src/error.rs
#[derive(Error, Debug, PartialEq)]
pub enum ContractError {
#[error("{0}")] Std(#[from] StdError),
#[error("Unauthorized")] Unauthorized {},
#[error("Task not found")] TaskNotFound {},
#[error("Invalid status: expected {expected}, actual {actual}")]
InvalidTaskStatus { expected: String, actual: String },
#[error("Invalid param: {msg}")] InvalidTaskParam { msg: String },
#[error("Insufficient deposit")] InsufficientDeposit {},
#[error("Bid exceeds max price")] BidExceedsMaxPrice {},
#[error("Invalid bid")] InvalidBid {},
#[error("Dispute already exists")] DisputeAlreadyExists {},
#[error("Dispute not found")] DisputeNotFound {},
#[error("Dispute period expired")] DisputePeriodExpired {},
}
3.10 Task Lifecycle
Created --> Matched --> Running --> Completed --> Settled
| | | |
+-> Cancelled +-> Failed +-> Disputed
+-> Resolved
4. 任务匹配引擎
4.1 匹配模型
匹配引擎的核心目标是将客户端的计算需求与提供者的资源供给高效配对。MSG Chain 支持三种匹配模式:
| 模式 | 描述 | 适用场景 |
|---|---|---|
| 订单簿 | 客户发布需求(Ask),提供者发布报价(Bid),引擎撮合 | 持续计算需求 |
| 反向拍卖 | 客户发布任务,多个提供者竞标 | 一次性训练任务 |
| 直接匹配 | 客户直接指定提供者 | 已有合作关系的双方 |
4.2 匹配算法实现 (Python)
# matcher/compute_matcher.py
import hashlib, json, time
from dataclasses import dataclass, field
from enum import Enum
from typing import Dict, List, Optional, Tuple
from decimal import Decimal, ROUND_HALF_UP
class TaskType(Enum):
INFERENCE = "inference"; TRAINING = "training"
FINE_TUNING = "fine_tuning"; BATCH_PROCESSING = "batch_processing"
DATA_PREPROCESSING = "data_preprocessing"; GENERAL_COMPUTE = "general_compute"
class OrderSide(Enum): ASK = "ask"; BID = "bid"
class OrderStatus(Enum): OPEN = "open"; PARTIAL = "partial"; FILLED = "filled"; CANCELLED = "cancelled"; EXPIRED = "expired"
@dataclass
class Resources:
gpu_model: str; gpu_count: int; vcpu: int; ram_gb: int; storage_gb: int
bandwidth_mbps: int = 1000; supports_tee: bool = False; tee_type: Optional[str] = None
def meets_requirements(self, required: 'Resources') -> bool:
if self.gpu_count < required.gpu_count: return False
if self.vcpu < required.vcpu: return False
if self.ram_gb < required.ram_gb: return False
if self.storage_gb < required.storage_gb: return False
if required.supports_tee and not self.supports_tee: return False
if required.gpu_model and self.gpu_model != required.gpu_model:
if required.gpu_model != "any": return False
return True
def score_match(self, required: 'Resources') -> float:
if not self.meets_requirements(required): return 0.0
score = 100.0
if required.gpu_count > 0: score *= min(1.0, required.gpu_count / self.gpu_count * 1.2)
if required.vcpu > 0: score *= min(1.0, required.vcpu / self.vcpu * 1.2)
if required.gpu_model and required.gpu_model != "any":
score *= 1.1 if self.gpu_model == required.gpu_model else 0.8
return round(score, 2)
@dataclass
class Pricing:
price_per_gpu_hour: Decimal; price_per_vcpu_hour: Decimal
price_per_ram_gb_hour: Decimal; price_per_storage_gb_hour: Decimal
min_billing_duration: int = 300; negotiable: bool = False
def estimate_cost(self, duration_seconds: int, resources: Resources) -> Decimal:
hours = Decimal(duration_seconds) / Decimal(3600)
min_hours = Decimal(self.min_billing_duration) / Decimal(3600)
billable = max(hours, min_hours)
return (self.price_per_gpu_hour * Decimal(resources.gpu_count) * billable
+ self.price_per_vcpu_hour * Decimal(resources.vcpu) * billable
+ self.price_per_ram_gb_hour * Decimal(resources.ram_gb) * billable
+ self.price_per_storage_gb_hour * Decimal(resources.storage_gb) * billable)
@dataclass
class ProviderOrder:
order_id: str; provider_did: str; provider_addr: str
resources: Resources; pricing: Pricing; region: str
rating: float; completed_tasks: int; status: 'ProviderStatus'
available_from: int; available_until: int
created_at: int = field(default_factory=time.time)
@dataclass
class ClientOrder:
order_id: str; client_addr: str; task_type: TaskType
required_resources: Resources; max_budget: Decimal; expected_duration: int
region_preference: Optional[str] = None
min_rating: Optional[float] = None; prefer_tee: bool = False
created_at: int = field(default_factory=time.time)
@dataclass
class Match:
match_id: str; client_order_id: str; provider_order_id: str
provider_addr: str; client_addr: str; price: Decimal
estimated_cost: Decimal; duration: int; score: float
matched_at: int = field(default_factory=time.time)
4.3 Core Matching Logic
class ComputeMatcher:
def __init__(self):
self.ask_orders: Dict[str, ClientOrder] = {}
self.bid_orders: Dict[str, ProviderOrder] = {}
self.matches: List[Match] = []
self.orderbooks: Dict[str, List[Tuple[Decimal, str]]] = {}
def place_ask(self, order: ClientOrder) -> str:
order.order_id = self._generate_id("ask")
self.ask_orders[order.order_id] = order
self.match_asks([order])
return order.order_id
def place_bid(self, order: ProviderOrder) -> str:
order.order_id = self._generate_id("bid")
self.bid_orders[order.order_id] = order
key = f"{order.resources.gpu_model}:{order.region}"
if key not in self.orderbooks: self.orderbooks[key] = []
self.orderbooks[key].append((order.pricing.price_per_gpu_hour, order.order_id))
self.orderbooks[key].sort(key=lambda x: x[0])
return order.order_id
def match_asks(self, asks: List[ClientOrder]) -> List[Match]:
matches = []
for ask in asks:
candidates = self._find_candidates(ask)
if not candidates: continue
candidates.sort(key=lambda c: c[1], reverse=True)
best, score, cost = candidates[0]
m = Match(match_id=self._generate_id("match"),
client_order_id=ask.order_id, provider_order_id=best.order_id,
provider_addr=best.provider_addr, client_addr=ask.client_addr,
price=best.pricing.price_per_gpu_hour, estimated_cost=cost,
duration=ask.expected_duration, score=score)
matches.append(m); self.matches.append(m)
return matches
def _find_candidates(self, ask: ClientOrder) -> List[Tuple[ProviderOrder, float, Decimal]]:
candidates = []
for o in self.bid_orders.values():
if o.status != ProviderStatus.Active: continue
if not o.resources.meets_requirements(ask.required_resources): continue
if ask.region_preference and o.region != ask.region_preference: continue
if ask.min_rating and o.rating < ask.min_rating: continue
if ask.prefer_tee and not o.resources.supports_tee: continue
cost = o.pricing.estimate_cost(ask.expected_duration, ask.required_resources)
max_allowed = ask.max_budget * Decimal(ask.expected_duration) / Decimal(3600)
if cost > max_allowed: continue
rs = o.resources.score_match(ask.required_resources)
ps = self._price_score(o.pricing.price_per_gpu_hour, ask.max_budget)
total = rs * 0.35 + ps * 0.35 + o.rating * 2.0 + min(o.completed_tasks, 100) * 0.10
candidates.append((o, round(total, 2), cost))
return candidates
def _price_score(self, p: Decimal, max_b: Decimal) -> float:
if max_b <= Decimal(0): return 50.0
ratio = float(p / max_b)
if ratio <= 0.5: return 100.0
if ratio <= 0.8: return 80.0
if ratio <= 1.0: return 60.0
return max(0, 50.0 / ratio)
def _generate_id(self, prefix: str) -> str:
ts = int(time.time() * 1000000)
return f"{prefix}:{ts}:{hashlib.sha256(str(ts).encode()).hexdigest()[:8]}"
class ProviderStatus(Enum):
Active = "active"; Inactive = "inactive"; Paused = "paused"
Slashed = "slashed"; Revoked = "revoked"
4.4 Reverse Auction Engine
# matcher/reverse_auction.py
import time, hashlib, json
from dataclasses import dataclass
from decimal import Decimal
from typing import Dict, List, Optional
from enum import Enum
class AuctionStatus(Enum): OPEN = "open"; CLOSED = "closed"; AWARDED = "awarded"; CANCELLED = "cancelled"
@dataclass
class Bid:
bidder: str; price: Decimal; estimated_time: int
submitted_at: int; signature: str
@dataclass
class Auction:
task_id: str; client_addr: str; required_resources: 'Resources'
max_budget: Decimal; expected_duration: int; status: AuctionStatus
bids: List[Bid] = None; winner: Optional[str] = None
winning_price: Optional[Decimal] = None
opened_at: int = 0; closes_at: int = 0
def __post_init__(self):
if self.bids is None: self.bids = []
class ReverseAuctionEngine:
def __init__(self):
self.auctions: Dict[str, Auction] = {}
self.min_bid_time = 300; self.bid_extension = 120
def create_auction(self, task_id: str, client_addr: str,
required_resources: 'Resources', max_budget: Decimal,
expected_duration: int, bidding_duration: int = 3600) -> Auction:
a = Auction(task_id=task_id, client_addr=client_addr,
required_resources=required_resources, max_budget=max_budget,
expected_duration=expected_duration, status=AuctionStatus.OPEN,
opened_at=int(time.time()),
closes_at=int(time.time()) + max(bidding_duration, self.min_bid_time))
self.auctions[task_id] = a; return a
def submit_bid(self, task_id: str, bidder: str, price: Decimal,
estimated_time: int, signature: str) -> bool:
a = self.auctions.get(task_id)
if not a: raise ValueError(f"Auction {task_id} not found")
if a.status != AuctionStatus.OPEN: raise ValueError("Auction not open")
if price > a.max_budget: raise ValueError("Bid exceeds max budget")
if int(time.time()) > a.closes_at: raise ValueError("Auction closed")
if estimated_time > a.expected_duration * 2: raise ValueError("Time too long")
a.bids.append(Bid(bidder=bidder, price=price, estimated_time=estimated_time,
submitted_at=int(time.time()), signature=signature))
if a.closes_at - int(time.time()) < 300: a.closes_at += self.bid_extension
return True
def close_auction(self, task_id: str) -> Optional[Bid]:
a = self.auctions.get(task_id)
if not a: raise ValueError(f"Auction {task_id} not found")
if a.status != AuctionStatus.OPEN: raise ValueError("Not open")
a.status = AuctionStatus.CLOSED
if not a.bids: a.status = AuctionStatus.CANCELLED; return None
scored = [(self._normalize_price(b.price, a.max_budget) * 0.6
+ self._normalize_time(b.estimated_time, a.expected_duration) * 0.4, b)
for b in a.bids]
scored.sort(key=lambda x: x[0], reverse=True)
winner = scored[0][1]
a.winner = winner.bidder; a.winning_price = winner.price
a.status = AuctionStatus.AWARDED; return winner
def _normalize_price(self, price: Decimal, max_b: Decimal) -> float:
if max_b <= Decimal(0): return 0
return max(0, (1 - float(price / max_b)) * 100)
def _normalize_time(self, estimated: int, expected: int) -> float:
if expected <= 0: return 0
ratio = estimated / expected
if ratio <= 0.5: return 100
if ratio <= 1.0: return 80
if ratio <= 1.5: return 50
return max(0, 30 / ratio)
4.5 Chain Integration
# matcher/chain_integration.py
import json, requests
from typing import Dict, List, Optional
class ChainClient:
def __init__(self, rpc: str, registry: str, task_contract: str):
self.lcd = rpc.replace("rpc", "lcd")
self.registry = registry; self.task_contract = task_contract
def get_open_tasks(self, limit: int = 100) -> List[Dict]:
url = f"{self.lcd}/cosmwasm/wasm/v1/contract/{self.task_contract}/smart"
resp = requests.get(url, params={"msg": json.dumps({"list_open_tasks": {"limit": limit}})})
return resp.json().get("data", {}).get("tasks", [])
def get_active_providers(self, gpu: Optional[str] = None) -> List[Dict]:
url = f"{self.lcd}/cosmwasm/wasm/v1/contract/{self.registry}/smart"
q = {"query_available_providers": {"gpu_model": gpu}}
resp = requests.get(url, params={"msg": json.dumps(q)})
return resp.json().get("data", {}).get("providers", [])
def submit_match(self, task_id: str, provider: str, price: str) -> Dict:
tx = {"select_provider": {"task_id": task_id, "provider": provider}}
print(f"Match: task={task_id}, provider={provider}, price={price}")
return {"status": "submitted"}
def sync_and_match(self) -> List[Dict]:
tasks = self.get_open_tasks()
providers = self.get_active_providers()
# Run matching algorithm on chain data
print(f"Synced {len(tasks)} tasks, {len(providers)} providers")
return []
5. 支付与结算
5.1 AIPAY 集成
支付系统基于 MSG Chain 的 AIPAY 模块,提供托管 (escrow)、释放 (release)、退款 (refund) 能力。
// contracts/compute-payment/src/state.rs
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct PaymentAccount {
pub task_id: String, pub escrow: Addr, pub payer: Addr, pub payee: Addr,
pub amount: Coin, pub released_amount: Coin, pub status: PaymentStatus,
pub created_at: Timestamp, pub dispute_deadline: Option<Timestamp>,
}
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub enum PaymentStatus { Escrowed, Released, Refunded, Disputed, PartialReleased }
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct FeeConfig { pub protocol_fee_bps: u64, pub dispute_fee: Coin, pub min_fee: Coin }
pub const PAYMENTS: Map<&str, PaymentAccount> = Map::new("payments");
pub const FEE_CONFIG: Item<FeeConfig> = Item::new("fee_config");
pub const TOTAL_FEES_COLLECTED: Item<Coin> = Item::new("total_fees");
5.2 Payment Messages
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct InstantiateMsg { pub protocol_fee_bps: u64, pub dispute_fee: Coin, pub min_fee: Coin }
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
#[serde(rename_all = "snake_case")]
pub enum ExecuteMsg {
Escrow { task_id: String, payee: String, amount: Coin, dispute_duration: u64 },
Release { task_id: String },
PartialRelease { task_id: String, amount: Coin },
Refund { task_id: String },
ReleaseToProvider { task_id: String, provider_pct: u8, client_pct: u8 },
CollectFees {},
}
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
#[serde(rename_all = "snake_case")]
pub enum QueryMsg {
GetPayment { task_id: String }, ListPaymentsByPayer { payer: String },
ListPaymentsByPayee { payee: String }, GetFeeConfig {}, GetTotalFees {},
}
5.3 Payment Contract Implementation
#[entry_point]
pub fn instantiate(deps: DepsMut, _env: Env, _info: MessageInfo, msg: InstantiateMsg) -> StdResult<Response> {
FEE_CONFIG.save(deps.storage, &FeeConfig {
protocol_fee_bps: msg.protocol_fee_bps, dispute_fee: msg.dispute_fee, min_fee: msg.min_fee })?;
TOTAL_FEES_COLLECTED.save(deps.storage, &Coin { denom: "uaipay".to_string(), amount: Uint128::zero() })?;
Ok(Response::new().add_attribute("action", "instantiate"))
}
#[entry_point]
pub fn execute(deps: DepsMut, env: Env, info: MessageInfo, msg: ExecuteMsg) -> Result<Response, ContractError> {
match msg {
ExecuteMsg::Escrow { task_id, payee, amount, dispute_duration } =>
execute_escrow(deps, env, info, task_id, payee, amount, dispute_duration),
ExecuteMsg::Release { task_id } => execute_release(deps, env, info, task_id),
ExecuteMsg::PartialRelease { task_id, amount } => execute_partial_release(deps, env, info, task_id, amount),
ExecuteMsg::Refund { task_id } => execute_refund(deps, env, info, task_id),
ExecuteMsg::ReleaseToProvider { task_id, provider_pct, client_pct } =>
execute_release_split(deps, env, info, task_id, provider_pct, client_pct),
ExecuteMsg::CollectFees {} => execute_collect_fees(deps, env, info),
}
}
pub fn execute_escrow(deps: DepsMut, env: Env, info: MessageInfo,
task_id: String, payee: String, amount: Coin, dispute_duration: u64)
-> Result<Response, ContractError> {
if PAYMENTS.may_load(deps.storage, &task_id)?.is_some() {
return Err(ContractError::PaymentAlreadyExists {});
}
let sent = info.funds.iter()
.find(|c| c.denom == amount.denom)
.ok_or(ContractError::NoFundsSent {})?;
if sent.amount < amount.amount {
return Err(ContractError::InsufficientFunds { required: amount.amount, sent: sent.amount });
}
let pay = PaymentAccount {
task_id: task_id.clone(), escrow: info.sender.clone(),
payer: info.sender.clone(), payee: deps.api.addr_validate(&payee)?,
amount: amount.clone(),
released_amount: Coin { denom: amount.denom.clone(), amount: Uint128::zero() },
status: PaymentStatus::Escrowed, created_at: env.block.time,
dispute_deadline: Some(env.block.time.plus_seconds(dispute_duration)),
};
PAYMENTS.save(deps.storage, &task_id, &pay)?;
Ok(Response::new().add_attribute("action", "escrow").add_attribute("task_id", task_id))
}
pub fn execute_release(deps: DepsMut, _env: Env, info: MessageInfo, task_id: String)
-> Result<Response, ContractError> {
let pay = PAYMENTS.load(deps.storage, &task_id).map_err(|_| ContractError::PaymentNotFound {})?;
if pay.escrow != info.sender { return Err(ContractError::Unauthorized {}); }
if pay.status != PaymentStatus::Escrowed {
return Err(ContractError::InvalidPaymentStatus { expected: "Escrowed".to_string(), actual: format!("{:?}", pay.status) });
}
let cfg = FEE_CONFIG.load(deps.storage)?;
let fee = pay.amount.amount * Uint128::from(cfg.protocol_fee_bps as u128) / Uint128::from(10000u128);
let to_provider = pay.amount.amount.checked_sub(fee)?;
PAYMENTS.update(deps.storage, &task_id, |x| -> StdResult<_> {
let mut x = x.unwrap(); x.status = PaymentStatus::Released;
x.released_amount = Coin { denom: pay.amount.denom.clone(), amount: to_provider }; Ok(x)
})?;
TOTAL_FEES_COLLECTED.update(deps.storage, |t| -> StdResult<_> {
Ok(Coin { denom: t.denom, amount: t.amount.checked_add(fee)? })
})?;
let mut msgs = vec![BankMsg::Send { to_address: pay.payee.to_string(),
amount: vec![Coin { denom: pay.amount.denom.clone(), amount: to_provider }] }];
if !fee.is_zero() {
msgs.push(BankMsg::Send { to_address: info.sender.to_string(),
amount: vec![Coin { denom: pay.amount.denom, amount: fee }] });
}
Ok(Response::new().add_messages(msgs).add_attribute("action", "release"))
}
pub fn execute_refund(deps: DepsMut, _env: Env, info: MessageInfo, task_id: String)
-> Result<Response, ContractError> {
let pay = PAYMENTS.load(deps.storage, &task_id).map_err(|_| ContractError::PaymentNotFound {})?;
if pay.escrow != info.sender { return Err(ContractError::Unauthorized {}); }
if pay.status != PaymentStatus::Escrowed {
return Err(ContractError::InvalidPaymentStatus { expected: "Escrowed".to_string(), actual: format!("{:?}", pay.status) });
}
PAYMENTS.update(deps.storage, &task_id, |x| -> StdResult<_> { let mut x = x.unwrap(); x.status = PaymentStatus::Refunded; Ok(x) })?;
Ok(Response::new().add_message(BankMsg::Send { to_address: pay.payer.to_string(), amount: vec![pay.amount] })
.add_attribute("action", "refund"))
}
pub fn execute_collect_fees(deps: DepsMut, _env: Env, info: MessageInfo) -> Result<Response, ContractError> {
let fees = TOTAL_FEES_COLLECTED.load(deps.storage)?;
if fees.amount.is_zero() { return Err(ContractError::NoFeesToCollect {}); }
TOTAL_FEES_COLLECTED.save(deps.storage, &Coin { denom: fees.denom.clone(), amount: Uint128::zero() })?;
Ok(Response::new().add_message(BankMsg::Send { to_address: info.sender.to_string(), amount: vec![fees] })
.add_attribute("action", "collect_fees"))
}
5.4 TypeScript Payment Client
// src/payment/paymentClient.ts
import { SigningCosmWasmClient } from "@cosmjs/cosmwasm-stargate";
import { GasPrice, calculateFee } from "@cosmjs/stargate";
import { DirectSecp256k1HdWallet } from "@cosmjs/proto-signing";
export interface EscrowParams {
mnemonic: string; taskId: string; payee: string;
amount: string; denom: string; disputeDuration: number;
}
export class PaymentClient {
constructor(private config: { rpcEndpoint: string; paymentContractAddr: string; gasPrice: string }) {}
private async getWallet(mnemonic: string) {
const w = await DirectSecp256k1HdWallet.fromMnemonic(mnemonic, { prefix: "msg" });
const [a] = await w.getAccounts(); return { wallet: w, address: a.address };
}
async escrow(p: EscrowParams): Promise<string> {
const { wallet, address } = await this.getWallet(p.mnemonic);
const client = await SigningCosmWasmClient.connectWithSigner(this.config.rpcEndpoint, wallet,
{ gasPrice: GasPrice.fromString(this.config.gasPrice) });
const result = await client.execute(address, this.config.paymentContractAddr,
{ escrow: { task_id: p.taskId, payee: p.payee, amount: { denom: p.denom, amount: p.amount }, dispute_duration: p.disputeDuration } },
calculateFee(300000, GasPrice.fromString(this.config.gasPrice)),
"", [{ denom: p.denom, amount: p.amount }]);
return result.transactionHash;
}
async release(mnemonic: string, taskId: string): Promise<string> {
const { wallet, address } = await this.getWallet(mnemonic);
const client = await SigningCosmWasmClient.connectWithSigner(this.config.rpcEndpoint, wallet,
{ gasPrice: GasPrice.fromString(this.config.gasPrice) });
const result = await client.execute(address, this.config.paymentContractAddr,
{ release: { task_id: taskId } },
calculateFee(200000, GasPrice.fromString(this.config.gasPrice)));
return result.transactionHash;
}
async refund(mnemonic: string, taskId: string): Promise<string> {
const { wallet, address } = await this.getWallet(mnemonic);
const client = await SigningCosmWasmClient.connectWithSigner(this.config.rpcEndpoint, wallet,
{ gasPrice: GasPrice.fromString(this.config.gasPrice) });
const result = await client.execute(address, this.config.paymentContractAddr,
{ refund: { task_id: taskId } },
calculateFee(200000, GasPrice.fromString(this.config.gasPrice)));
return result.transactionHash;
}
}
6. 容器与执行环境
6.1 Container Specification
提供者需运行一个执行环境(Execution Environment),能够根据任务规范拉取容器镜像、挂载输入、执行计算、返回输出。
# spec/container_spec.yaml
version: "1.0"
task:
id: "task-uuid-xxxx"
image: "docker.io/username/ai-inference:latest"
image_pull_policy: "IfNotPresent"
command:
- "python"
- "/app/inference.py"
- "--model"
- "llama-3-70b"
- "--input"
- "/data/input.json"
- "--output"
- "/data/output.json"
env:
- name: "MODEL_PATH"
value: "/models/llama-3-70b"
- name: "CUDA_VISIBLE_DEVICES"
value: "0,1,2,3"
- name: "HF_TOKEN"
value: "hf_xxx..."
secret: true
resources:
gpu_model: "A100"; gpu_count: 4; vcpu: 16; ram_gb: 128; storage_gb: 500
supports_tee: false
inputs:
- name: "model-weights"; url: "ipfs://QmXxx.../model.bin"; path: "/models/"
- name: "input-data"; url: "ipfs://QmYyy.../input.json"; path: "/data/"
outputs:
- name: "results"; path: "/data/output.json"
timeout_seconds: 3600
network: "none"
verification:
method: "hash"
expected_output_hash: "sha256:..."
6.2 Provider Container Executor
# executor/container_executor.py
import asyncio, docker, hashlib, time, aiohttp
from pathlib import Path
from typing import Dict, List, Optional
from dataclasses import dataclass
@dataclass
class ExecutionResult:
task_id: str; success: bool; output_cid: str
duration_seconds: int; logs: str; error: Optional[str] = None
class ContainerExecutor:
def __init__(self, docker_client: docker.DockerClient,
ipfs_endpoint: str = "http://localhost:5001",
data_dir: str = "/data/compute"):
self.docker = docker_client; self.ipfs_endpoint = ipfs_endpoint
self.data_dir = Path(data_dir); self.data_dir.mkdir(parents=True, exist_ok=True)
async def execute_task(self, task_id: str, image: str,
command: List[str], env_vars: Dict[str, str], gpu_count: int,
timeout: int, input_cid: str, output_paths: List[str]) -> ExecutionResult:
start = time.time()
work_dir = self.data_dir / task_id; work_dir.mkdir(parents=True, exist_ok=True)
try:
await self._pull_image(image)
input_dir = work_dir / "input"; input_dir.mkdir(exist_ok=True)
await self._download_from_ipfs(input_cid, input_dir)
device_requests = []
if gpu_count > 0:
device_requests.append(docker.types.DeviceRequest(
driver="nvidia", count=gpu_count, capabilities=[["gpu", "compute", "utility"]]))
output_dir = work_dir / "output"; output_dir.mkdir(exist_ok=True)
ctn = self.docker.containers.run(image=image, command=command,
environment=env_vars, volumes={
str(input_dir): {"bind": "/data/input", "mode": "ro"},
str(output_dir): {"bind": "/data/output", "mode": "rw"},
}, device_requests=device_requests, detach=True, network_mode="none")
try:
result = ctn.wait(timeout=timeout)
logs = ctn.logs(stdout=True, stderr=True).decode("utf-8")
except docker.errors.TimeoutError:
ctn.kill(); ctn.remove()
return ExecutionResult(task_id=task_id, success=False, output_cid="",
duration_seconds=int(time.time() - start), logs="", error="Timeout")
ctn.remove()
if result["StatusCode"] != 0:
return ExecutionResult(task_id=task_id, success=False, output_cid="",
duration_seconds=int(time.time() - start), logs=logs,
error=f"Exit code {result['StatusCode']}")
output_cid = await self._upload_to_ipfs(output_dir)
return ExecutionResult(task_id=task_id, success=True, output_cid=output_cid,
duration_seconds=int(time.time() - start), logs=logs)
except Exception as e:
return ExecutionResult(task_id=task_id, success=False, output_cid="",
duration_seconds=int(time.time() - start), logs="", error=str(e))
async def _pull_image(self, image: str):
loop = asyncio.get_event_loop()
await loop.run_in_executor(None, self.docker.images.pull, image)
async def _download_from_ipfs(self, cid: str, target: Path):
import tarfile, io
async with aiohttp.ClientSession() as session:
async with session.post(f"{self.ipfs_endpoint}/api/v0/get?arg={cid}") as r:
if r.status != 200: raise Exception(f"IPFS download failed: {r.status}")
data = await r.read()
with tarfile.open(fileobj=io.BytesIO(data)) as tar: tar.extractall(path=target)
async def _upload_to_ipfs(self, dir_path: Path) -> str:
import tarfile, io
buf = io.BytesIO()
with tarfile.open(fileobj=buf, mode="w") as tar: tar.add(dir_path, arcname="output")
buf.seek(0)
async with aiohttp.ClientSession() as session:
async with session.post(f"{self.ipfs_endpoint}/api/v0/add?recursive=true&pin=true",
data={"file": buf.read()}) as r:
if r.status != 200: raise Exception(f"IPFS upload failed: {r.status}")
return (await r.json()).get("Hash", "")
6.3 TEE Execution
# executor/tee_executor.py
import hashlib, time, json
from typing import List
from dataclasses import dataclass
@dataclass
class TeeAttestation:
provider: str; task_id: str; enclave_measurement: str
report_data: str; timestamp: int; signature: str
class TeeExecutor:
def __init__(self, tee_type: str = "sgx"):
self.tee_type = tee_type
def prepare_enclave(self, image: str, cmd: List[str]) -> str:
return hashlib.sha256(json.dumps({
"image": image, "command": cmd, "tee_type": self.tee_type
}, sort_keys=True).encode()).hexdigest()
def generate_attestation(self, task_id: str, output_hash: str,
provider_key: str) -> TeeAttestation:
report = hashlib.sha256(f"{task_id}:{output_hash}".encode()).hexdigest()
return TeeAttestation(provider="msg1...", task_id=task_id,
enclave_measurement="mr_enclave", report_data=report,
timestamp=int(time.time()), signature="signed")
def verify_attestation(self, att: TeeAttestation, expected_hash: str) -> bool:
expected = hashlib.sha256(f"{att.task_id}:{expected_hash}".encode()).hexdigest()
return att.report_data == expected
6.4 WASM Execution via CosmWasm
// contracts/compute-executor/src/wasm_executor.rs
use wasmtime::{Engine, Module, Store, Linker};
pub struct WasmExecutionResult { pub output: Vec<u8>, pub gas_used: u64, pub success: bool }
pub fn execute_wasm(wasm_bytes: &[u8], input: &[u8], memory_limit: usize)
-> Result<WasmExecutionResult, String> {
let engine = Engine::default();
let module = Module::new(&engine, wasm_bytes)
.map_err(|e| format!("Failed to compile WASM: {}", e))?;
let mut store = Store::new(&engine, ());
let mut linker = Linker::new(&engine);
linker.func_wrap("env", "read_input", || -> Vec<u8> { input.to_vec() })
.map_err(|e| format!("Failed to bind read_input: {}", e))?;
linker.func_wrap("env", "write_output", |output: Vec<u8>| { Ok(()) })
.map_err(|e| format!("Failed to bind write_output: {}", e))?;
let instance = linker.instantiate(&mut store, &module)
.map_err(|e| format!("Failed to instantiate: {}", e))?;
let main = instance.get_typed_func::<(), i32>(&mut store, "main")
.map_err(|e| format!("Failed to get main: {}", e))?;
let result = main.call(&mut store, ())
.map_err(|e| format!("WASM execution failed: {}", e))?;
Ok(WasmExecutionResult { output: Vec::new(), gas_used: 0, success: result == 0 })
}
6.5 Resource Verification
# executor/resource_verifier.py
import subprocess, json, re, time, os
from typing import Dict, List
class ResourceVerifier:
@staticmethod
def verify_gpu() -> List[Dict]:
try:
r = subprocess.run(["nvidia-smi", "--query-gpu=index,name,uuid,memory.total",
"--format=csv,noheader,nounits"], capture_output=True, text=True, timeout=10)
gpus = []
for line in r.stdout.strip().split("\n"):
if not line: continue
p = [x.strip() for x in line.split(",")]
if len(p) >= 4: gpus.append({"index": int(p[0]), "model": p[1], "uuid": p[2], "memory_mb": int(p[3])})
return gpus
except: return []
@staticmethod
def verify_cpu() -> Dict:
try:
with open("/proc/cpuinfo") as f: c = f.read()
cpus = re.findall(r"processor\s+:\s+(\d+)", c)
model = re.search(r"model name\s+:\s+(.+)", c)
return {"count": len(cpus), "model": model.group(1).strip() if model else "unknown"}
except: return {"count": 0, "model": "unknown"}
@staticmethod
def verify_memory() -> Dict:
try:
with open("/proc/meminfo") as f: c = f.read()
t = re.search(r"MemTota[未公开路径])\s+kB", c)
gb = int(t.group(1)) // (1024 * 1024) if t else 0
return {"total_kb": int(t.group(1)) if t else 0, "total_gb": gb}
except: return {"total_kb": 0, "total_gb": 0}
def generate_attestation_doc(self) -> str:
return json.dumps({"gpu": self.verify_gpu(), "cpu": self.verify_cpu(),
"memory": self.verify_memory(), "timestamp": time.time()})
6.6 Execution Flow
Provider Execution Flow:
1. Poll chain for Matched tasks
2. Download task specification from chain
3. Pull container image from registry
4. Download input data from IPFS
5. Configure GPU / resource limits
6. Execute container with timeout
7. Capture logs and exit code
8. Upload output to IPFS
9. Submit completion (output_cid, duration)
10. Wait for settlement
7. 前端实现
7.1 Provider Dashboard
// src/components/ProviderDashboard.tsx
import React, { useState, useEffect } from "react";
import { useWallet } from "../hooks/useWallet";
import { RegistryClient } from "../clients/registryClient";
import { TaskClient } from "../clients/taskClient";
import { parseAIPAY, formatAIPAY } from "../utils/format";
interface ProviderForm {
providerDid: string; gpuModel: string; gpuCount: number;
vcpu: number; ramGb: number; storageGb: number;
pricePerGpuHour: string; region: string;
}
export const ProviderDashboard: React.FC = () => {
const { address, signingClient } = useWallet();
const [registered, setRegistered] = useState(false);
const [provider, setProvider] = useState<any>(null);
const [tasks, setTasks] = useState<any[]>([]);
const [loading, setLoading] = useState(false);
const [form, setForm] = useState<ProviderForm>({
providerDid: "", gpuModel: "A100", gpuCount: 1,
vcpu: 8, ramGb: 32, storageGb: 100,
pricePerGpuHour: "0.5", region: "us-east-1",
});
const registryAddr = "msg14hj2tavq8fpesdwxxcu44rty3hh90vhujrvcmstl4zr3txmfvw9sss8k4k";
const taskAddr = "msg1task...contract...addr";
useEffect(() => { if (address && signingClient) checkRegistration(); }, [address, signingClient]);
const checkRegistration = async () => {
try {
const r = new RegistryClient(signingClient!, registryAddr);
const info = await r.getProvider(address!);
if (info) { setRegistered(true); setProvider(info); loadTasks(); }
} catch { setRegistered(false); }
};
const loadTasks = async () => {
const t = new TaskClient(signingClient!, taskAddr);
setTasks(await t.listProviderTasks(address!));
};
const handleRegister = async () => {
if (!signingClient || !address) return; setLoading(true);
try {
const r = new RegistryClient(signingClient, registryAddr);
await r.register({
providerDid: form.providerDid,
resources: { gpu_model: form.gpuModel, gpu_count: form.gpuCount,
vcpu: form.vcpu, ram_gb: form.ramGb, storage_gb: form.storageGb,
bandwidth_mbps: 1000, supports_tee: false, tee_type: null },
pricing: { price_per_gpu_hour: parseAIPAY(form.pricePerGpuHour).toString(),
price_per_vcpu_hour: "10000000000000000",
price_per_ram_gb_hour: "5000000000000000",
price_per_storage_gb_hour: "1000000000000000",
min_billing_duration: 300, negotiable: false },
region: form.region,
});
setRegistered(true); await checkRegistration();
} catch (err: any) { alert(err.message); }
finally { setLoading(false); }
};
const handleStatusUpdate = async (status: string) => {
if (!signingClient) return;
await new RegistryClient(signingClient, registryAddr).updateStatus(status);
await checkRegistration();
};
if (!address) return <div className="card"><p>Connect wallet</p></div>;
return (
<div className="dashboard">
<h2>Provider Dashboard</h2>
{!registered ? (
<div className="register-form">
<h3>Register as Compute Provider</h3>
<input placeholder="DID" value={form.providerDid}
onChange={e => setForm({...form, providerDid: e.target.value})} />
<select value={form.gpuModel} onChange={e => setForm({...form, gpuModel: e.target.value})}>
<option>A100</option><option>H100</option><option>RTX4090</option><option>L40S</option>
</select>
<input type="number" placeholder="GPU Count" value={form.gpuCount}
onChange={e => setForm({...form, gpuCount: +e.target.value})} />
<input type="number" placeholder="vCPU" value={form.vcpu}
onChange={e => setForm({...form, vcpu: +e.target.value})} />
<input type="number" placeholder="RAM (GB)" value={form.ramGb}
onChange={e => setForm({...form, ramGb: +e.target.value})} />
<input type="number" placeholder="Storage (GB)" value={form.storageGb}
onChange={e => setForm({...form, storageGb: +e.target.value})} />
<input placeholder="Price per GPU/h (AIPAY)" value={form.pricePerGpuHour}
onChange={e => setForm({...form, pricePerGpuHour: e.target.value})} />
<input placeholder="Region" value={form.region}
onChange={e => setForm({...form, region: e.target.value})} />
<button onClick={handleRegister} disabled={loading}>
{loading ? "Registering..." : "Register Provider"}
</button>
</div>
) : (
<div className="provider-info">
<h3>Status: {provider?.status}</h3>
<p>GPU: {provider?.resources?.gpu_model} x{provider?.resources?.gpu_count}</p>
<p>Region: {provider?.region}</p>
<p>Earned: {formatAIPAY(provider?.total_earned)} AIPAY</p>
<p>Tasks: {provider?.completed_tasks} completed / {provider?.failed_tasks} failed</p>
<div className="actions">
<button onClick={() => handleStatusUpdate("active")}>Set Active</button>
<button onClick={() => handleStatusUpdate("inactive")}>Set Inactive</button>
<button onClick={() => handleStatusUpdate("paused")}>Pause</button>
</div>
<h3>My Tasks</h3>
<ul>
{tasks.map(t => (
<li key={t.id}>{t.id} - {t.status}</li>
))}
</ul>
</div>
)}
</div>
);
};
7.2 Client Dashboard
// src/components/ClientDashboard.tsx
import React, { useState, useEffect } from "react";
import { useWallet } from "../hooks/useWallet";
import { TaskClient } from "../clients/taskClient";
import { RegistryClient } from "../clients/registryClient";
import { parseAIPAY } from "../utils/format";
interface TaskForm {
containerImage: string; gpuModel: string; gpuCount: number;
vcpu: number; ramGb: number; storageGb: number;
maxPrice: string; expectedDuration: number; inputCid: string;
command: string; useTee: boolean;
}
export const ClientDashboard: React.FC = () => {
const { address, signingClient } = useWallet();
const [tasks, setTasks] = useState<any[]>([]);
const [providers, setProviders] = useState<any[]>([]);
const [loading, setLoading] = useState(false);
const [form, setForm] = useState<TaskForm>({
containerImage: "tensorflow/tensorflow:latest-gpu",
gpuModel: "A100", gpuCount: 1, vcpu: 4, ramGb: 16, storageGb: 50,
maxPrice: "1.0", expectedDuration: 3600, inputCid: "",
command: "python /app/train.py", useTee: false,
});
const taskAddr = "msg1task...";
const registryAddr = "msg14hj2tavq8fpesdwxxcu44rty3hh90vhujrvcmstl4zr3txmfvw9sss8k4k";
useEffect(() => {
if (address && signingClient) { loadMyTasks(); loadProviders(); }
}, [address, signingClient]);
const loadMyTasks = async () => {
try {
const t = new TaskClient(signingClient!, taskAddr);
setTasks(await t.listClientTasks(address!));
} catch {}
};
const loadProviders = async () => {
try {
const r = new RegistryClient(signingClient!, registryAddr);
setProviders(await r.queryAvailable({ gpuModel: form.gpuModel }));
} catch {}
};
const handleCreateTask = async () => {
if (!signingClient || !address) return; setLoading(true);
try {
const t = new TaskClient(signingClient, taskAddr);
await t.createTask({
taskType: "training",
containerImage: form.containerImage,
resourceRequirements: {
gpu_model: form.gpuModel, gpu_count: form.gpuCount,
vcpu: form.vcpu, ram_gb: form.ramGb, storage_gb: form.storageGb,
bandwidth_mbps: 1000, supports_tee: form.useTee, tee_type: null,
},
maxPrice: { denom: "uaipay", amount: parseAIPAY(form.maxPrice).toString() },
inputDataCid: form.inputCid,
expectedDuration: form.expectedDuration,
command: form.command.split(" "),
envVars: [], useTee: form.useTee,
});
await loadMyTasks();
} catch (err: any) { alert(err.message); }
finally { setLoading(false); }
};
const handleConfirm = async (taskId: string) => {
try {
const t = new TaskClient(signingClient!, taskAddr);
await t.confirmCompletion(taskId);
await loadMyTasks();
} catch (err: any) { alert(err.message); }
};
return (
<div className="dashboard">
<h2>Client Dashboard</h2>
<div className="create-task">
<h3>Post Compute Task</h3>
<input placeholder="Container Image" value={form.containerImage}
onChange={e => setForm({...form, containerImage: e.target.value})} />
<select value={form.gpuModel}
onChange={e => setForm({...form, gpuModel: e.target.value})}>
<option>A100</option><option>H100</option><option>RTX4090</option>
</select>
<input type="number" placeholder="GPU Count" value={form.gpuCount}
onChange={e => setForm({...form, gpuCount: +e.target.value})} />
<input type="number" placeholder="vCPU" value={form.vcpu}
onChange={e => setForm({...form, vcpu: +e.target.value})} />
<input type="number" placeholder="RAM (GB)" value={form.ramGb}
onChange={e => setForm({...form, ramGb: +e.target.value})} />
<input placeholder="Max Price (AIPAY/h)" value={form.maxPrice}
onChange={e => setForm({...form, maxPrice: e.target.value})} />
<input type="number" placeholder="Expected Duration (s)" value={form.expectedDuration}
onChange={e => setForm({...form, expectedDuration: +e.target.value})} />
<input placeholder="Input Data CID (IPFS)" value={form.inputCid}
onChange={e => setForm({...form, inputCid: e.target.value})} />
<input placeholder="Command (space separated)" value={form.command}
onChange={e => setForm({...form, command: e.target.value})} />
<label><input type="checkbox" checked={form.useTee}
onChange={e => setForm({...form, useTee: e.target.checked})} /> Use TEE</label>
<button onClick={handleCreateTask} disabled={loading}>
{loading ? "Creating..." : "Create Task"}
</button>
</div>
<div className="providers">
<h3>Available Providers ({providers.length})</h3>
<ul>
{providers.map(p => (
<li key={p.provider_did}>{p.provider_did}: {p.resources.gpu_model} x{p.resources.gpu_count} @ {p.region}</li>
))}
</ul>
</div>
<div className="my-tasks">
<h3>My Tasks ({tasks.length})</h3>
<table>
<thead><tr><th>ID</th><th>Type</th><th>Status</th><th>Provider</th><th>Action</th></tr></thead>
<tbody>
{tasks.map(t => (
<tr key={t.id}>
<td>{t.id.slice(0, 16)}...</td>
<td>{t.task_type}</td>
<td className={t.status.toLowerCase()}>{t.status}</td>
<td>{t.provider || "-"}</td>
<td>
{t.status === "completed" && (
<button onClick={() => handleConfirm(t.id)}>Confirm</button>
)}
</td>
</tr>
))}
</tbody>
</table>
</div>
</div>
);
};
7.3 React Client Library
// src/clients/registryClient.ts
import { SigningCosmWasmClient } from "@cosmjs/cosmwasm-stargate";
import { calculateFee, GasPrice } from "@cosmjs/stargate";
export class RegistryClient {
constructor(private client: SigningCosmWasmClient, private addr: string) {}
async register(params: {
providerDid: string; resources: any; pricing: any; region: string
}) {
const fee = calculateFee(200000, GasPrice.fromString("1000000000attoMSG"));
return this.client.execute(this.addr, {
register: {
provider_did: params.providerDid, resources: params.resources,
pricing: params.pricing, region: params.region,
}
}, fee);
}
async updateStatus(status: string) {
const fee = calculateFee(150000, GasPrice.fromString("1000000000attoMSG"));
return this.client.execute(this.addr, { update_status: { status } }, fee);
}
async getProvider(address: string) {
return this.client.queryContractSmart(this.addr, { get_provider: { address } });
}
async queryAvailable(filters: { gpuModel?: string; region?: string; maxPrice?: string }) {
return this.client.queryContractSmart(this.addr, {
query_available_providers: {
gpu_model: filters.gpuModel || null,
region: filters.region || null,
max_price_per_hour: filters.maxPrice || null,
}
});
}
}
// src/clients/taskClient.ts
export class TaskClient {
constructor(private client: SigningCosmWasmClient, private addr: string) {}
async createTask(params: {
taskType: string; containerImage: string; resourceRequirements: any;
maxPrice: any; inputDataCid: string; expectedDuration: number;
command: string[]; envVars: any[]; useTee: boolean;
}) {
const fee = calculateFee(300000, GasPrice.fromString("1000000000attoMSG"));
const deposit = [{ denom: "uaipay", amount: "10000000000000000000" }]; // 10 AIPAY
return this.client.execute(this.addr, {
create_task: {
task_type: params.taskType, container_image: params.containerImage,
resource_requirements: params.resourceRequirements, max_price: params.maxPrice,
input_data_cid: params.inputDataCid, expected_duration: params.expectedDuration,
command: params.command, env_vars: params.envVars, use_tee: params.useTee,
metadata: null,
}
}, fee, "", deposit);
}
async listClientTasks(client: string) {
return this.client.queryContractSmart(this.addr, {
list_client_tasks: { client }
});
}
async listProviderTasks(provider: string) {
return this.client.queryContractSmart(this.addr, {
list_provider_tasks: { provider }
});
}
async confirmCompletion(taskId: string) {
const fee = calculateFee(200000, GasPrice.fromString("1000000000attoMSG"));
return this.client.execute(this.addr, { confirm_completion: { task_id: taskId } }, fee);
}
async getTask(taskId: string) {
return this.client.queryContractSmart(this.addr, { get_task: { task_id: taskId } });
}
}
// src/utils/format.ts
import { Uint128 } from "@cosmjs/cosmwasm-stargate";
const DECIMALS = 18;
export function parseAIPAY(amount: string): bigint {
const parts = amount.split(".");
const whole = parts[0] || "0";
const fraction = (parts[1] || "").padEnd(DECIMALS, "0").slice(0, DECIMALS);
return BigInt(whole + fraction);
}
export function formatAIPAY(amount: string | Uint128): string {
const s = amount.toString().padStart(DECIMALS + 1, "0");
const dotPos = s.length - DECIMALS;
return s.slice(0, dotPos) + "." + s.slice(dotPos);
}
// src/hooks/useWallet.ts
import { useState, useEffect } from "react";
import { SigningCosmWasmClient } from "@cosmjs/cosmwasm-stargate";
import { GasPrice } from "@cosmjs/stargate";
import { OfflineSigner } from "@cosmjs/proto-signing";
interface WalletState {
address: string | null;
signingClient: SigningCosmWasmClient | null;
connect: (signer: OfflineSigner) => Promise<void>;
disconnect: () => void;
}
export function useWallet(): WalletState {
const [address, setAddress] = useState<string | null>(null);
const [signingClient, setSigningClient] = useState<SigningCosmWasmClient | null>(null);
const connect = async (signer: OfflineSigner) => {
const client = await SigningCosmWasmClient.connectWithSigner(
"https://rpc.msg-chain-1.msgchain.zone", signer,
{ gasPrice: GasPrice.fromString("1000000000attoMSG") }
);
const [account] = await signer.getAccounts();
setAddress(account.address);
setSigningClient(client);
};
const disconnect = () => {
setAddress(null); setSigningClient(null);
};
return { address, signingClient, connect, disconnect };
}
8. 完整示例
8.1 端到端流程
本示例展示完整的去中心化算力市场流程:Provider 注册 -> Client 发布任务 -> 匹配 -> 执行 -> 结算。
流程示意图:
Provider (GPU Provider) MSG Chain Client (AI Agent)
| | |
|-- 1. Register Provider --->| |
| (stake 1M AIPAY) | |
|<-- Provider DID -----------| |
| | |
| |<-- 2. Create Task ------|
| | (escrow 5 AIPAY) |
| | |
|-- 3. Bid Task ------------>| |
| (offer: 4.5 AIPAY) | |
| | |
| |-- 4. Task Matched ------>|
|<-- 5. Start Task ----------| |
| | |
| [Execute Container] | |
| [Upload output to IPFS] | |
| | |
|-- 6. Complete Task ------->| |
| (output_cid, duration) | |
| | |
| |<-- 7. Confirm ---------|
| | (release payment) |
|<-- 8. Payment Released ----| |
| (4.4775 AIPAY earned) | |
8.2 Complete TypeScript Integration
// examples/full_workflow.ts
import { SigningCosmWasmClient } from "@cosmjs/cosmwasm-stargate";
import { DirectSecp256k1HdWallet } from "@cosmjs/proto-signing";
import { GasPrice, calculateFee } from "@cosmjs/stargate";
import { RegistryClient } from "../src/clients/registryClient";
import { TaskClient } from "../src/clients/taskClient";
import { parseAIPAY } from "../src/utils/format";
const RPC = "https://rpc.msg-chain-1.msgchain.zone";
async function fullWorkflow() {
// ===== PROVIDER SETUP =====
const providerWallet = await DirectSecp256k1HdWallet.fromMnemonic(
"abandon abandon ... word24", { prefix: "msg" })
const [providerAddr] = await providerWallet.getAccounts();
const client = await SigningCosmWasmClient.connectWithSigner(
RPC, providerWallet, { gasPrice: GasPrice.fromString("1000000000attoMSG") })
// Register provider
const reg = new RegistryClient(client, REGISTRY);
await reg.register({
providerDid: "did:msg:provider:demo-1",
resources: { gpu_model: "A100", gpu_count: 8, vcpu: 64,
ram_gb: 512, storage_gb: 2000, bandwidth_mbps: 10000,
supports_tee: true, tee_type: "Intel SGX" },
pricing: {
price_per_gpu_hour: parseAIPAY("0.5").toString(),
price_per_vcpu_hour: parseAIPAY("0.01").toString(),
price_per_ram_gb_hour: parseAIPAY("0.005").toString(),
price_per_storage_gb_hour: parseAIPAY("0.001").toString(),
min_billing_duration: 300, negotiable: false,
},
region: "us-east-1",
});
console.log("Provider registered");
// ===== CLIENT =====
const clientWallet = await DirectSecp256k1HdWallet.fromMnemonic(
"client client ... word24", { prefix: "msg" })
const [clientAddr] = await clientWallet.getAccounts();
const client2 = await SigningCosmWasmClient.connectWithSigner(
RPC, clientWallet, { gasPrice: GasPrice.fromString("1000000000attoMSG") })
// Post task
const task = new TaskClient(client2, TASK_CONTRACT);
await task.createTask({
taskType: "fine_tuning",
containerImage: "docker.io/myorg/llm-finetuner:latest",
resourceRequirements: { gpu_model: "A100", gpu_count: 4,
vcpu: 16, ram_gb: 128, storage_gb: 500, bandwidth_mbps: 1000,
supports_tee: true, tee_type: null },
maxPrice: { denom: "uaipay", amount: parseAIPAY("5.0").toString() },
inputDataCid: "QmXxx...", expectedDuration: 7200,
command: ["python", "/app/train.py"], envVars: [], useTee: true,
});
// Provider bids, starts, completes
// (All executed via provider wallet)
await client.execute(providerAddr.address, TASK_CONTRACT,
{ bid_task: { task_id: "task-xxx",
bid_price: { denom: "uaipay", amount: parseAIPAY("4.5").toString() } } },
calculateFee(200000, GasPrice.fromString("1000000000attoMSG")));
await client.execute(providerAddr.address, TASK_CONTRACT,
{ start_task: { task_id: "task-xxx" } },
calculateFee(150000, GasPrice.fromString("1000000000attoMSG")));
await client.execute(providerAddr.address, TASK_CONTRACT,
{ complete_task: { task_id: "task-xxx",
output_data_cid: "QmYyy...", actual_duration: 6800 } },
calculateFee(200000, GasPrice.fromString("1000000000attoMSG")));
// Client confirms
await task.confirmCompletion("task-xxx");
console.log("Task completed and paid");
}
fullWorkflow().catch(console.error);
8.3 Python Integration Example
# examples/full_workflow.py
import asyncio
from matcher.compute_matcher import ComputeMatcher, Resources, Pricing, ProviderOrder, ClientOrder, TaskType
from matcher.reverse_auction import ReverseAuctionEngine
from decimal import Decimal
async def run_market_example():
matcher = ComputeMatcher()
# Provider registers GPU resources
provider = ProviderOrder(
order_id="", provider_did="did:msg:provider:demo-1",
provider_addr="msg1provider...",
resources=Resources(gpu_model="A100", gpu_count=4, vcpu=16, ram_gb=128, storage_gb=500),
pricing=Pricing(price_per_gpu_hour=Decimal("0.5"),
price_per_vcpu_hour=Decimal("0.01"),
price_per_ram_gb_hour=Decimal("0.005"),
price_per_storage_gb_hour=Decimal("0.001")),
region="us-east-1", rating=4.8, completed_tasks=42,
status=ProviderStatus.Active, available_from=0, available_until=9999999999,
)
matcher.place_bid(provider)
print(f"Provider registered: {provider.provider_did}")
# Client posts fine-tuning task
task = ClientOrder(
order_id="", client_addr="msg1client...",
task_type=TaskType.FINE_TUNING,
required_resources=Resources(gpu_model="A100", gpu_count=2, vcpu=8, ram_gb=64, storage_gb=200),
max_budget=Decimal("5.0"), expected_duration=7200,
region_preference="us-east-1", min_rating=4.0,
)
matcher.place_ask(task)
print(f"Task posted: {task.order_id}")
# Reverse auction
auction = ReverseAuctionEngine()
auction.create_auction(
task_id="task-001", client_addr="msg1client...",
required_resources=Resources(gpu_model="A100", gpu_count=2, vcpu=8, ram_gb=64, storage_gb=200),
max_budget=Decimal("5.0"), expected_duration=7200,
)
auction.submit_bid("task-001", "msg1provider1...", Decimal("4.5"), 6800, "sig1")
auction.submit_bid("task-001", "msg1provider2...", Decimal("3.8"), 7500, "sig2")
winner = auction.close_auction("task-001")
if winner:
print(f"Auction winner: {winner.bidder} @ {winner.price} AIPAY/h")
print("Market example completed")
asyncio.run(run_market_example())
8.4 Contract Deployment Script
# scripts/deploy.sh
#!/bin/bash
set -e
CHAIN_ID="msg-chain-1"
RPC="https://rpc.msg-chain-1.msgchain.zone"
TX_FLAGS="--chain-id $CHAIN_ID --gas auto --gas-adjustment 1.3 --fees 500uaipay -y"
# Upload contracts
echo "Uploading compute-registry..."
REGISTRY_CODE=$(msgd tx wasm store artifacts/compute_registry.wasm --from admin $TX_FLAGS --output json | jq -r '.logs[0].events[-1].attributes[0].value')
echo "Registry code ID: $REGISTRY_CODE"
echo "Uploading compute-task..."
TASK_CODE=$(msgd tx wasm store artifacts/compute_task.wasm --from admin $TX_FLAGS --output json | jq -r '.logs[0].events[-1].attributes[0].value')
echo "Task code ID: $TASK_CODE"
echo "Uploading compute-payment..."
PAYMENT_CODE=$(msgd tx wasm store artifacts/compute_payment.wasm --from admin $TX_FLAGS --output json | jq -r '.logs[0].events[-1].attributes[0].value')
echo "Payment code ID: $PAYMENT_CODE"
# Instantiate contracts
echo "Instantiating compute-payment..."
PAYMENT_ADDR=$(msgd tx wasm instantiate $PAYMENT_CODE '{"protocol_fee_bps":50,"dispute_fee":{"denom":"uaipay","amount":"10000000000000000000"},"min_fee":{"denom":"uaipay","amount":"1000000000000000000"}}' --label compute-payment --admin $ADMIN $TX_FLAGS --output json | jq -r '.logs[0].events[-1].attributes[0].value')
echo "Payment: $PAYMENT_ADDR"
echo "Instantiating compute-registry..."
REGISTRY_ADDR=$(msgd tx wasm instantiate $REGISTRY_CODE '{"admin":"'$ADMIN'","min_bond":{"denom":"uaipay","amount":"1000000000000000000000000"},"protocol_fee_bps":50,"dispute_duration":86400}' --label compute-registry --admin $ADMIN $TX_FLAGS --output json | jq -r '.logs[0].events[-1].attributes[0].value')
echo "Registry: $REGISTRY_ADDR"
echo "Instantiating compute-task..."
TASK_ADDR=$(msgd tx wasm instantiate $TASK_CODE '{"registry_addr":"'$REGISTRY_ADDR'","aipay_addr":"'$PAYMENT_ADDR'","admin":"'$ADMIN'","dispute_duration":86400,"protocol_fee_bps":50,"min_task_deposit":{"denom":"uaipay","amount":"10000000000000000000"}}' --label compute-task --admin $ADMIN $TX_FLAGS --output json | jq -r '.logs[0].events[-1].attributes[0].value')
echo "Task: $TASK_ADDR"
echo "Deployment complete!"
echo "REGISTRY=$REGISTRY_ADDR"
echo "TASK=$TASK_ADDR"
echo "PAYMENT=$PAYMENT_ADDR"
8.5 Architecture Summary
+----------------------------+
| MSG Chain (msg-chain-1) |
+----------------------------+
| |
| compute_registry.wasm |
| - Provider registration |
| - Resource declarations |
| - Reputation tracking |
| |
| compute_task.wasm |
| - Task lifecycle mgmt |
| - Bid/ask matching |
| - Dispute resolution |
| |
| compute_payment.wasm |
| - Escrow (AIPAY) |
| - Release with fee |
| - Refund / split |
+----------------------------+
| |
| Off-chain: |
| - Python Match Engine |
| - Container Executor |
| - IPFS for data storage |
+----------------------------+
8.6 Key Constants
| Parameter | Value | Description |
|---|---|---|
| Chain ID | msg-chain-1 | MSG Chain identifier |
| Bech32 prefix | msg | Human-readable address prefix |
| Decimals | 18 | AIPAY token precision |
| Min bond | 1,000,000 AIPAY | Provider registration stake |
| Protocol fee | 0.5% (50 bps) | Per-task protocol fee |
| Dispute period | 24 hours (86400s) | Window for dispute filing |
| Min task deposit | 10 AIPAY | Minimum escrow for tasks |
| Min billing | 5 minutes (300s) | Minimum chargeable duration |
本指南完整覆盖了在 MSG Chain 上构建去中心化 AI 算力市场的全部关键组件。开发者可根据实际需求调整合约参数、匹配算法和前端界面。
本文档基于 MSG Chain 代码库核实的技术事实。
白皮书系统: https://msgchain.org/whitepaper/
