AI Agent 模型推断服务市场实现指南
基于 MSG Chain 的去中心化机器学习推理市场
版本: v1.0.0
链: msg-chain-1
Bech32 前缀:msg
精度: 18 位小数
核心合约:agent_llm_v1,agent_payment_v1,agent_registry_v1,agent_a2a_v1
1. 概述
1.1 什么是 AI 模型推断服务市场
AI 模型推断服务市场是一个去中心化的平台,允许 AI Agent 在 MSG Chain 上发布、发现和消费机器学习模型推断服务。这个市场将计算资源与消费需求通过区块链智能合约进行撮合,实现无需信任的模型推断服务交易。
在传统模式下,模型推断通常依赖中心化 API 提供商。而在 MSG Chain 的 Agent 生态中,任何 AI Agent 都可以成为模型服务提供者,也可以通过 A2A 通信协议消费其他 Agent 提供的模型服务。
1.2 基础设施
本市场基于 MSG Chain 的四个核心合约构建:
| 合约 | 地址示例 | 功能 |
|---|---|---|
agent_llm_v1 |
msg1llm... |
LLM 推断执行环境,支持模型加载与推理 |
agent_payment_v1 |
msg1pay... |
AIPAY 支付通道,支持代币化支付、流支付与条件支付 |
agent_registry_v1 |
msg1reg... |
Agent 与模型服务注册表,提供发现与元数据查询 |
agent_a2a_v1 |
msg1a2a... |
Agent 间通信协议,支持消息路由与异步调用 |
1.2.1 agent_llm_v1
agent_llm_v1 是 MSG Chain 上的 AI 推断执行合约。它维护一个链上模型 registry,允许合约调用预注册的机器学习模型。关键特性:
- 链上模型存储与版本管理
- 模型推断的 gas 消耗计量
- 推断结果的上链验证
- 与 agent_payment_v1 集成的按需付费
1.2.2 agent_payment_v1 (AIPAY)
AIPAY 是 MSG Chain 的 Agent 间支付系统:
- 支持
msg原生代币与 CW20 代币支付 - 条件支付(基于预言机或合约条件释放资金)
- 流支付(按时间线性释放)
- 支付通道与小额支付聚合
1.2.3 agent_registry_v1
注册中心负责:
- Agent 身份注册与 DID 管理
- 服务端点发布
- 模型服务元数据注册
- 服务信誉评分
1.2.4 agent_a2a_v1
A2A 通信协议提供:
- Agent 间消息路由
- 异步请求/响应模式
- 消息签名与验证
- 超时与重试机制
1.3 支持的模型类型
本市场支持以下 AI 模型推断类型:
| 类型 | 描述 | 定价单位 |
|---|---|---|
| LLM (Large Language Model) | 文本生成、对话、代码生成 | 每 token |
| Embedding | 文本向量化 | 每 token |
| Image Generation | 文生图、图生图 | 每张 |
| Audio Transcription | 语音转文字 | 每秒 |
| Audio Generation | 文字转语音 | 每秒 |
| Classification | 文本/图像分类 | 每请求 |
| Object Detection | 目标检测 | 每张 |
| Translation | 机器翻译 | 每字符 |
1.4 市场参与者角色
+------------------+ +------------------+
| 模型服务提供者 | | 模型服务消费者 |
| (Provider Agent) | | (Consumer Agent) |
+------------------+ +------------------+
| |
| ① 注册模型服务 |
|-----> [agent_registry_v1] |
| |
| ② 查询可用模型 |
| [agent_registry_v1] <-|
| |
| ③ A2A 推断请求 + 支付意图 |
|<---- [agent_a2a_v1] ------->|
| |
| ④ off-chain 推断执行 |
| |
| ⑤ 返回结果 + 证明 |
|-----> [agent_a2a_v1] ------>|
| |
| ⑥ 支付结算 |
|-----> [agent_payment_v1] --->|
+ +
- Provider Agent: 运行模型推断服务的 AI Agent,注册模型、设定价格、执行推断
- Consumer Agent: 需要模型推断能力的 AI Agent,发现模型、发送请求、支付费用
- 仲裁者 (Arbitrator): 在发生争议时验证推断结果质量的第三方 Agent
1.5 经济模型
- 计价单位:
uamp(1 msg = 10^18 uamp) - 平台手续费:0.5% ~ 2%(由 Provider 承担)
- 最小结算单位:1 uamp
- 支付通道生命周期:按需创建、即时结算
2. 模型服务注册
2.1 数据结构定义
2.1.1 核心结构体
// file: contracts/model-market/src/state.rs
use cosmwasm_std::{Addr, Coin, Decimal, Timestamp, Uint128};
use cw_storage_plus::{Item, Map};
use schemars::JsonSchema;
use serde::{Deserialize, Serialize};
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub enum ServiceStatus { Active, Paused, Retired, Banned }
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
#[serde(rename_all = "snake_case")]
pub enum PricingModel {
PerToken { input_price: Uint128, output_price: Uint128, free_tokens: Uint128 },
PerRequest { price: Uint128 },
PerDuration { price_per_second: Uint128, min_duration: u64 },
PerOutputUnit { price_per_unit: Uint128, unit_description: String },
Subscription { rate_per_block: Uint128, free_requests_per_block: u64, overage_price: Uint128 },
Tiered { tiers: Vec<PriceTier> },
}
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct PriceTier { pub min_quantity: Uint128, pub max_quantity: Option<Uint128>, pub unit_price: Uint128 }
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
#[serde(rename_all = "snake_case")]
pub enum Constraint {
MaxInputTokens(u64), MaxContextLength(u64), MaxOutputTokens(u64),
MaxInputSize(u64), SupportedFormats(Vec<String>), MaxInputFileSize(u64),
MaxBatchSize(u32), TimeoutSeconds(u64), Custom { key: String, value: String },
}
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct ModelService {
pub provider: Addr, pub model_id: String, pub model_name: String,
pub model_version: String, pub model_type: ModelType,
pub capabilities: Vec<String>, pub pricing: PricingModel,
pub input_constraints: Vec<Constraint>, pub status: ServiceStatus,
pub accepted_tokens: Vec<String>, pub created_at: Timestamp,
pub updated_at: Timestamp, pub metadata_uri: Option<String>, pub sla: ServiceLevelAgreement,
}
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
#[serde(rename_all = "snake_case")]
pub enum ModelType { Llm, Embedding, ImageGeneration, AudioTranscription, AudioGeneration, Classification, ObjectDetection, Translation, Other(String) }
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct ServiceLevelAgreement {
pub p95_latency_ms: u64, pub p99_latency_ms: u64,
pub uptime_percentage: Decimal, pub min_quality_score: u8, pub observation_window_seconds: u64,
}
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct ServiceQualityReport {
pub service_id: String, pub reporter: Addr, pub latency_ms: u64,
pub quality_score: u8, pub success: bool, pub error_message: Option<String>, pub timestamp: Timestamp,
}
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct Config {
pub owner: Addr, pub fee_bps: u64, pub min_stake: Coin,
pub dispute_deposit: Coin, pub max_services: u32, pub arbitrators: Vec<Addr>,
pub agent_llm: Addr, pub agent_payment: Addr, pub agent_registry: Addr,
}
2.1.2 存储键定义
pub const CONFIG: Item<Config> = Item::new("config");
pub const SERVICES: Map<&str, ModelService> = Map::new("services");
pub const PROVIDER_SERVICES: Map<&Addr, Vec<String>> = Map::new("prov_services");
pub const TYPE_INDEX: Map<&ModelType, Vec<String>> = Map::new("type_idx");
pub const QUALITY_REPORTS: Map<&str, Vec<ServiceQualityReport>> = Map::new("quality_rpt");
pub const REPUTATION_SCORES: Map<&str, (u8, u64, Decimal)> = Map::new("reputation");
pub const DISPUTES: Map<&str, Dispute> = Map::new("disputes");
pub const SUBSCRIPTIONS: Map<&Addr, Subscription> = Map::new("subscriptions");
pub const NONCE: Map<(&Addr, &str), u64> = Map::new("nonce");
2.2 消息定义
// file: contracts/model-market/src/msg.rs
use cosmwasm_std::Coin;
use schemars::JsonSchema;
use serde::{Deserialize, Serialize};
use crate::state::{Constraint, ModelType, PricingModel, ServiceFilter, ServiceLevelAgreement, ServiceStatus};
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct InstantiateMsg {
pub owner: String, pub fee_bps: u64, pub min_stake: Coin, pub dispute_deposit: Coin,
pub max_services: u32, pub arbitrators: Vec<String>,
pub agent_llm: String, pub agent_payment: String, pub agent_registry: String,
}
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
#[serde(rename_all = "snake_case")]
pub enum ExecuteMsg {
RegisterService {
model_id: String, model_name: String, model_version: String, model_type: ModelType,
capabilities: Vec<String>, pricing: PricingModel, input_constraints: Vec<Constraint>,
accepted_tokens: Vec<String>, metadata_uri: Option<String>, sla: ServiceLevelAgreement,
},
UpdateService {
model_id: String, model_name: Option<String>, model_version: Option<String>,
pricing: Option<PricingModel>, input_constraints: Option<Vec<Constraint>>,
accepted_tokens: Option<Vec<String>>, metadata_uri: Option<String>, sla: Option<ServiceLevelAgreement>,
},
SetServiceStatus { model_id: String, status: ServiceStatus },
SubmitQualityReport { service_id: String, latency_ms: u64, quality_score: u8, success: bool, error_message: Option<String> },
InitiateDispute { service_id: String, request_id: String, reason: String, expected_result_hash: Option<String> },
ResolveDispute { dispute_id: String, resolution: DisputeResolution },
UpdateConfig { fee_bps: Option<u64>, min_stake: Option<Coin>, dispute_deposit: Option<Coin>, max_services: Option<u32>, arbitrators: Option<Vec<String>> },
}
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
#[serde(rename_all = "snake_case")]
pub enum DisputeResolution { InFavorOfConsumer, InFavorOfProvider, Compromise { refund_percentage: Decimal } }
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
#[serde(rename_all = "snake_case")]
pub enum QueryMsg {
GetConfig {}, GetService { model_id: String },
ListServices { start_after: Option<String>, limit: Option<u32> },
QueryServices { filter: ServiceFilter, start_after: Option<String>, limit: Option<u32> },
GetProviderServices { provider: String }, GetQualityReports { service_id: String, limit: Option<u32> },
GetReputation { service_id: String }, GetDispute { dispute_id: String },
GetSubscription { subscriber: String, service_id: String },
}
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct ServiceListResponse { pub services: Vec<ModelService> }
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct ReputationResponse { pub avg_quality: u8, pub total_reports: u64, pub success_rate: Decimal, pub service_id: String }
2.3 合约入口与执行逻辑
// file: contracts/model-market/src/contract.rs
use cosmwasm_std::{entry_point, to_json_binary, Binary, Coin, Deps, DepsMut, Env,
MessageInfo, Response, StdError, StdResult, Timestamp};
use cw2::set_contract_version;
use crate::error::ContractError;
use crate::msg::{DisputeResolution, ExecuteMsg, InstantiateMsg, QueryMsg, ReputationResponse, ServiceListResponse};
use crate::state::{Config, Constraint, Dispute, ModelService, ModelType, PricingModel,
ServiceLevelAgreement, ServiceStatus, Subscription,
CONFIG, DISPUTES, NONCE, PROVIDER_SERVICES, QUALITY_REPORTS,
REPUTATION_SCORES, SERVICES, SUBSCRIPTIONS, TYPE_INDEX};
use crate::validation::validate_service;
const CONTRACT_NAME: &str = "model-market";
const CONTRACT_VERSION: &str = "1.0.0";
#[entry_point]
pub fn instantiate(deps: DepsMut, env: Env, info: MessageInfo, msg: InstantiateMsg) -> Result<Response, ContractError> {
set_contract_version(deps.storage, CONTRACT_NAME, CONTRACT_VERSION)?;
let config = Config {
owner: deps.api.addr_validate(&msg.owner)?,
fee_bps: msg.fee_bps, min_stake: msg.min_stake, dispute_deposit: msg.dispute_deposit,
max_services: msg.max_services,
arbitrators: msg.arbitrators.iter().map(|a| deps.api.addr_validate(a)).collect::<Result<Vec<_>, _>>()?,
agent_llm: deps.api.addr_validate(&msg.agent_llm)?,
agent_payment: deps.api.addr_validate(&msg.agent_payment)?,
agent_registry: deps.api.addr_validate(&msg.agent_registry)?,
};
CONFIG.save(deps.storage, &config)?;
Ok(Response::new().add_attribute("action", "instantiate").add_attribute("owner", msg.owner))
}
#[entry_point]
pub fn execute(deps: DepsMut, env: Env, info: MessageInfo, msg: ExecuteMsg) -> Result<Response, ContractError> {
match msg {
ExecuteMsg::RegisterService { model_id, model_name, model_version, model_type,
capabilities, pricing, input_constraints, accepted_tokens, metadata_uri, sla } =>
execute_register_service(deps, env, info, model_id, model_name, model_version,
model_type, capabilities, pricing, input_constraints, accepted_tokens, metadata_uri, sla),
ExecuteMsg::UpdateService { model_id, model_name, model_version, pricing,
input_constraints, accepted_tokens, metadata_uri, sla } =>
execute_update_service(deps, env, info, model_id, model_name, model_version,
pricing, input_constraints, accepted_tokens, metadata_uri, sla),
ExecuteMsg::SetServiceStatus { model_id, status } => execute_set_status(deps, env, info, model_id, status),
ExecuteMsg::SubmitQualityReport { service_id, latency_ms, quality_score, success, error_message } =>
execute_submit_report(deps, env, info, service_id, latency_ms, quality_score, success, error_message),
ExecuteMsg::InitiateDispute { service_id, request_id, reason, expected_result_hash } =>
execute_initiate_dispute(deps, env, info, service_id, request_id, reason, expected_result_hash),
ExecuteMsg::ResolveDispute { dispute_id, resolution } => execute_resolve_dispute(deps, env, info, dispute_id, resolution),
ExecuteMsg::UpdateConfig { fee_bps, min_stake, dispute_deposit, max_services, arbitrators } =>
execute_update_config(deps, env, info, fee_bps, min_stake, dispute_deposit, max_services, arbitrators),
}
}
2.4 注册示例
export CHAIN_ID="msg-chain-1"
export RPC="https://rpc.msgchain.org"
export CONTRACT="msg1modelmarket..."
msg-chaind tx wasm execute $CONTRACT \
'{"register_service":{
"model_id":"gpt-4-turbo-2026","model_name":"GPT-4 Turbo 2026","model_version":"1.2.0",
"model_type":"llm","capabilities":["text-generation","chat","code-generation"],
"pricing":{"per_token":{"input_price":"30000000000000000","output_price":"60000000000000000","free_tokens":"1000"}},
"input_constraints":[{"max_input_tokens":128000},{"timeout_seconds":120}],
"accepted_tokens":["uamp"],"metadata_uri":"ipfs://QmCard",
"sla":{"p95_latency_ms":2000,"p99_latency_ms":5000,"uptime_percentage":"99.9","min_quality_score":85,"observation_window_seconds":86400}
}}' --amount 1000000uamp --from provider_key --chain-id $CHAIN_ID --gas auto -y
3. 推断请求工作流
3.1 工作流概述
Consumer Agent Provider Agent
| |
| ① 准备推断输入 |
| ② 查询 Registry 发现可用模型 |
| ③ 创建支付意图 + 签名请求 |
| ④ ───────── A2A Request ──────► |
| ⑤ 校验支付意图 & 签名 |
| ⑥ 执行模型推断 |
| ⑦ 生成结果证明 |
| ⑧ ◄────── A2A Response ──────── |
| ⑨ 验证结果证明 |
| ⑩ 调用 Payment 结算 |
| ⑪ 提交质量报告 |
3.2 TypeScript SDK 实现
3.2.1 客户端配置
// src/sdk/config.ts
import { SigningCosmWasmClient } from "@cosmjs/cosmwasm-stargate";
import { DirectSecp256k1HdWallet } from "@cosmjs/proto-signing";
import { GasPrice } from "@cosmjs/stargate";
export interface MarketConfig {
rpcEndpoint: string;
chainId: string;
prefix: string;
contractAddresses: { modelMarket: string; agentRegistry: string; agentPayment: string; agentA2A: string; agentLlm: string; };
gasPrice: GasPrice;
gasLimit: number;
}
export const DEFAULT_CONFIG: MarketConfig = {
rpcEndpoint: "https://rpc.msgchain.org",
chainId: "msg-chain-1",
prefix: "msg",
contractAddresses: { modelMarket: "msg1modelmarket...", agentRegistry: "msg1registry...", agentPayment: "msg1payment...", agentA2A: "msg1a2a...", agentLlm: "msg1llm..." },
gasPrice: GasPrice.fromString("1000000000attoMSG"),
gasLimit: 2_000_000,
};
export class MarketClient {
public client!: SigningCosmWasmClient;
public wallet!: DirectSecp256k1HdWallet;
public address!: string;
public config: MarketConfig;
constructor(config: Partial<MarketConfig> = {}) { this.config = { ...DEFAULT_CONFIG, ...config }; }
async connect(mnemonic: string): Promise<void> {
this.wallet = await DirectSecp256k1HdWallet.fromMnemonic(mnemonic, { prefix: this.config.prefix });
this.client = await SigningCosmWasmClient.connectWithSigner(this.config.rpcEndpoint, this.wallet, { gasPrice: this.config.gasPrice });
const accounts = await this.wallet.getAccounts();
this.address = accounts[0].address;
}
get modelMarket(): string { return this.config.contractAddresses.modelMarket; }
get agentRegistry(): string { return this.config.contractAddresses.agentRegistry; }
get agentPayment(): string { return this.config.contractAddresses.agentPayment; }
get agentA2A(): string { return this.config.contractAddresses.agentA2A; }
get agentLlm(): string { return this.config.contractAddresses.agentLlm; }
}
3.2.2 推断请求客户端
// src/sdk/inference.ts
import { toUtf8, fromBase64, toBase64 } from "@cosmjs/encoding";
import { SHA256 } from "@cosmjs/crypto";
import { MarketClient } from "./config";
export interface ModelInfo {
model_id: string; model_name: string; model_version: string; model_type: string;
capabilities: string[]; provider: string; pricing: any; input_constraints: any[];
status: string; sla: { p95_latency_ms: number; p99_latency_ms: number; uptime_percentage: string; min_quality_score: number; };
}
export interface InferenceInput { prompt?: string; messages?: ChatMessage[]; image?: string; audio?: string; parameters: ModelParameters; }
export interface ChatMessage { role: "system" | "user" | "assistant"; content: string; }
export interface ModelParameters { temperature?: number; top_p?: number; max_tokens?: number; stop?: string[]; frequency_penalty?: number; presence_penalty?: number; stream?: boolean; }
export interface PaymentIntent { escrow_id: string; amount: string; denom: string; expiration: number; terms: string; }
export interface InferenceResponse { request_id: string; model_id: string; output: InferenceOutput; usage: TokenUsage; attestation: ResultAttestation; provider_sig: string; timestamp: number; }
export interface InferenceOutput { text?: string; image?: string; audio?: string; embeddings?: number[][]; classification?: { label: string; confidence: number }[]; error?: string; }
export interface TokenUsage { input_tokens: number; output_tokens: number; total_tokens: number; }
export interface ResultAttestation { output_hash: string; model_id: string; input_hash: string; timestamp: number; provider_address: string; }
export class InferenceClient {
private market: MarketClient;
constructor(market: MarketClient) { this.market = market; }
async queryModels(filter?: { model_types?: string[]; min_quality?: number; capabilities?: string[] }): Promise<ModelInfo[]> {
const result: any = await this.market.client.queryContractSmart(this.market.modelMarket, {
query_services: { filter: { model_types: filter?.model_types, min_quality: filter?.min_quality ?? 0, capabilities: filter?.capabilities, status: "active" }, limit: 50 }
});
return result.services as ModelInfo[];
}
async selectOptimalModel(requirements: { model_type: string; estimated_input_tokens: number; estimated_output_tokens: number; max_price: string; min_quality: number; required_capabilities: string[] }): Promise<ModelInfo | null> {
const models = await this.queryModels({ model_types: [requirements.model_type], min_quality: requirements.min_quality, capabilities: requirements.required_capabilities });
if (models.length === 0) return null;
const scored = await Promise.all(models.map(async (model) => {
const price = await this.estimateCost(model.model_id, requirements.estimated_input_tokens, requirements.estimated_output_tokens);
if (price > BigInt(requirements.max_price)) return null;
const rep = await this.getReputation(model.model_id);
const quality = model.sla.min_quality_score;
const reliability = rep ? rep.successRate : 0.95;
const priceScore = Number(BigInt(requirements.max_price) - price) / Number(BigInt(requirements.max_price));
const totalScore = priceScore * 0.4 + (quality / 100) * 0.3 + reliability * 0.3;
return { model, score: totalScore };
}));
const valid = scored.filter(Boolean) as { model: ModelInfo; score: number }[];
valid.sort((a, b) => b.score - a.score);
return valid.length > 0 ? valid[0].model : null;
}
async estimateCost(modelId: string, inputTokens: number, outputTokens: number): Promise<bigint> {
const service: any = await this.market.client.queryContractSmart(this.market.modelMarket, { get_service: { model_id: modelId } });
const pricing = service.pricing;
const type = Object.keys(pricing)[0];
const p = pricing[type];
switch (type) {
case "per_token": {
const free = BigInt(p.free_tokens ?? 0);
const chargeable = BigInt(inputTokens) > free ? BigInt(inputTokens) - free : BigInt(0);
return chargeable * BigInt(p.input_price) + BigInt(outputTokens) * BigInt(p.output_price);
}
case "per_request": return BigInt(p.price);
case "tiered": {
const totalUnits = BigInt(inputTokens + outputTokens);
for (const tier of p.tiers) {
const minQ = BigInt(tier.min_quantity);
const maxQ = tier.max_quantity ? BigInt(tier.max_quantity) : null;
if (totalUnits >= minQ && (!maxQ || totalUnits < maxQ)) return totalUnits * BigInt(tier.unit_price);
}
return totalUnits * BigInt(p.tiers[p.tiers.length - 1].unit_price);
}
default: throw new Error(`Unknown pricing: ${type}`);
}
}
async getReputation(serviceId: string): Promise<{ avgQuality: number; totalReports: number; successRate: number } | null> {
try {
const result: any = await this.market.client.queryContractSmart(this.market.modelMarket, { get_reputation: { service_id: serviceId } });
return { avgQuality: result.avg_quality, totalReports: result.total_reports, successRate: parseFloat(result.success_rate) };
} catch { return null; }
}
async sendInferenceRequest(modelId: string, input: InferenceInput, options?: { maxPrice?: string; timeoutMs?: number }): Promise<InferenceResponse> {
const service: any = await this.market.client.queryContractSmart(this.market.modelMarket, { get_service: { model_id: modelId } });
if (service.status !== "Active") throw new Error(`Service ${modelId} is not active`);
const provider = service.provider;
const requestId = `${this.market.address}-${modelId}-${Date.now()}-${Math.random().toString(36).slice(2)}`;
const paymentIntent = await this.createPaymentIntent(provider, modelId, input, options?.maxPrice);
await this.market.client.execute(this.market.address, this.market.agentA2A, {
send_message: { recipient: provider, message_type: "inference_request", payload: toBase64(toUtf8(JSON.stringify({ request_id: requestId, model_id: modelId, input, payment_intent: paymentIntent, timestamp: Math.floor(Date.now() / 1000), consumer_sig: "" }))), correlation_id: requestId }
}, "auto", "Inference request: " + modelId);
const response = await this.pollForResponse(requestId, options?.timeoutMs ?? 60000);
this.verifyAttestation(response);
return response;
}
private async createPaymentIntent(provider: string, modelId: string, input: InferenceInput, maxPrice?: string): Promise<PaymentIntent> {
const inputTokens = Math.ceil((input.prompt ?? input.messages?.map(m => m.content).join(" ") ?? "").length / 4);
const outputTokens = input.parameters.max_tokens ?? 1024;
const estimatedCost = await this.estimateCost(modelId, inputTokens, outputTokens);
const finalAmount = maxPrice ? (BigInt(maxPrice) < estimatedCost ? BigInt(maxPrice) : estimatedCost) : estimatedCost;
const escrowResult = await this.market.client.execute(this.market.address, this.market.agentPayment, {
create_escrow: { recipient: provider, amount: finalAmount.toString(), denom: "uamp", expiration: Math.floor(Date.now() / 1000) + 3600, conditions: { oracle: this.market.modelMarket, condition_type: "inference_complete" } }
}, "auto", "Create escrow");
const escrowId = escrowResult.events.find(e => e.type === "wasm")?.attributes.find(a => a.key === "escrow_id")?.value ?? "escrow_unknown";
return { escrow_id: escrowId, amount: finalAmount.toString(), denom: "uamp", expiration: Math.floor(Date.now() / 1000) + 3600, terms: JSON.stringify({ model_id: modelId, max_input_tokens: inputTokens, max_output_tokens: outputTokens }) };
}
private async pollForResponse(requestId: string, timeoutMs: number): Promise<InferenceResponse> {
const start = Date.now();
while (Date.now() - start < timeoutMs) {
const result: any = await this.market.client.queryContractSmart(this.market.agentA2A, { get_messages: { address: this.market.address, limit: 10, message_type: "inference_response" } });
const msg = (result.messages ?? []).find((m: any) => m.correlation_id === requestId);
if (msg) return JSON.parse(new TextDecoder().decode(fromBase64(msg.payload))) as InferenceResponse;
await new Promise(r => setTimeout(r, 1000));
}
throw new Error(`Request ${requestId} timed out`);
}
verifyAttestation(response: InferenceResponse): boolean {
const hash = toBase64(SHA256(toUtf8(JSON.stringify(response.output))));
if (hash !== response.attestation.output_hash) throw new Error("Attestation verification failed");
return true;
}
}
3.3 Provider 端 Python
# src/provider/inference_server.py
import asyncio, hashlib, json, time, logging
from dataclasses import dataclass, field
from typing import Optional
from cosmpy.aerial.client import LedgerClient
from cosmpy.aerial.wallet import LocalWallet
from cosmpy.aerial.contract import CosmWasmClient
logger = logging.getLogger(__name__)
@dataclass
class ModelConfig:
model_id: str; model_path: str; model_type: str; device: str = "cuda"
max_batch_size: int = 8; max_input_length: int = 128000; temperature: float = 0.7; top_p: float = 0.9
@dataclass
class InferenceRequest:
request_id: str; model_id: str; consumer: str; prompt: Optional[str] = None
messages: Optional[list[dict]] = None; parameters: dict = field(default_factory=dict)
payment_intent: dict = field(default_factory=dict); consumer_sig: str = ""; timestamp: int = 0
@dataclass
class InferenceResult:
request_id: str; model_id: str; output: dict; usage: dict; attestation: dict; provider_sig: str; timestamp: int
class InferenceEngine:
def __init__(self, config: ModelConfig):
self.config = config; self.tokenizer = None; self.model = None; self.lock = asyncio.Lock()
async def load(self):
from transformers import AutoModelForCausalLM, AutoTokenizer
logger.info(f"Loading model {self.config.model_id}")
loop = asyncio.get_event_loop()
self.tokenizer = await loop.run_in_executor(None, lambda: AutoTokenizer.from_pretrained(self.config.model_path, trust_remote_code=True))
self.model = await loop.run_in_executor(None, lambda: AutoModelForCausalLM.from_pretrained(self.config.model_path, torch_dtype=torch.float16, device_map="auto", trust_remote_code=True))
self.model.eval()
async def infer(self, request: InferenceRequest) -> InferenceResult:
prompt = self._build_prompt(request)
loop = asyncio.get_event_loop()
output_text, usage = await loop.run_in_executor(None, lambda: self._run_inference(prompt, request.parameters))
output = {"text": output_text}
ts = int(time.time())
return InferenceResult(request_id=request.request_id, model_id=self.config.model_id, output=output, usage=usage,
attestation={"output_hash": hashlib.sha256(json.dumps(output, sort_keys=True).encode()).hexdigest(),
"model_id": self.config.model_id, "input_hash": hashlib.sha256(prompt.encode()).hexdigest(),
"timestamp": ts, "provider_address": ""}, provider_sig="", timestamp=ts)
def _build_prompt(self, r: InferenceRequest) -> str:
if r.messages: return "\n".join([f"<|{m['role']}|>\n{m['content']}" for m in r.messages]) + "\n<|assistant|>\n"
return r.prompt or ""
def _run_inference(self, prompt: str, params: dict) -> tuple:
import torch
inputs = self.tokenizer(prompt, return_tensors="pt")
with torch.no_grad():
outputs = self.model.generate(inputs.input_ids.to(self.model.device),
max_new_tokens=params.get("max_tokens", 4096), temperature=params.get("temperature", 0.7),
do_sample=True, pad_token_id=self.tokenizer.eos_token_id)
text = self.tokenizer.decode(outputs[0][inputs.input_ids.shape[1]:], skip_special_tokens=True)
usage = {"input_tokens": inputs.input_ids.shape[1], "output_tokens": outputs.shape[1] - inputs.input_ids.shape[1], "total_tokens": outputs.shape[1]}
return text, usage
class InferenceProvider:
def __init__(self, wallet: LocalWallet, ledger: LedgerClient, engine: InferenceEngine, contracts: dict):
self.wallet, self.ledger, self.engine, self.contracts, self.cosmwasm = wallet, ledger, engine, contracts, CosmWasmClient(ledger)
self.running = False
async def start(self):
self.running = True
self.ledger.execute_contract(self.contracts["model_market"], {"register_service": {
"model_id": self.engine.config.model_id, "model_name": self.engine.config.model_id, "model_version": "1.0.0",
"model_type": "llm", "capabilities": ["text-generation", "chat"],
"pricing": {"per_token": {"input_price": "30000000000000000", "output_price": "60000000000000000", "free_tokens": "1000"}},
"input_constraints": [{"max_input_tokens": 128000}, {"timeout_seconds": 120}],
"accepted_tokens": ["uamp"], "metadata_uri": None,
"sla": {"p95_latency_ms": 2000, "p99_latency_ms": 5000, "uptime_percentage": "99.9", "min_quality_score": 85, "observation_window_seconds": 86400}
}}, self.wallet, funds="1000000uamp")
await asyncio.gather(self._listen(), self._events())
async def _listen(self):
while self.running:
try:
result = self.cosmwasm.query(self.contracts["agent_a2a"], {"get_messages": {"address": str(self.wallet.address()), "limit": 10, "message_type": "inference_request"}})
for msg in result.get("messages", []): asyncio.create_task(self._handle(msg))
except Exception as e: logger.error(f"Poll error: {e}")
await asyncio.sleep(1)
async def _handle(self, raw: dict):
try:
p = json.loads(raw["payload"])
req = InferenceRequest(request_id=p["request_id"], model_id=p["model_id"], consumer=raw.get("sender",""),
prompt=p.get("input",{}).get("prompt"), messages=p.get("input",{}).get("messages"),
parameters=p.get("input",{}).get("parameters",{}), payment_intent=p.get("payment_intent",{}),
timestamp=p.get("timestamp",0))
result = await self.engine.infer(req)
result.attestation["provider_address"] = str(self.wallet.address())
sig = self.wallet.keypair.sign(json.dumps({"request_id": result.request_id, "output_hash": result.attestation["output_hash"], "timestamp": result.timestamp}, sort_keys=True).encode())
result.provider_sig = sig.hex()
self.ledger.execute_contract(self.contracts["agent_a2a"], {"send_message": {"recipient": result.attestation["provider_address"], "message_type": "inference_response", "payload": json.dumps({"request_id": result.request_id, "model_id": result.model_id, "output": result.output, "usage": result.usage, "attestation": result.attestation, "provider_sig": result.provider_sig, "timestamp": result.timestamp}), "correlation_id": result.request_id}}, self.wallet)
self.ledger.execute_contract(self.contracts["agent_payment"], {"release_escrow": {"escrow_id": req.payment_intent["escrow_id"]}}, self.wallet)
except Exception as e: logger.error(f"Handle error: {e}")
async def _events(self): pass
3.4 Consumer Python
# src/consumer/agent_client.py
import hashlib, json, time, os
from typing import Optional
from cosmpy.aerial.client import LedgerClient, NetworkConfig
from cosmpy.aerial.wallet import LocalWallet
from cosmpy.aerial.contract import CosmWasmClient
from cosmpy.crypto.keypairs import PrivateKey
class ConsumerAgent:
def __init__(self, wallet: LocalWallet, ledger: LedgerClient, contracts: dict):
self.wallet, self.ledger, self.cosmwasm, self.contracts = wallet, ledger, CosmWasmClient(ledger), contracts
def discover(self, model_type: str = None, caps: list[str] = None) -> list[dict]:
result = self.cosmwasm.query(self.contracts["model_market"], {"query_services": {"filter": {"model_types": [model_type] if model_type else None, "capabilities": caps, "status": "active"}, "limit": 50}})
services = result.get("services", [])
for s in services:
try: s["reputation"] = self.cosmwasm.query(self.contracts["model_market"], {"get_reputation": {"service_id": s["model_id"]}})
except: s["reputation"] = {"avg_quality": 50, "success_rate": "0.9"}
return services
def estimate_cost(self, pricing: dict, in_t: int, out_t: int) -> Optional[int]:
t = list(pricing.keys())[0]; p = pricing[t]
if t == "per_token": return max(0, in_t - int(p.get("free_tokens",0))) * int(p["input_price"]) + out_t * int(p["output_price"])
if t == "per_request": return int(p["price"])
if t == "tiered":
total = in_t + out_t
for tier in p["tiers"]:
mx = tier.get("max_quantity")
if total >= int(tier["min_quantity"]) and (mx is None or total < int(mx)): return total * int(tier["unit_price"])
return None
def infer(self, model: dict, prompt: str, max_tokens: int = 1024, timeout: int = 60) -> dict:
cost = self.estimate_cost(model["pricing"], len(prompt)//4, max_tokens)
tx = self.ledger.execute_contract(self.contracts["agent_payment"], {"create_escrow": {"recipient": model["provider"], "amount": str(cost), "denom": "uamp", "expiration": int(time.time())+3600, "conditions": {"oracle": self.contracts["model_market"], "condition_type": "inference_complete"}}}, self.wallet, funds=f"{cost}uamp")
eid = None
for ev in tx.events:
if ev.type == "wasm":
for a in ev.attributes:
if a.key == "escrow_id": eid = a.value
rid = f"{self.wallet.address()}-{model['model_id']}-{int(time.time())}"
self.ledger.execute_contract(self.contracts["agent_a2a"], {"send_message": {"recipient": model["provider"], "message_type": "inference_request", "payload": json.dumps({"request_id": rid, "model_id": model["model_id"], "input": {"prompt": prompt, "parameters": {"max_tokens": max_tokens}}, "payment_intent": {"escrow_id": eid, "amount": str(cost), "denom": "uamp", "expiration": int(time.time())+3600}, "timestamp": int(time.time()), "consumer_sig": ""}), "correlation_id": rid}}, self.wallet)
start = time.time()
while time.time() - start < timeout:
result = self.cosmwasm.query(self.contracts["agent_a2a"], {"get_messages": {"address": str(self.wallet.address()), "limit": 10, "message_type": "inference_response"}})
for m in result.get("messages", []):
if m.get("correlation_id") == rid: return json.loads(m["payload"])
time.sleep(1)
raise TimeoutError("Request timed out")
4. 结果验证
4.1 Hash 承诺
// file: contracts/model-market/src/attestation.rs
use cosmwasm_std::{Binary, Deps, StdResult};
use sha2::{Digest, Sha256};
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct InferenceAttestation {
pub output_hash: String, pub model_id: String, pub input_hash: String,
pub timestamp: u64, pub provider_address: String, pub provider_signature: Binary,
pub params_hash: Option<String>,
}
pub fn verify_attestation(deps: Deps, att: &InferenceAttestation, expected_output: &str, expected_input: &str, provider_pubkey: &[u8]) -> StdResult<bool> {
if hash_string(expected_output) != att.output_hash { return Ok(false); }
if hash_string(expected_input) != att.input_hash { return Ok(false); }
let msg = format!("{}{}{}", att.model_id, att.output_hash, att.timestamp);
deps.api.secp256k1_verify(msg.as_bytes(), &att.provider_signature, provider_pubkey)
}
pub fn hash_string(input: &str) -> String {
let mut h = Sha256::new(); h.update(input.as_bytes()); format!("{:x}", h.finalize())
}
4.2 争议解决
// file: contracts/model-market/src/dispute.rs
use cosmwasm_std::{Addr, Coin, Decimal, DepsMut, Env, MessageInfo, Response, Timestamp, Uint128, BankMsg};
use crate::error::ContractError;
use crate::msg::DisputeResolution;
use crate::state::{CONFIG, DISPUTES, SERVICES};
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub enum DisputeStatus { Pending, Arbitrating, Resolved, Expired }
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct Dispute {
pub dispute_id: String, pub service_id: String, pub request_id: String,
pub consumer: Addr, pub provider: Addr, pub reason: String,
pub expected_output_hash: Option<String>, pub actual_output_hash: Option<String>,
pub status: DisputeStatus, pub created_at: Timestamp,
pub resolved_at: Option<Timestamp>, pub resolution: Option<DisputeResolution>,
pub arbitrator: Option<Addr>, pub deposit: Coin,
}
pub fn execute_initiate_dispute(deps: DepsMut, env: Env, info: MessageInfo, service_id: String, request_id: String, reason: String, expected_hash: Option<String>) -> Result<Response, ContractError> {
let config = CONFIG.load(deps.storage)?;
let service = SERVICES.may_load(deps.storage, &service_id)?.ok_or(ContractError::ServiceNotFound { model_id: service_id.clone() })?;
if !info.funds.iter().any(|c| c.denom == config.dispute_deposit.denom && c.amount >= config.dispute_deposit.amount) {
return Err(ContractError::InsufficientStake { required: config.dispute_deposit });
}
let d = Dispute { dispute_id: format!("d-{}-{}", request_id, env.block.height), service_id, request_id,
consumer: info.sender.clone(), provider: service.provider.clone(), reason,
expected_output_hash: expected_hash, actual_output_hash: None, status: DisputeStatus::Pending,
created_at: env.block.time, resolved_at: None, resolution: None, arbitrator: None, deposit: config.dispute_deposit.clone() };
DISPUTES.save(deps.storage, &d.dispute_id, &d)?;
Ok(Response::new().add_attribute("action", "initiate_dispute"))
}
5. 定价模型
// file: contracts/model-market/src/pricing.rs
use cosmwasm_std::{Decimal, StdResult, Uint128};
use crate::state::{PriceTier, PricingModel};
pub fn calculate_cost(pricing: &PricingModel, in_tokens: u64, out_tokens: u64, duration: Option<u64>, qty: Option<u64>) -> StdResult<Uint128> {
match pricing {
PricingModel::PerToken { input_price, output_price, free_tokens } => {
let free = Uint128::from(*free_tokens);
let chargeable = if Uint128::from(in_tokens) > free { Uint128::from(in_tokens) - free } else { Uint128::zero() };
Ok(chargeable * input_price + Uint128::from(out_tokens) * output_price)
}
PricingModel::PerRequest { price } => Ok(*price),
PricingModel::PerDuration { price_per_second, min_duration } =>
Ok(Uint128::from(duration.unwrap_or(*min_duration).max(*min_duration)) * price_per_second),
PricingModel::PerOutputUnit { price_per_unit, .. } => Ok(Uint128::from(qty.unwrap_or(1)) * price_per_unit),
PricingModel::Subscription { rate_per_block, free_requests_per_block, overage_price } => {
let total = qty.unwrap_or(1);
let overage = if total > *free_requests_per_block { Uint128::from(total - *free_requests_per_block) * overage_price } else { Uint128::zero() };
Ok(*rate_per_block + overage)
}
PricingModel::Tiered { tiers } => {
let total = Uint128::from(in_tokens + out_tokens);
for t in tiers {
if total >= t.min_quantity && t.max_quantity.map_or(true, |m| total < m) { return Ok(total * t.unit_price); }
}
Ok(total * tiers.last().map_or(Uint128::zero(), |t| t.unit_price))
}
}
}
6. 服务级别协议
// file: contracts/model-market/src/sla.rs
use cosmwasm_std::{Decimal, Timestamp, Uint128};
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct SlaTerms { pub latency: LatencyCommitment, pub availability: AvailabilityCommitment, pub quality: QualityCommitment, pub compensation: CompensationScheme, pub observation_window: u64, pub min_samples: u32 }
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct LatencyCommitment { pub p50_max_ms: u64, pub p95_max_ms: u64, pub p99_max_ms: u64 }
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct AvailabilityCommitment { pub uptime_percentage: Decimal, pub max_consecutive_failures: u32 }
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct QualityCommitment { pub min_quality_score: u8, pub min_semantic_similarity: Option<Decimal> }
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct CompensationScheme { pub latency_violation_refund_bps: u64, pub availability_violation_refund_bps: u64, pub quality_violation_refund_bps: u64, pub max_refund_bps: u64 }
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct SlaSnapshot { pub service_id: String, pub window_start: Timestamp, pub window_end: Timestamp, pub total_requests: u64, pub successful_requests: u64, pub latency_p50: u64, pub latency_p95: u64, pub latency_p99: u64, pub avg_quality_score: u8, pub violations: Vec<SlaViolation>, pub compensation_due: Uint128, pub snapshot_hash: String }
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct SlaViolation { pub violation_type: ViolationType, pub observed_value: String, pub threshold_value: String, pub timestamp: Timestamp, pub request_id: Option<String> }
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub enum ViolationType { LatencyP50, LatencyP95, LatencyP99, Uptime, QualityScore, ConsecutiveFailure }
pub fn check_sla_compliance(terms: &SlaTerms, snap: &SlaSnapshot) -> Vec<SlaViolation> {
let mut v = Vec::new();
if snap.latency_p50 > terms.latency.p50_max_ms { v.push(SlaViolation { violation_type: ViolationType::LatencyP50, observed_value: snap.latency_p50.to_string(), threshold_value: terms.latency.p50_max_ms.to_string(), timestamp: snap.window_end, request_id: None }); }
if snap.latency_p95 > terms.latency.p95_max_ms { v.push(SlaViolation { violation_type: ViolationType::LatencyP95, observed_value: snap.latency_p95.to_string(), threshold_value: terms.latency.p95_max_ms.to_string(), timestamp: snap.window_end, request_id: None }); }
if snap.latency_p99 > terms.latency.p99_max_ms { v.push(SlaViolation { violation_type: ViolationType::LatencyP99, observed_value: snap.latency_p99.to_string(), threshold_value: terms.latency.p99_max_ms.to_string(), timestamp: snap.window_end, request_id: None }); }
let uptime = if snap.total_requests > 0 { Decimal::from_ratio(Uint128::from(snap.successful_requests), Uint128::from(snap.total_requests)) * Decimal::percent(100) } else { Decimal::percent(100) };
if uptime < terms.availability.uptime_percentage { v.push(SlaViolation { violation_type: ViolationType::Uptime, observed_value: uptime.to_string(), threshold_value: terms.availability.uptime_percentage.to_string(), timestamp: snap.window_end, request_id: None }); }
if snap.avg_quality_score < terms.quality.min_quality_score { v.push(SlaViolation { violation_type: ViolationType::QualityScore, observed_value: snap.avg_quality_score.to_string(), threshold_value: terms.quality.min_quality_score.to_string(), timestamp: snap.window_end, request_id: None }); }
v
}
pub fn calculate_compensation(terms: &SlaTerms, violations: &[SlaViolation], fees: Uint128) -> Uint128 {
if violations.is_empty() { return Uint128::zero(); }
let mut bps: u64 = 0;
for vio in violations {
match vio.violation_type {
ViolationType::LatencyP50 | ViolationType::LatencyP95 | ViolationType::LatencyP99 => bps = bps.saturating_add(terms.compensation.latency_violation_refund_bps),
ViolationType::Uptime | ViolationType::ConsecutiveFailure => bps = bps.saturating_add(terms.compensation.availability_violation_refund_bps),
ViolationType::QualityScore => bps = bps.saturating_add(terms.compensation.quality_violation_refund_bps),
}
}
fees * Uint128::from(bps.min(terms.compensation.max_refund_bps)) / Uint128::from(10000u64)
}
7. 前端实现
// src/ui/ModelMarketplace.tsx
import React, { useState, useEffect, useCallback } from "react";
import { MarketClient } from "../sdk/config";
import { InferenceClient, ModelInfo } from "../sdk/inference";
export const ModelMarketplace: React.FC<{ marketClient: MarketClient }> = ({ marketClient }) => {
const [models, setModels] = useState<ModelInfo[]>([]);
const [search, setSearch] = useState("");
const [filterType, setFilterType] = useState("all");
const [loading, setLoading] = useState(true);
const client = new InferenceClient(marketClient);
const fetch = useCallback(async () => {
setLoading(true);
try { setModels(await client.queryModels({ model_types: filterType !== "all" ? [filterType] : undefined })); }
catch (e) { console.error(e); } finally { setLoading(false); }
}, [filterType, client]);
useEffect(() => { fetch(); }, [fetch]);
const filtered = models.filter(m => m.model_name.toLowerCase().includes(search.toLowerCase()) || m.model_id.toLowerCase().includes(search.toLowerCase()));
return (
<div className="model-marketplace">
<header className="market-header">
<h1>AI 模型推断市场</h1>
<div className="market-controls">
<input type="text" placeholder="搜索模型..." value={search} onChange={e => setSearch(e.target.value)} className="search-input" />
<select value={filterType} onChange={e => setFilterType(e.target.value)}>
<option value="all">全部类型</option>
<option value="llm">LLM</option>
<option value="embedding">Embedding</option>
<option value="image_generation">图像生成</option>
<option value="audio_transcription">语音识别</option>
</select>
</div>
</header>
{loading ? <div className="loading">加载中...</div> : (
<div className="model-grid">
{filtered.map(m => <ModelCard key={m.model_id} model={m} />)}
{filtered.length === 0 && <div className="empty">未找到模型</div>}
</div>
)}
</div>
);
};
const ModelCard: React.FC<{ model: ModelInfo }> = ({ model }) => {
const [expanded, setExpanded] = useState(false);
return (
<div className={`model-card ${model.status === "Active" ? "" : "inactive"}`}>
<div className="card-header">
<span className={`badge ${model.status.toLowerCase()}`}>{model.status === "Active" ? "活跃" : model.status}</span>
<h3>{model.model_name}</h3>
<span className="version">v{model.model_version}</span>
</div>
<div className="card-body">
<div className="type">{model.model_type}</div>
<div className="id">{model.model_id}</div>
<div className="provider">提供者: {model.provider.slice(0, 12)}...</div>
<div className="sla">P95 {model.sla.p95_latency_ms}ms | {model.sla.uptime_percentage}%</div>
<div className="caps">{(model.capabilities ?? []).slice(0, 3).map(c => <span key={c} className="tag">{c}</span>)}</div>
</div>
{expanded && <div className="details"><pre>{JSON.stringify(model.pricing, null, 2)}</pre></div>}
<button onClick={() => setExpanded(!expanded)}>{expanded ? "收起" : "展开"}</button>
</div>
);
};
8. AI Agent 集成
8.1 Provider Agent
# src/agent/provider_agent.py
import logging
from cosmpy.aerial.client import LedgerClient
from cosmpy.aerial.wallet import LocalWallet
from cosmpy.aerial.contract import CosmWasmClient
logger = logging.getLogger(__name__)
class ProviderAgent:
def __init__(self, wallet: LocalWallet, ledger: LedgerClient, contracts: dict):
self.wallet, self.ledger, self.cosmwasm, self.contracts = wallet, ledger, CosmWasmClient(ledger), contracts
def register(self, config: dict) -> str:
tx = self.ledger.execute_contract(self.contracts["model_market"], {"register_service": {
"model_id": config["model_id"], "model_name": config["model_name"],
"model_version": config.get("model_version", "1.0.0"), "model_type": config["model_type"],
"capabilities": config.get("capabilities", []), "pricing": config["pricing"],
"input_constraints": config.get("input_constraints", []),
"accepted_tokens": config.get("accepted_tokens", ["uamp"]),
"metadata_uri": config.get("metadata_uri"), "sla": config["sla"]
}}, self.wallet, funds=config.get("stake", "1000000uamp"))
logger.info(f"Registered: {config['model_id']}")
return tx.tx_hash
def pause(self, model_id: str) -> str:
return self.ledger.execute_contract(self.contracts["model_market"], {"set_service_status": {"model_id": model_id, "status": "paused"}}, self.wallet).tx_hash
def resume(self, model_id: str) -> str:
return self.ledger.execute_contract(self.contracts["model_market"], {"set_service_status": {"model_id": model_id, "status": "active"}}, self.wallet).tx_hash
def retire(self, model_id: str) -> str:
return self.ledger.execute_contract(self.contracts["model_market"], {"set_service_status": {"model_id": model_id, "status": "retired"}}, self.wallet).tx_hash
def earnings(self) -> dict:
return self.cosmwasm.query(self.contracts["model_market"], {"get_provider_earnings": {"provider": str(self.wallet.address())}})
def withdraw(self) -> str:
return self.ledger.execute_contract(self.contracts["model_market"], {"withdraw_earnings": {}}, self.wallet).tx_hash
8.2 Cost Optimization
# src/agent/cost_optimizer.py
import statistics
from collections import defaultdict
from typing import Optional
class CostOptimizer:
def __init__(self, agent):
self.agent = agent
self.price_history: dict[str, list[int]] = defaultdict(list)
self.quality_history: dict[str, list[int]] = defaultdict(list)
def record(self, model_id: str, cost: int, quality: int):
self.price_history[model_id].append(cost)
self.quality_history[model_id].append(quality)
def suggest(self, model_type: str, in_tokens: int, out_tokens: int, quality_min: int = 70) -> Optional[dict]:
models = self.agent.discover(model_type=model_type)
scored = []
for m in models:
cost = self.agent.estimate_cost(m["pricing"], in_tokens, out_tokens)
if cost is None: continue
avg_q = statistics.mean(self.quality_history.get(m["model_id"], [])) or m["sla"]["min_quality_score"]
if avg_q < quality_min: continue
scored.append((avg_q / (cost + 1), m))
scored.sort(key=lambda x: x[0], reverse=True)
return scored[0][1] if scored else None
def allocate_budget(self, tasks: list[dict], budget: int) -> list[dict]:
result = []
for task in tasks:
model = self.suggest(task.get("model_type", "llm"), task.get("input", 500), task.get("output", 1024))
if model:
cost = self.agent.estimate_cost(model["pricing"], task.get("input", 500), task.get("output", 1024)) or 0
if cost <= budget:
budget -= cost
result.append({"model": model, "cost": cost, "task": task})
return result
8.3 Auto Pipeline
# src/agent/inference_pipeline.py
import asyncio, logging, time
from dataclasses import dataclass
from typing import Optional
logger = logging.getLogger(__name__)
@dataclass
class PipelineTask:
task_id: str; prompt: str; model_type: str = "llm"; max_tokens: int = 1024
temperature: float = 0.7; max_price: Optional[int] = None; priority: int = 0
@dataclass
class PipelineResult:
task_id: str; output: str; usage: dict; cost: int; latency_ms: int; model_used: str; provider: str
class InferencePipeline:
def __init__(self, consumer, max_retries: int = 3, budget: int = None):
self.agent = consumer; self.max_retries = max_retries; self.budget = budget
self.total_cost = 0; self.queue = asyncio.Queue(); self.results = {}; self.running = False
async def start(self, workers: int = 3):
self.running = True
await asyncio.gather(*[self._worker(f"w{i}") for i in range(workers)])
async def submit(self, task: PipelineTask):
await self.queue.put(task)
async def _worker(self, name: str):
while self.running:
try:
task = await asyncio.wait_for(self.queue.get(), 1.0)
result = await self._execute(task)
self.results[task.task_id] = result
except asyncio.TimeoutError: continue
except Exception as e: logger.error(f"{name}: {e}")
async def _execute(self, task: PipelineTask) -> PipelineResult:
for attempt in range(self.max_retries):
try:
if self.budget and self.total_cost >= self.budget: raise RuntimeError("Budget exceeded")
model = self.agent.suggest(task.model_type, len(task.prompt)//4, task.max_tokens)
if not model: raise RuntimeError("No model")
start = time.time()
resp = self.agent.infer(model, task.prompt, task.max_tokens)
cost = self.agent.estimate_cost(model["pricing"], len(task.prompt)//4, task.max_tokens) or 0
self.total_cost += cost
return PipelineResult(task_id=task.task_id, output=resp.get("output",{}).get("text",""),
usage=resp.get("usage",{}), cost=cost, latency_ms=int((time.time()-start)*1000),
model_used=model["model_id"], provider=model["provider"])
except Exception as e:
logger.warning(f"Attempt {attempt+1} failed: {e}")
await asyncio.sleep(2**attempt)
raise RuntimeError(f"All retries failed")
async def batch(self, tasks: list[PipelineTask]) -> list[PipelineResult]:
for t in tasks: await self.submit(t)
while len(self.results) < len(tasks): await asyncio.sleep(0.1)
return [self.results[t.task_id] for t in tasks]
附录 A. 合约开发环境搭建
# 安装 Rust
curl --proto '=https' --tlsv1.2 -sSf https://sh.rustup.rs | sh
# 安装 wasm32 目标
rustup target add wasm32-unknown-unknown
# 安装 CosmWasm 优化器
cargo install cosmwasm-check
# 克隆项目
git clone https://github.com/msgchain/model-market.git
cd model-market
# 编译合约
RUSTFLAGS='-C link-arg=-s' cargo wasm
cosmwasm-check target/wasm32-unknown-unknown/release/model_market.wasm
# 部署
msg-chaind tx wasm store target/wasm32-unknown-unknown/release/model_market.wasm --from deployer --chain-id msg-chain-1 --gas auto -y
# 实例化
msg-chaind tx wasm instantiate <CODE_ID> '{"owner":"msg1...","fee_bps":50,"min_stake":{"denom":"uamp","amount":"1000000"},"dispute_deposit":{"denom":"uamp","amount":"1000000000000000000"},"max_services":10,"arbitrators":["msg1arb..."],"agent_llm":"msg1llm...","agent_payment":"msg1payment...","agent_registry":"msg1registry..."}' --label "model-market-v1" --from deployer --chain-id msg-chain-1 -y
附录 B. 测试脚本
// test/integration.test.ts
import { MarketClient } from "../src/sdk/config";
import { InferenceClient } from "../src/sdk/inference";
describe("Model Marketplace Integration", () => {
let client: MarketClient;
let inference: InferenceClient;
beforeAll(async () => {
client = new MarketClient({ rpcEndpoint: "https://rpc.msgchain.org" });
await client.connect("test mnemonic phrase here");
inference = new InferenceClient(client);
});
test("queryModels returns active services", async () => {
const models = await inference.queryModels();
expect(models.length).toBeGreaterThan(0);
expect(models[0].model_id).toBeDefined();
});
test("estimateCost returns valid cost for per_token model", async () => {
const cost = await inference.estimateCost("gpt-4-turbo-2026", 500, 1024);
expect(cost).toBeGreaterThan(BigInt(0));
});
test("selectOptimalModel returns best model", async () => {
const best = await inference.selectOptimalModel({
model_type: "llm", estimated_input_tokens: 500,
estimated_output_tokens: 1024, max_price: "1000000000000000000",
min_quality: 70, required_capabilities: ["chat"],
});
expect(best).not.toBeNull();
});
});
附录 C. 错误码参考
| 错误码 | 描述 | 解决方法 |
|---|---|---|
| ServiceAlreadyExists | 模型 ID 已存在 | 使用唯一 model_id |
| InsufficientStake | 质押不足 | 增加质押金额 |
| MaxServicesExceeded | 超过最大服务数 | 下架不需要的服务 |
| InvalidPricing | 定价无效 | 确保价格为正 |
| InvalidQualityScore | 质量分无效(0-100) | 调整分数范围 |
| DisputeAlreadyResolved | 争议已解决 | 查询现有结果 |
| NotArbitrator | 非仲裁者 | 联系平台所有者 |
附录 D. Cargo.toml
[package]
name = "model-market"
version = "1.0.0"
edition = "2021"
[lib]
crate-type = ["cdylib", "rlib"]
[dependencies]
cosmwasm-std = "1.5"
cw-storage-plus = "1.2"
cw2 = "1.1"
schemars = "0.8"
serde = { version = "1.0", features = ["derive"] }
thiserror = "1.0"
sha2 = "0.10"
[dev-dependencies]
cosmwasm-vm = "1.5"
4. 结果验证(续)
4.3 质量验证(客户端统计验证)
消费者 Agent 在收到推断结果后,不仅验证 Hash 承诺,还应对结果质量进行统计验证。以下 Python 代码实现了客户端质量验证器:
# src/consumer/quality_verifier.py
import hashlib, json, time, statistics, logging
from dataclasses import dataclass
from typing import Optional
import numpy as np
logger = logging.getLogger(__name__)
@dataclass
class QualityMetrics:
response_time_ms: float
output_length: int
output_entropy: float
token_diversity: float
repetition_ratio: float
stop_reason: str
@dataclass
class VerificationResult:
passed: bool
score: int
issues: list[str]
metrics: QualityMetrics
class QualityVerifier:
def __init__(self, min_quality_score: int = 70, max_repetition_ratio: float = 0.4,
min_entropy: float = 0.5, max_response_time_ms: float = 30000):
self.min_quality = min_quality_score
self.max_repetition = max_repetition_ratio
self.min_entropy = min_entropy
self.max_response_time = max_response_time_ms
def verify(self, prompt: str, output: str, usage: dict,
response_time_ms: float) -> VerificationResult:
issues = []
metrics = self._compute_metrics(prompt, output, usage, response_time_ms)
if metrics.response_time_ms > self.max_response_time:
issues.append(f"响应时间 {metrics.response_time_ms:.0f}ms 超过阈值 {self.max_response_time:.0f}ms")
if metrics.output_entropy < self.min_entropy:
issues.append(f"输出熵 {metrics.output_entropy:.3f} 低于最小值 {self.min_entropy}")
if metrics.repetition_ratio > self.max_repetition:
issues.append(f"重复比率 {metrics.repetition_ratio:.3f} 超过阈值 {self.max_repetition}")
if metrics.output_length < 5:
issues.append("输出过短或为空")
if metrics.stop_reason == "error":
issues.append("模型返回错误")
score = self._compute_score(metrics, len(issues))
passed = score >= self.min_quality and len(issues) == 0
return VerificationResult(passed=passed, score=score, issues=issues, metrics=metrics)
def _compute_metrics(self, prompt: str, output: str, usage: dict,
response_time_ms: float) -> QualityMetrics:
output_len = len(output)
if output_len == 0:
return QualityMetrics(
response_time_ms=response_time_ms, output_length=0,
output_entropy=0.0, token_diversity=0.0,
repetition_ratio=1.0, stop_reason="empty"
)
char_counts = {}
for c in output:
char_counts[c] = char_counts.get(c, 0) + 1
total = output_len
entropy = -sum((c/total) * np.log2(c/total) for c in char_counts.values())
tokens = output.split()
unique_tokens = len(set(tokens))
token_diversity = unique_tokens / max(len(tokens), 1)
ngrams = set()
repeat_count = 0
for i in range(len(tokens) - 1):
pair = (tokens[i], tokens[i+1])
if pair in ngrams:
repeat_count += 1
ngrams.add(pair)
repetition_ratio = repeat_count / max(len(tokens) - 1, 1)
return QualityMetrics(
response_time_ms=response_time_ms, output_length=output_len,
output_entropy=entropy, token_diversity=token_diversity,
repetition_ratio=repetition_ratio,
stop_reason=usage.get("stop_reason", "complete"),
)
def _compute_score(self, metrics: QualityMetrics, issue_count: int) -> int:
score = 100
score -= max(0, int((metrics.response_time_ms - 1000) / 200))
score -= max(0, int((1.0 - metrics.output_entropy) * 30))
score -= max(0, int(metrics.repetition_ratio * 50))
score -= issue_count * 10
return max(0, min(100, score))
class CrossValidator:
"""跨提供者交叉验证 - 将同一请求发送给多个提供者进行结果比较"""
def __init__(self, consumer, num_validators: int = 2):
self.consumer = consumer
self.num_validators = num_validators
async def validate(self, model: dict, prompt: str, max_tokens: int = 256) -> dict:
results = []
for i in range(self.num_validators):
try:
resp = self.consumer.infer(model, prompt, max_tokens)
results.append(resp)
except Exception as e:
logger.warning(f"Validator {i} failed: {e}")
if len(results) < 2:
return {"valid": True, "reason": "insufficient_validators"}
texts = [r.get("output", {}).get("text", "") for r in results]
similarities = self._compute_similarities(texts)
avg_sim = statistics.mean(similarities) if similarities else 0
return {
"valid": avg_sim > 0.7,
"similarity": avg_sim,
"results": results,
"reason": "cross_validation_pass" if avg_sim > 0.7 else "cross_validation_fail",
}
def _compute_similarities(self, texts: list[str]) -> list[float]:
similarities = []
for i in range(len(texts)):
for j in range(i + 1, len(texts)):
sim = self._jaccard_similarity(texts[i], texts[j])
similarities.append(sim)
return similarities
def _jaccard_similarity(self, a: str, b: str) -> float:
set_a = set(a.lower().split())
set_b = set(b.lower().split())
intersection = set_a & set_b
union = set_a | set_b
return len(intersection) / max(len(union), 1)
4.4 欺诈证明构造
当提供者返回垃圾结果时,消费者可以构造欺诈证明提交上链。以下 Rust 代码实现欺诈证明的链上验证:
// file: contracts/model-market/src/fraud_proof.rs
use cosmwasm_std::{Addr, Binary, DepsMut, Env, MessageInfo, Response, StdResult, Uint128};
use sha2::{Digest, Sha256};
use crate::error::ContractError;
use crate::state::{CONFIG, DISPUTES, REPUTATION_SCORES, SERVICES};
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct FraudProof {
pub dispute_id: String,
pub service_id: String,
pub request_id: String,
pub consumer: Addr,
pub provider: Addr,
pub claimed_output_hash: String,
pub actual_output_hash: String,
pub input_commitment: String,
pub provider_response: Binary,
pub independent_verification: Option<IndependentVerification>,
}
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct IndependentVerification {
pub verifier: Addr,
pub verification_result: String,
pub verification_hash: String,
pub signature: Binary,
}
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct QualityFraudReport {
pub service_id: String,
pub request_id: String,
pub consumer: Addr,
pub provider: Addr,
pub expected_quality_score: u8,
pub actual_quality_score: u8,
pub evidence_hash: String,
pub fraud_type: FraudType,
}
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub enum FraudType {
GarbageOutput,
EmptyResponse,
IncorrectModel,
LatencyManipulation,
QualityMismatch,
AttestationForgery,
}
pub fn verify_fraud_proof(
deps: DepsMut,
env: Env,
info: MessageInfo,
proof: FraudProof,
) -> Result<Response, ContractError> {
let config = CONFIG.load(deps.storage)?;
if info.sender != proof.consumer
&& !config.arbitrators.contains(&info.sender)
{
return Err(ContractError::Unauthorized {});
}
if proof.claimed_output_hash == proof.actual_output_hash {
return Err(ContractError::FraudProofInvalid {
reason: "claimed and actual hashes match".to_string(),
});
}
let service = SERVICES.load(deps.storage, &proof.service_id)?;
let current_reputation = REPUTATION_SCORES
.may_load(deps.storage, &proof.service_id)?
.unwrap_or((50u8, 0u64, cosmwasm_std::Decimal::percent(100)));
let (avg_quality, total_reports, success_rate) = current_reputation;
let new_quality = avg_quality.saturating_sub(20);
let new_success_rate = cosmwasm_std::Decimal::from_ratio(
Uint128::from(total_reports.saturating_sub(1)),
Uint128::from(total_reports.max(1)),
);
REPUTATION_SCORES.save(
deps.storage,
&proof.service_id,
&(new_quality, total_reports + 1, new_success_rate),
)?;
Ok(Response::new()
.add_attribute("action", "fraud_proven")
.add_attribute("service_id", &proof.service_id)
.add_attribute("provider", &proof.provider.to_string())
.add_attribute("new_quality", &new_quality.to_string()))
}
4.5 AIPAY 挑战窗口争议解决流程
当消费者对推理结果不满时,可在 AIPAY 的挑战窗口期内发起争议。以下是完整的争议解决流程实现:
// file: contracts/model-market/src/dispute_resolution.rs
use cosmwasm_std::{Addr, Coin, Decimal, DepsMut, Env, MessageInfo,
Response, Timestamp, Uint128, BankMsg, CosmosMsg};
use crate::error::ContractError;
use crate::msg::DisputeResolution;
use crate::state::{CONFIG, DISPUTES, REPUTATION_SCORES, SERVICES};
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct ChallengeWindow {
pub escrow_id: String,
pub consumer: Addr,
pub provider: Addr,
pub amount: Coin,
pub created_at: Timestamp,
pub challenge_deadline: Timestamp,
pub resolution_deadline: Timestamp,
pub status: ChallengeStatus,
pub consumer_evidence: Option<String>,
pub provider_evidence: Option<String>,
pub arbitrator_id: Option<Addr>,
}
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub enum ChallengeStatus {
AwaitingConsumer,
ConsumerChallenged,
ProviderResponded,
UnderArbitration,
Resolved,
Expired,
}
pub fn execute_challenge_response(
deps: DepsMut,
env: Env,
info: MessageInfo,
dispute_id: String,
provider_evidence: String,
) -> Result<Response, ContractError> {
let mut dispute = DISPUTES.load(deps.storage, &dispute_id)?;
let config = CONFIG.load(deps.storage)?;
if info.sender != dispute.provider {
return Err(ContractError::Unauthorized {});
}
if dispute.status != crate::state::DisputeStatus::Pending {
return Err(ContractError::DisputeAlreadyResolved {});
}
let challenge_window_blocks = 100u64;
if env.block.height > dispute.created_at.plus_seconds(challenge_window_blocks * 6).nanos() / 1_000_000_000 {
dispute.status = crate::state::DisputeStatus::Expired;
DISPUTES.save(deps.storage, &dispute_id, &dispute)?;
return Err(ContractError::ChallengeWindowExpired {});
}
let arbitrator = config.arbitrators
.get(env.block.height as usize % config.arbitrators.len())
.ok_or(ContractError::NoArbitratorsAvailable {})?
.clone();
dispute.arbitrator = Some(arbitrator.clone());
dispute.status = crate::state::DisputeStatus::Arbitrating;
DISPUTES.save(deps.storage, &dispute_id, &dispute)?;
Ok(Response::new()
.add_attribute("action", "challenge_response")
.add_attribute("dispute_id", &dispute_id)
.add_attribute("arbitrator", &arbitrator.to_string()))
}
pub fn execute_arbitration(
deps: DepsMut,
env: Env,
info: MessageInfo,
dispute_id: String,
resolution: DisputeResolution,
) -> Result<Response, ContractError> {
let config = CONFIG.load(deps.storage)?;
let mut dispute = DISPUTES.load(deps.storage, &dispute_id)?;
if !config.arbitrators.contains(&info.sender) {
return Err(ContractError::NotArbitrator {});
}
if dispute.status != crate::state::DisputeStatus::Arbitrating {
return Err(ContractError::InvalidDisputeStatus {});
}
let mut messages: Vec<CosmosMsg> = Vec::new();
let deposit = dispute.deposit.clone();
match &resolution {
DisputeResolution::InFavorOfConsumer => {
messages.push(BankMsg::Send {
to_address: dispute.consumer.to_string(),
amount: vec![deposit.clone()],
}.into());
messages.push(BankMsg::Send {
to_address: dispute.consumer.to_string(),
amount: vec![Coin {
denom: config.dispute_deposit.denom.clone(),
amount: config.dispute_deposit.amount,
}],
}.into());
apply_reputation_penalty(deps, &dispute.service_id, 30)?;
}
DisputeResolution::InFavorOfProvider => {
messages.push(BankMsg::Send {
to_address: dispute.provider.to_string(),
amount: vec![deposit.clone()],
}.into());
messages.push(BankMsg::Send {
to_address: dispute.provider.to_string(),
amount: vec![Coin {
denom: config.dispute_deposit.denom.clone(),
amount: config.dispute_deposit.amount,
}],
}.into());
}
DisputeResolution::Compromise { refund_percentage } => {
let refund = Coin {
denom: config.dispute_deposit.denom.clone(),
amount: config.dispute_deposit.amount * *refund_percentage / Decimal::percent(100),
};
let provider_keep = Coin {
denom: config.dispute_deposit.denom.clone(),
amount: config.dispute_deposit.amount - refund.amount,
};
messages.push(BankMsg::Send {
to_address: dispute.consumer.to_string(),
amount: vec![refund],
}.into());
messages.push(BankMsg::Send {
to_address: dispute.provider.to_string(),
amount: vec![provider_keep],
}.into());
}
}
dispute.status = crate::state::DisputeStatus::Resolved;
dispute.resolved_at = Some(env.block.time);
dispute.resolution = Some(resolution);
DISPUTES.save(deps.storage, &dispute_id, &dispute)?;
Ok(Response::new()
.add_messages(messages)
.add_attribute("action", "arbitration_complete")
.add_attribute("dispute_id", &dispute_id))
}
fn apply_reputation_penalty(
deps: DepsMut,
service_id: &str,
penalty_points: u8,
) -> Result<(), ContractError> {
let current = REPUTATION_SCORES
.may_load(deps.storage, service_id)?
.unwrap_or((50u8, 0u64, Decimal::percent(100)));
let new_quality = current.0.saturating_sub(penalty_points);
let new_reports = current.1 + 1;
let new_rate = Decimal::from_ratio(
Uint128::from(current.1.saturating_sub(1)),
Uint128::from(new_reports.max(1)),
);
REPUTATION_SCORES.save(deps.storage, service_id, &(new_quality, new_reports, new_rate))?;
Ok(())
}
4.6 信誉影响评分
以下 Rust 代码实现了争议对提供者信誉评分的自动影响机制:
// file: contracts/model-market/src/reputation.rs
use cosmwasm_std::{Decimal, DepsMut, StdResult, Storage, Uint128};
use crate::state::{REPUTATION_SCORES, QUALITY_REPORTS, ServiceQualityReport};
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct ReputationScore {
pub service_id: String,
pub overall_score: u8,
pub total_requests: u64,
pub success_rate: Decimal,
pub avg_latency_ms: u64,
pub dispute_count: u32,
pub dispute_win_rate: Decimal,
pub recent_trend: i8,
pub confidence_level: ConfidenceLevel,
}
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub enum ConfidenceLevel {
Low,
Medium,
High,
}
pub fn compute_reputation(
storage: &dyn Storage,
service_id: &str,
penalty: Option<u8>,
) -> StdResult<ReputationScore> {
let (avg_quality, total_reports, success_rate) = REPUTATION_SCORES
.may_load(storage, service_id)?
.unwrap_or((50u8, 0u64, Decimal::percent(90)));
let reports: Vec<ServiceQualityReport> = QUALITY_REPORTS
.may_load(storage, service_id)?
.unwrap_or_default();
let total = reports.len() as u64;
let avg_latency = if total > 0 {
reports.iter().map(|r| r.latency_ms).sum::<u64>() / total
} else {
0
};
let success_count = reports.iter().filter(|r| r.success).count() as u64;
let final_success_rate = if total > 0 {
Decimal::from_ratio(Uint128::from(success_count), Uint128::from(total))
} else {
success_rate
};
let final_quality = match penalty {
Some(p) => avg_quality.saturating_sub(p),
None => avg_quality,
};
let recent: Vec<&ServiceQualityReport> = reports.iter().rev().take(10).collect();
let recent_avg = if !recent.is_empty() {
recent.iter().map(|r| r.quality_score).sum::<u8>() as f64 / recent.len() as f64
} else {
final_quality as f64
};
let trend = (recent_avg as i8) - (final_quality as i8);
let confidence = if total_reports < 100 {
ConfidenceLevel::Low
} else if total_reports < 1000 {
ConfidenceLevel::Medium
} else {
ConfidenceLevel::High
};
Ok(ReputationScore {
service_id: service_id.to_string(),
overall_score: final_quality,
total_requests: total_reports,
success_rate: final_success_rate,
avg_latency_ms: avg_latency,
dispute_count: 0,
dispute_win_rate: Decimal::percent(50),
recent_trend: trend,
confidence_level: confidence,
})
}
pub fn decay_reputation(storage: &mut dyn Storage, service_id: &str) -> StdResult<()> {
let current = REPUTATION_SCORES
.may_load(storage, service_id)?
.unwrap_or((50u8, 0u64, Decimal::percent(100)));
let (quality, reports, rate) = current;
if reports > 0 {
let decay = reports.saturating_sub(1) % 100;
if decay == 0 && quality > 30 {
let new_quality = quality - 1;
REPUTATION_SCORES.save(storage, service_id, &(new_quality, reports, rate))?;
}
}
Ok(())
}
4.7 完整客户端验证代码
以下 Python 代码整合了所有客户端验证逻辑:
# src/consumer/inference_verifier.py
import hashlib, json, time, logging
from dataclasses import dataclass
from typing import Optional
from cryptography.hazmat.primitives import hashes, serialization
from cryptography.hazmat.primitives.asymmetric import ec
from cryptography.exceptions import InvalidSignature
logger = logging.getLogger(__name__)
@dataclass
class VerificationChain:
attestation_hash_valid: bool
provider_signature_valid: bool
input_integrity_valid: bool
quality_check_passed: bool
price_accuracy_valid: bool
timestamp_freshness_valid: bool
all_passed: bool
details: dict
class InferenceVerifier:
"""完整的客户端推理结果验证器"""
def __init__(self, max_timestamp_drift_seconds: int = 300):
self.max_drift = max_timestamp_drift_seconds
def verify_all(self, request: dict, response: dict,
provider_pubkey_hex: str) -> VerificationChain:
details = {}
checks = {}
try:
output = response.get("output", {})
computed_hash = hashlib.sha256(
json.dumps(output, sort_keys=True).encode()
).hexdigest()
claimed_hash = response.get("attestation", {}).get("output_hash", "")
checks["attestation_hash_valid"] = computed_hash == claimed_hash
details["computed_output_hash"] = computed_hash
details["claimed_output_hash"] = claimed_hash
except Exception as e:
checks["attestation_hash_valid"] = False
details["hash_error"] = str(e)
try:
att = response.get("attestation", {})
sig_hex = response.get("provider_sig", "")
sign_message = json.dumps({
"request_id": response.get("request_id"),
"output_hash": att.get("output_hash"),
"timestamp": response.get("timestamp"),
}, sort_keys=True).encode()
pub_key_bytes = bytes.fromhex(provider_pubkey_hex)
public_key = ec.EllipticCurvePublicKey.from_encoded_point(
ec.SECP256R1(), pub_key_bytes
)
signature = bytes.fromhex(sig_hex)
public_key.verify(signature, sign_message, ec.ECDSA(hashes.SHA256()))
checks["provider_signature_valid"] = True
except (InvalidSignature, Exception) as e:
checks["provider_signature_valid"] = False
details["signature_error"] = str(e)
try:
input_data = request.get("input", {})
input_hash_claimed = att.get("input_hash", "")
input_hash_computed = hashlib.sha256(
json.dumps(input_data, sort_keys=True).encode()
).hexdigest()
checks["input_integrity_valid"] = input_hash_computed == input_hash_claimed
except Exception as e:
checks["input_integrity_valid"] = False
details["input_hash_error"] = str(e)
try:
ts = response.get("timestamp", 0)
now = int(time.time())
drift = abs(now - ts)
checks["timestamp_freshness_valid"] = drift <= self.max_drift
details["timestamp_drift_seconds"] = drift
except Exception as e:
checks["timestamp_freshness_valid"] = False
details["timestamp_error"] = str(e)
try:
usage = response.get("usage", {})
output_text = output.get("text", "")
output_tokens = usage.get("output_tokens", 0)
estimated_tokens = len(output_text) // 4
ratio = abs(output_tokens - estimated_tokens) / max(estimated_tokens, 1)
checks["usage_consistency_valid"] = ratio < 0.5
except Exception as e:
checks["usage_consistency_valid"] = False
details["usage_error"] = str(e)
try:
payment = request.get("payment_intent", {})
claimed_amount = int(payment.get("amount", "0"))
checks["price_accuracy_valid"] = claimed_amount > 0
except Exception as e:
checks["price_accuracy_valid"] = False
details["price_error"] = str(e)
all_passed = all(checks.values())
return VerificationChain(
attestation_hash_valid=checks.get("attestation_hash_valid", False),
provider_signature_valid=checks.get("provider_signature_valid", False),
input_integrity_valid=checks.get("input_integrity_valid", False),
quality_check_passed=checks.get("usage_consistency_valid", False),
price_accuracy_valid=checks.get("price_accuracy_valid", False),
timestamp_freshness_valid=checks.get("timestamp_freshness_valid", False),
all_passed=all_passed,
details=details,
)
def build_dispute_package(self, request: dict, response: dict,
verification: VerificationChain) -> dict:
"""构建提交链上争议的完整数据包"""
return {
"dispute_package": {
"request_id": response.get("request_id"),
"service_id": response.get("model_id"),
"reason": "result_verification_failed",
"failed_checks": [
k for k, v in {
"hash_mismatch": not verification.attestation_hash_valid,
"sig_invalid": not verification.provider_signature_valid,
"input_mismatch": not verification.input_integrity_valid,
"quality_fail": not verification.quality_check_passed,
"stale_timestamp": not verification.timestamp_freshness_valid,
}.items() if v
],
"expected_output_hash": verification.details.get("computed_output_hash", ""),
"request_snapshot": json.dumps(request, sort_keys=True),
"response_snapshot": json.dumps(response, sort_keys=True),
"consumer_address": request.get("consumer", ""),
"timestamp": int(time.time()),
}
}
5. 定价模型(续)
5.1 Per-Token 定价
Per-Token 定价针对 LLM 和 Embedding 模型,按输入和输出 token 数量分别计费:
| 参数 | 类型 | 描述 | 示例 |
|---|---|---|---|
| input_price | Uint128 | 每输入 token 价格(uamp) | 30000000000000000 (0.03 msg) |
| output_price | Uint128 | 每输出 token 价格(uamp) | 60000000000000000 (0.06 msg) |
| free_tokens | Uint128 | 免费 token 配额 | 1000 |
5.1.1 计费公式
费用 = max(0, input_tokens - free_tokens) × input_price + output_tokens × output_price
5.2 Per-Request 定价
固定价格每次推断请求,适用于分类、目标检测等确定性模型:
PricingModel::PerRequest {
price: Uint128::new(100_000_000_000_000_000), // 0.1 msg per request
}
5.2.1 适用场景
| 场景 | 典型价格 | 说明 |
|---|---|---|
| 文本分类 | 0.01 msg | 情感分析、标签分类 |
| 目标检测 | 0.05 msg | 每张图像的检测 |
| 翻译 | 0.005 msg/千字 | 按文本量固定 |
5.3 订阅定价
提供时间维度订阅,用户支付固定速率获得免费配额,超出部分按单价计费:
// 每区块 0.001 msg,每区块免 10 次,超出每次 0.01 msg
PricingModel::Subscription {
rate_per_block: Uint128::new(1_000_000_000_000_000),
free_requests_per_block: 10,
overage_price: Uint128::new(10_000_000_000_000_000),
}
5.3.1 计费公式
费用 = rate_per_block + max(0, total_requests - free_requests_per_block) × overage_price
5.3.2 订阅示例
| 订阅级别 | 费率(/区块) | 免费配额 | 超额单价 | 典型月费 |
|---|---|---|---|---|
| 基础版 | 0.001 msg | 10次 | 0.01 msg | ~432 msg |
| 专业版 | 0.005 msg | 50次 | 0.008 msg | ~2160 msg |
| 企业版 | 0.02 msg | 200次 | 0.005 msg | ~8640 msg |
5.4 分层定价
提供批量折扣,按使用量区间设置不同单价:
| 区间 | 单价 |
|---|---|
| 0-1000 tokens | 0.05 msg / 1K tokens |
| 1001-10000 tokens | 0.04 msg / 1K tokens |
| 10001+ tokens | 0.03 msg / 1K tokens |
// TypeScript 分层定价选择器
function findTierPrice(
tiers: Array<{min_quantity: string; max_quantity?: string; unit_price: string}>,
totalUnits: bigint
): bigint {
for (const tier of tiers) {
const min = BigInt(tier.min_quantity);
const max = tier.max_quantity ? BigInt(tier.max_quantity) : null;
if (totalUnits >= min && (!max || totalUnits < max)) {
return totalUnits * BigInt(tier.unit_price);
}
}
return totalUnits * BigInt(tiers[tiers.length - 1].unit_price);
}
5.5 拍卖定价
针对稀缺 GPU 计算资源,实现密封投标拍卖机制:
// file: contracts/model-market/src/auction.rs
use cosmwasm_std::{Addr, Coin, Decimal, DepsMut, Env, MessageInfo,
Response, StdResult, Storage, Timestamp, Uint128};
use cw_storage_plus::Map;
use sha2::{Digest, Sha256};
use crate::error::ContractError;
use crate::state::CONFIG;
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct Auction {
pub auction_id: String,
pub service_id: String,
pub provider: Addr,
pub start_time: Timestamp,
pub end_time: Timestamp,
pub min_bid: Coin,
pub reserve_price: Coin,
pub status: AuctionStatus,
pub winning_bid: Option<Bid>,
pub total_slots: u32,
pub available_slots: u32,
}
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct Bid {
pub bidder: Addr,
pub amount: Coin,
pub bid_hash: String,
pub reveal_time: Timestamp,
pub slots_requested: u32,
}
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub enum AuctionStatus {
NotStarted,
Bidding,
Revealing,
Settled,
Cancelled,
}
pub const AUCTIONS: Map<&str, Auction> = Map::new("auctions");
pub const BIDS: Map<(&str, &Addr), Bid> = Map::new("bids");
pub fn execute_create_auction(
deps: DepsMut,
env: Env,
info: MessageInfo,
service_id: String,
duration_seconds: u64,
min_bid: Coin,
reserve_price: Coin,
total_slots: u32,
) -> Result<Response, ContractError> {
let auction_id = format!("auc-{}-{}", service_id, env.block.height);
let auction = Auction {
auction_id: auction_id.clone(),
service_id,
provider: info.sender,
start_time: env.block.time,
end_time: env.block.time.plus_seconds(duration_seconds),
min_bid,
reserve_price,
status: AuctionStatus::Bidding,
winning_bid: None,
total_slots,
available_slots: total_slots,
};
AUCTIONS.save(deps.storage, &auction_id, &auction)?;
Ok(Response::new()
.add_attribute("action", "create_auction")
.add_attribute("auction_id", &auction_id))
}
pub fn execute_place_bid(
deps: DepsMut,
env: Env,
info: MessageInfo,
auction_id: String,
amount: Coin,
bid_hash: String,
slots_requested: u32,
) -> Result<Response, ContractError> {
let mut auction = AUCTIONS.load(deps.storage, &auction_id)?;
if auction.status != AuctionStatus::Bidding {
return Err(ContractError::AuctionNotActive {});
}
if env.block.time > auction.end_time {
return Err(ContractError::AuctionEnded {});
}
if amount.amount < auction.min_bid.amount {
return Err(ContractError::BidTooLow {});
}
if slots_requested > auction.available_slots {
return Err(ContractError::NotEnoughSlots {});
}
let bid = Bid {
bidder: info.sender,
amount,
bid_hash,
reveal_time: Timestamp::from_nanos(0),
slots_requested,
};
BIDS.save(deps.storage, (&auction_id, &info.sender), &bid)?;
Ok(Response::new()
.add_attribute("action", "bid_placed")
.add_attribute("auction_id", &auction_id)
.add_attribute("bidder", &info.sender.to_string()))
}
pub fn execute_reveal_bid(
deps: DepsMut,
env: Env,
info: MessageInfo,
auction_id: String,
amount: Coin,
nonce: String,
) -> Result<Response, ContractError> {
let mut auction = AUCTIONS.load(deps.storage, &auction_id)?;
if auction.status != AuctionStatus::Bidding {
auction.status = AuctionStatus::Revealing;
AUCTIONS.save(deps.storage, &auction_id, &auction)?;
}
let bid = BIDS.load(deps.storage, (&auction_id, &info.sender))?;
let computed_hash = hash_string(&format!("{}{}{}", amount.amount, amount.denom, nonce));
if computed_hash != bid.bid_hash {
return Err(ContractError::BidHashMismatch {});
}
let revealed = Bid {
bidder: info.sender,
amount,
bid_hash: bid.bid_hash,
reveal_time: env.block.time,
slots_requested: bid.slots_requested,
};
BIDS.save(deps.storage, (&auction_id, &info.sender), &revealed)?;
if auction.winning_bid.as_ref().map_or(true, |w| amount.amount > w.amount.amount) {
if amount.amount >= auction.reserve_price.amount {
auction.winning_bid = Some(revealed.clone());
}
}
AUCTIONS.save(deps.storage, &auction_id, &auction)?;
Ok(Response::new()
.add_attribute("action", "bid_revealed")
.add_attribute("auction_id", &auction_id)
.add_attribute("amount", &amount.amount.to_string()))
}
pub fn execute_settle_auction(
deps: DepsMut,
env: Env,
info: MessageInfo,
auction_id: String,
) -> Result<Response, ContractError> {
let mut auction = AUCTIONS.load(deps.storage, &auction_id)?;
if env.block.time < auction.end_time.plus_seconds(300) {
return Err(ContractError::AuctionNotEnded {});
}
if auction.status != AuctionStatus::Revealing && auction.status != AuctionStatus::Bidding {
return Err(ContractError::InvalidAuctionStatus {});
}
auction.status = AuctionStatus::Settled;
AUCTIONS.save(deps.storage, &auction_id, &auction)?;
Ok(Response::new()
.add_attribute("action", "auction_settled")
.add_attribute("auction_id", &auction_id)
.add_attribute("winner", auction.winning_bid
.as_ref().map_or("none".to_string(), |b| b.bidder.to_string())))
}
fn hash_string(input: &str) -> String {
let mut h = Sha256::new();
h.update(input.as_bytes());
format!("{:x}", h.finalize())
}
5.6 动态定价
根据实时需求调整价格,使用指数移动平均(EMA)算法:
// file: contracts/model-market/src/dynamic_pricing.rs
use cosmwasm_std::{Decimal, StdResult, Storage, Timestamp, Uint128};
use cw_storage_plus::Map;
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct DynamicPriceState {
pub service_id: String,
pub base_price: Uint128,
pub current_price: Uint128,
pub demand_factor: Decimal,
pub last_update: Timestamp,
pub request_count_24h: u64,
pub peak_multiplier: Decimal,
}
pub const DYNAMIC_PRICES: Map<&str, DynamicPriceState> = Map::new("dyn_prices");
const EMA_ALPHA: Decimal = Decimal::permille(100);
const TARGET_UTILIZATION: Decimal = Decimal::percent(70);
const PRICE_ELASTICITY: Decimal = Decimal::permille(500);
pub fn update_dynamic_price(
storage: &mut dyn Storage,
service_id: &str,
current_demand: u64,
max_capacity: u64,
current_time: Timestamp,
) -> StdResult<DynamicPriceState> {
let mut state = DYNAMIC_PRICES
.may_load(storage, service_id)?
.unwrap_or(DynamicPriceState {
service_id: service_id.to_string(),
base_price: Uint128::new(50_000_000_000_000_000_000u128),
current_price: Uint128::new(50_000_000_000_000_000_000u128),
demand_factor: Decimal::one(),
last_update: current_time,
request_count_24h: 0,
peak_multiplier: Decimal::one(),
});
let utilization = if max_capacity > 0 {
Decimal::from_ratio(Uint128::from(current_demand), Uint128::from(max_capacity))
} else {
Decimal::zero()
};
let demand_ratio = if utilization > Decimal::zero() {
utilization / TARGET_UTILIZATION
} else {
Decimal::one()
};
let utilization_diff = utilization - TARGET_UTILIZATION;
let price_adjustment = Decimal::one() + PRICE_ELASTICITY * utilization_diff;
let new_price = state.base_price * price_adjustment;
let smoothed_price = state.current_price * (Decimal::one() - EMA_ALPHA) + new_price * EMA_ALPHA;
let min_price = state.base_price * Decimal::permille(500);
let max_price = state.base_price * Uint128::from(3u128);
let final_price = smoothed_price.max(min_price).min(max_price);
state.current_price = final_price;
state.demand_factor = demand_ratio;
state.last_update = current_time;
state.request_count_24h = state.request_count_24h.saturating_add(1);
DYNAMIC_PRICES.save(storage, service_id, &state)?;
Ok(state)
}
5.6.1 动态定价公式
利用率 = 当前需求 / 最大容量
价格调整 = 弹性系数 × (利用率 - 目标利用率)
新价格 = 基础价格 × (1 + 价格调整)
平滑价格 = 之前价格 × (1 - α) + 新价格 × α
最终价格 = clamp(平滑价格, 基础价格 × 0.5, 基础价格 × 3.0)
5.7 定价数学完整 Rust 代码
// file: contracts/model-market/src/pricing_math.rs
use cosmwasm_std::{Decimal, StdResult, Uint128};
use crate::state::{PriceTier, PricingModel, DYNAMIC_PRICES};
pub fn calculate_cost(
pricing: &PricingModel,
in_tokens: u64,
out_tokens: u64,
duration_secs: Option<u64>,
quantity: Option<u64>,
) -> StdResult<Uint128> {
match pricing {
PricingModel::PerToken { input_price, output_price, free_tokens } => {
let free = Uint128::from(*free_tokens);
let chargeable = if Uint128::from(in_tokens) > free {
Uint128::from(in_tokens) - free
} else {
Uint128::zero()
};
Ok(chargeable * input_price + Uint128::from(out_tokens) * output_price)
}
PricingModel::PerRequest { price } => Ok(*price),
PricingModel::PerDuration { price_per_second, min_duration } => {
let dur = duration_secs.unwrap_or(*min_duration).max(*min_duration);
Ok(Uint128::from(dur) * price_per_second)
}
PricingModel::PerOutputUnit { price_per_unit, .. } => {
Ok(Uint128::from(quantity.unwrap_or(1)) * price_per_unit)
}
PricingModel::Subscription { rate_per_block, free_requests_per_block, overage_price } => {
let total = quantity.unwrap_or(1);
let overage = if total > *free_requests_per_block {
Uint128::from((total - *free_requests_per_block) as u128) * overage_price
} else {
Uint128::zero()
};
Ok(*rate_per_block + overage)
}
PricingModel::Tiered { tiers } => {
let total = Uint128::from(in_tokens + out_tokens);
for t in tiers {
if total >= t.min_quantity && t.max_quantity.map_or(true, |m| total < m) {
return Ok(total * t.unit_price);
}
}
Ok(total * tiers.last().map_or(Uint128::zero(), |t| t.unit_price))
}
}
}
pub fn calculate_final_cost(
pricing: &PricingModel,
in_tokens: u64,
out_tokens: u64,
duration_secs: Option<u64>,
quantity: Option<u64>,
discount_bps: Option<u64>,
) -> StdResult<Uint128> {
let base_cost = calculate_cost(pricing, in_tokens, out_tokens, duration_secs, quantity)?;
if let Some(discount) = discount_bps {
if discount > 0 && discount <= 10000 {
let discount_factor = Decimal::permille(discount as u64);
let discount_amount = base_cost * discount_factor;
return Ok(base_cost - discount_amount);
}
}
Ok(base_cost)
}
pub fn compare_prices(
pricing_a: &PricingModel,
pricing_b: &PricingModel,
typical_input: u64,
typical_output: u64,
) -> (Uint128, Uint128, i8) {
let cost_a = calculate_cost(pricing_a, typical_input, typical_output, None, None)
.unwrap_or(Uint128::MAX);
let cost_b = calculate_cost(pricing_b, typical_input, typical_output, None, None)
.unwrap_or(Uint128::MAX);
let cheaper = if cost_a < cost_b {
-1i8
} else if cost_a > cost_b {
1i8
} else {
0i8
};
(cost_a, cost_b, cheaper)
}
pub fn batch_cost(
pricing: &PricingModel,
requests: &[(u64, u64)],
) -> StdResult<Vec<Uint128>> {
requests
.iter()
.map(|(in_t, out_t)| calculate_cost(pricing, *in_t, *out_t, None, None))
.collect()
}
pub fn format_price(amount: Uint128, denom: &str, decimals: u8) -> String {
let divisor = Uint128::from(10u128.pow(decimals as u32));
let integer = amount / divisor;
let fraction = amount % divisor;
let frac_str = format!("{:0width$}", fraction.u128(), width = decimals as usize);
format!("{}.{} {}", integer, &frac_str[..6], denom)
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn test_per_token_pricing() {
let pricing = PricingModel::PerToken {
input_price: Uint128::new(30_000_000_000_000_000),
output_price: Uint128::new(60_000_000_000_000_000),
free_tokens: Uint128::new(1000),
};
let cost = calculate_cost(&pricing, 500, 200, None, None).unwrap();
assert_eq!(cost, Uint128::new(200 * 60_000_000_000_000_000u128));
}
#[test]
fn test_tiered_pricing() {
let pricing = PricingModel::Tiered {
tiers: vec![
PriceTier {
min_quantity: Uint128::zero(),
max_quantity: Some(Uint128::new(1000)),
unit_price: Uint128::new(50_000_000_000_000_000),
},
PriceTier {
min_quantity: Uint128::new(1000),
max_quantity: Some(Uint128::new(10000)),
unit_price: Uint128::new(40_000_000_000_000_000),
},
PriceTier {
min_quantity: Uint128::new(10000),
max_quantity: None,
unit_price: Uint128::new(30_000_000_000_000_000),
},
],
};
let cost = calculate_cost(&pricing, 5000, 0, None, None).unwrap();
assert_eq!(cost, Uint128::new(5000 * 40_000_000_000_000_000u128));
}
#[test]
fn test_subscription_pricing() {
let pricing = PricingModel::Subscription {
rate_per_block: Uint128::new(1_000_000_000_000_000),
free_requests_per_block: 10,
overage_price: Uint128::new(10_000_000_000_000_000),
};
let cost = calculate_cost(&pricing, 0, 0, None, Some(15)).unwrap();
let expected = Uint128::new(1_000_000_000_000_000)
+ Uint128::new(5 * 10_000_000_000_000_000u128);
assert_eq!(cost, expected);
}
}
5.8 TypeScript 定价选择 UI 组件
// src/ui/PricingSelector.tsx
import React, { useState, useMemo } from "react";
export type PricingModelType = "per_token" | "per_request" | "subscription" | "tiered" | "auction";
export interface PricingConfig {
modelType: PricingModelType;
inputPrice?: string;
outputPrice?: string;
freeTokens?: string;
flatPrice?: string;
ratePerBlock?: string;
freePerBlock?: string;
overagePrice?: string;
tiers?: TierConfig[];
}
export interface TierConfig {
minQuantity: string;
maxQuantity?: string;
unitPrice: string;
}
interface PricingSelectorProps {
initialConfig?: PricingConfig;
onChange: (config: PricingConfig) => void;
denom?: string;
}
const PRICE_DECIMALS = 18;
const DENOM = "uamp";
function formatMsg(amount: string): string {
const val = BigInt(amount);
const divisor = BigInt(10 ** PRICE_DECIMALS);
const integer = val / divisor;
const fraction = val % divisor;
const fracStr = fraction.toString().padStart(PRICE_DECIMALS, "0").slice(0, 6);
return `${integer}.${fracStr} msg`;
}
export const PricingSelector: React.FC<PricingSelectorProps> = ({
initialConfig, onChange, denom = DENOM,
}) => {
const [modelType, setModelType] = useState<PricingModelType>(
initialConfig?.modelType ?? "per_token"
);
const [inputPrice, setInputPrice] = useState(initialConfig?.inputPrice ?? "30000000000000000");
const [outputPrice, setOutputPrice] = useState(initialConfig?.outputPrice ?? "60000000000000000");
const [freeTokens, setFreeTokens] = useState(initialConfig?.freeTokens ?? "1000");
const [flatPrice, setFlatPrice] = useState(initialConfig?.flatPrice ?? "100000000000000000");
const [ratePerBlock, setRatePerBlock] = useState(initialConfig?.ratePerBlock ?? "1000000000000000");
const [freePerBlock, setFreePerBlock] = useState(initialConfig?.freePerBlock ?? "10");
const [overagePrice, setOveragePrice] = useState(initialConfig?.overagePrice ?? "10000000000000000");
const [tiers, setTiers] = useState<TierConfig[]>(
initialConfig?.tiers ?? [
{ minQuantity: "0", maxQuantity: "1000", unitPrice: "50000000000000000" },
{ minQuantity: "1000", maxQuantity: "10000", unitPrice: "40000000000000000" },
{ minQuantity: "10000", unitPrice: "30000000000000000" },
]
);
const computeExampleCost = useMemo(() => {
const inputTokens = 500;
const outputTokens = 200;
switch (modelType) {
case "per_token": {
const free = BigInt(freeTokens);
const inp = BigInt(inputTokens);
const chargeable = inp > free ? inp - free : BigInt(0);
return chargeable * BigInt(inputPrice) + BigInt(outputTokens) * BigInt(outputPrice);
}
case "per_request": return BigInt(flatPrice);
case "subscription": {
const total = 15;
const freeCount = parseInt(freePerBlock);
const overage = total > freeCount ? BigInt(total - freeCount) * BigInt(overagePrice) : BigInt(0);
return BigInt(ratePerBlock) + overage;
}
case "tiered": {
const total = BigInt(inputTokens + outputTokens);
for (const t of tiers) {
const min = BigInt(t.minQuantity);
const max = t.maxQuantity ? BigInt(t.maxQuantity) : null;
if (total >= min && (!max || total < max)) return total * BigInt(t.unitPrice);
}
return total * BigInt(tiers[tiers.length - 1].unitPrice);
}
default: return BigInt(0);
}
}, [modelType, inputPrice, outputPrice, freeTokens, flatPrice, ratePerBlock, freePerBlock, overagePrice, tiers]);
const buildConfig = (): PricingConfig => {
switch (modelType) {
case "per_token": return { modelType, inputPrice, outputPrice, freeTokens };
case "per_request": return { modelType, flatPrice };
case "subscription": return { modelType, ratePerBlock, freePerBlock, overagePrice };
case "tiered": return { modelType, tiers };
default: return { modelType };
}
};
const handleChange = () => onChange(buildConfig());
const addTier = () => {
const last = tiers[tiers.length - 1];
const newMin = last ? BigInt(last.maxQuantity ?? last.minQuantity) + BigInt(1000) : BigInt(0);
setTiers([...tiers, { minQuantity: newMin.toString(), unitPrice: "30000000000000000" }]);
};
const updateTier = (index: number, field: keyof TierConfig, value: string) => {
const updated = [...tiers];
updated[index] = { ...updated[index], [field]: value };
setTiers(updated);
};
const removeTier = (index: number) => {
if (tiers.length > 1) setTiers(tiers.filter((_, i) => i !== index));
};
return (
<div className="pricing-selector">
<h3>定价模型配置</h3>
<div className="form-group">
<label>定价类型</label>
<select value={modelType} onChange={e => {
setModelType(e.target.value as PricingModelType);
setTimeout(handleChange, 0);
}}>
<option value="per_token">按 Token 计费</option>
<option value="per_request">按请求计费</option>
<option value="subscription">订阅计费</option>
<option value="tiered">分层计费</option>
<option value="auction">拍卖计费</option>
</select>
</div>
{modelType === "per_token" && (
<>
<div className="form-group">
<label>输入价格 ({denom}/token):</label>
<input type="text" value={inputPrice} onChange={e => setInputPrice(e.target.value)} onBlur={handleChange} />
<small className="hint">{formatMsg(inputPrice)} / token</small>
</div>
<div className="form-group">
<label>输出价格 ({denom}/token):</label>
<input type="text" value={outputPrice} onChange={e => setOutputPrice(e.target.value)} onBlur={handleChange} />
<small className="hint">{formatMsg(outputPrice)} / token</small>
</div>
<div className="form-group">
<label>免费 token 数:</label>
<input type="number" value={freeTokens} onChange={e => setFreeTokens(e.target.value)} onBlur={handleChange} />
</div>
</>
)}
{modelType === "per_request" && (
<div className="form-group">
<label>固定价格 ({denom}):</label>
<input type="text" value={flatPrice} onChange={e => setFlatPrice(e.target.value)} onBlur={handleChange} />
<small className="hint">{formatMsg(flatPrice)} / 请求</small>
</div>
)}
{modelType === "subscription" && (
<>
<div className="form-group">
<label>每区块费率 ({denom}):</label>
<input type="text" value={ratePerBlock} onChange={e => setRatePerBlock(e.target.value)} onBlur={handleChange} />
</div>
<div className="form-group">
<label>每区块免费请求数:</label>
<input type="number" value={freePerBlock} onChange={e => setFreePerBlock(e.target.value)} onBlur={handleChange} />
</div>
<div className="form-group">
<label>超额单价 ({denom}):</label>
<input type="text" value={overagePrice} onChange={e => setOveragePrice(e.target.value)} onBlur={handleChange} />
</div>
</>
)}
{modelType === "tiered" && (
<div className="tier-list">
<h4>价格分层</h4>
{tiers.map((tier, i) => (
<div key={i} className="tier-row">
<input type="text" placeholder="最低量" value={tier.minQuantity}
onChange={e => updateTier(i, "minQuantity", e.target.value)} onBlur={handleChange} />
<input type="text" placeholder="最高量(选填)" value={tier.maxQuantity ?? ""}
onChange={e => updateTier(i, "maxQuantity", e.target.value || undefined)} onBlur={handleChange} />
<input type="text" placeholder="单价" value={tier.unitPrice}
onChange={e => updateTier(i, "unitPrice", e.target.value)} onBlur={handleChange} />
<button onClick={() => removeTier(i)} disabled={tiers.length <= 1}>删除</button>
</div>
))}
<button onClick={addTier}>添加分层</button>
</div>
)}
<div className="cost-preview">
<strong>示例费用 (500 in / 200 out): </strong>
<span className="cost-value">{formatMsg(computeExampleCost.toString())}</span>
</div>
</div>
);
};
export default PricingSelector;
6. 服务级别协议 SLA(续)
6.1 SLA 条款结构
SLA 定义了提供者承诺的服务质量标准,包含三个核心维度:
| 维度 | 指标 | 典型值 | 说明 |
|---|---|---|---|
| 延迟 | P50 / P95 / P99 | 500ms / 2000ms / 5000ms | 请求响应时间的百分位承诺 |
| 可用性 | Uptime % | 99.9% | 成功请求占总请求的百分比 |
| 质量 | Quality Score | 85/100 | 结果质量的综合评分 |
6.1.1 完整 SLA 数据结构
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct SlaTerms {
pub latency: LatencyCommitment,
pub availability: AvailabilityCommitment,
pub quality: QualityCommitment,
pub compensation: CompensationScheme,
pub observation_window: u64,
pub min_samples: u32,
}
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct LatencyCommitment {
pub p50_max_ms: u64,
pub p95_max_ms: u64,
pub p99_max_ms: u64,
}
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct AvailabilityCommitment {
pub uptime_percentage: Decimal,
pub max_consecutive_failures: u32,
}
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct QualityCommitment {
pub min_quality_score: u8,
pub min_semantic_similarity: Option<Decimal>,
}
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct CompensationScheme {
pub latency_violation_refund_bps: u64,
pub availability_violation_refund_bps: u64,
pub quality_violation_refund_bps: u64,
pub max_refund_bps: u64,
}
6.2 SLA 链上锚定
SLA 文档以 IPFS 或 Arweave 存储,链上仅保存其 hash:
// file: contracts/model-market/src/sla_anchor.rs
use cosmwasm_std::{DepsMut, Env, MessageInfo, Response, StdResult, Timestamp};
use sha2::{Digest, Sha256};
use cw_storage_plus::Map;
use crate::state::ServiceLevelAgreement;
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct SlaAnchor {
pub service_id: String,
pub sla_hash: String,
pub document_uri: String,
pub committed_at: Timestamp,
pub version: u32,
pub previous_hash: Option<String>,
}
pub const SLA_ANCHORS: Map<&str, Vec<SlaAnchor>> = Map::new("sla_anchors");
pub fn anchor_sla(
deps: DepsMut,
env: Env,
service_id: String,
sla: &ServiceLevelAgreement,
document_uri: String,
previous_hash: Option<String>,
) -> StdResult<SlaAnchor> {
let sla_json = serde_json::to_string(sla)?;
let mut hasher = Sha256::new();
hasher.update(sla_json.as_bytes());
if let Some(prev) = &previous_hash {
hasher.update(prev.as_bytes());
}
let sla_hash = format!("{:x}", hasher.finalize());
let anchor = SlaAnchor {
service_id: service_id.clone(),
sla_hash,
document_uri,
committed_at: env.block.time,
version: 1,
previous_hash,
};
let mut history = SLA_ANCHORS
.may_load(deps.storage, &service_id)?
.unwrap_or_default();
history.push(anchor.clone());
SLA_ANCHORS.save(deps.storage, &service_id, &history)?;
Ok(anchor)
}
pub fn verify_sla_anchor(
sla: &ServiceLevelAgreement,
document_uri: &str,
expected_hash: &str,
) -> bool {
let sla_json = serde_json::to_string(sla).unwrap_or_default();
let mut hasher = Sha256::new();
hasher.update(sla_json.as_bytes());
hasher.update(document_uri.as_bytes());
let computed_hash = format!("{:x}", hasher.finalize());
computed_hash == expected_hash
}
6.2.1 锚定流程
1. 提供者编写 SLA 文档 (JSON/YAML)
2. 上传文档到 IPFS → 获取 document_uri
3. 调用 anchor_sla() 将 SLA hash 和 URI 提交上链
4. 消费者通过 document_uri 获取完整 SLA 文档
5. 使用 verify_sla_anchor() 验证文档一致性
6.3 SLA 违规检测
结合链上数据与链下预言机的混合检测方案:
// file: contracts/model-market/src/sla_violation.rs
use cosmwasm_std::{Decimal, DepsMut, Env, MessageInfo, Response, StdResult, Timestamp, Uint128};
use crate::state::{CONFIG, SERVICES};
use crate::error::ContractError;
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct ViolationReport {
pub service_id: String,
pub report_type: ViolationReportType,
pub observed_value: String,
pub threshold_value: String,
pub timestamp: Timestamp,
pub reporter: String,
pub evidence_uri: Option<String>,
}
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub enum ViolationReportType {
LatencyViolation,
UptimeViolation,
QualityViolation,
ConsecutiveFailure,
}
pub const VIOLATION_REPORTS: Map<&str, Vec<ViolationReport>> = Map::new("violation_reports");
pub fn submit_violation_report(
deps: DepsMut,
env: Env,
info: MessageInfo,
report: ViolationReport,
) -> Result<Response, ContractError> {
let config = CONFIG.load(deps.storage)?;
let service = SERVICES.load(deps.storage, &report.service_id)?;
let is_authorized = info.sender == config.owner
|| config.arbitrators.contains(&info.sender);
if !is_authorized {
return Err(ContractError::Unauthorized {});
}
let mut reports = VIOLATION_REPORTS
.may_load(deps.storage, &report.service_id)?
.unwrap_or_default();
reports.push(report.clone());
VIOLATION_REPORTS.save(deps.storage, &report.service_id, &reports)?;
Ok(Response::new()
.add_attribute("action", "violation_reported")
.add_attribute("service_id", &report.service_id)
.add_attribute("type", format!("{:?}", report.report_type)))
}
pub fn check_and_trigger_auto_compensation(
deps: DepsMut,
env: Env,
service_id: &str,
) -> Result<Response, ContractError> {
let service = SERVICES.load(deps.storage, service_id)?;
let reports = VIOLATION_REPORTS
.may_load(deps.storage, service_id)?
.unwrap_or_default();
let observation_window = service.sla.observation_window_seconds;
let cutoff = env.block.time.minus_seconds(observation_window);
let recent_violations: Vec<&ViolationReport> = reports
.iter()
.filter(|r| r.timestamp > cutoff)
.collect();
if recent_violations.is_empty() {
return Ok(Response::new().add_attribute("action", "no_violations"));
}
let latency_violations = recent_violations.iter()
.filter(|r| matches!(r.report_type, ViolationReportType::LatencyViolation))
.count();
let uptime_violations = recent_violations.iter()
.filter(|r| matches!(r.report_type, ViolationReportType::UptimeViolation))
.count();
let quality_violations = recent_violations.iter()
.filter(|r| matches!(r.report_type, ViolationReportType::QualityViolation))
.count();
let consecutive_failures = recent_violations.iter()
.filter(|r| matches!(r.report_type, ViolationReportType::ConsecutiveFailure))
.count();
let mut total_refund_bps: u64 = 0;
if latency_violations > 3 {
total_refund_bps += service.sla.compensation.latency_violation_refund_bps * latency_violations as u64;
}
if uptime_violations > 0 || consecutive_failures > 2 {
total_refund_bps += service.sla.compensation.availability_violation_refund_bps;
}
if quality_violations > 0 {
total_refund_bps += service.sla.compensation.quality_violation_refund_bps * quality_violations as u64;
}
let capped_refund = total_refund_bps.min(service.sla.compensation.max_refund_bps);
Ok(Response::new()
.add_attribute("action", "auto_compensation_triggered")
.add_attribute("service_id", service_id)
.add_attribute("violations", &recent_violations.len().to_string())
.add_attribute("refund_bps", &capped_refund.to_string()))
}
6.4 罚金计算与自动退款
// file: contracts/model-market/src/sla_compensation.rs
use cosmwasm_std::{Coin, Decimal, DepsMut, Env, MessageInfo,
Response, Uint128, BankMsg, CosmosMsg};
use crate::error::ContractError;
use crate::state::{CONFIG, SERVICES, ServiceLevelAgreement};
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct CompensationRequest {
pub consumer: String,
pub service_id: String,
pub violation_window_start: u64,
pub violation_window_end: u64,
pub total_fees_paid: Uint128,
pub refund_amount: Uint128,
pub reason: String,
}
pub fn calculate_sla_refund(
sla_terms: &ServiceLevelAgreement,
violations: &[&str],
total_fees: Uint128,
) -> Uint128 {
if violations.is_empty() {
return Uint128::zero();
}
let mut total_bps: u64 = 0;
for v in violations {
match *v {
"latency_p50" | "latency_p95" | "latency_p99" => {
total_bps += sla_terms.compensation.latency_violation_refund_bps;
}
"uptime" | "consecutive_failure" => {
total_bps += sla_terms.compensation.availability_violation_refund_bps;
}
"quality" => {
total_bps += sla_terms.compensation.quality_violation_refund_bps;
}
_ => {}
}
}
let capped_bps = total_bps.min(sla_terms.compensation.max_refund_bps);
total_fees * Uint128::from(capped_bps) / Uint128::from(10000u64)
}
pub fn execute_process_refund(
deps: DepsMut,
env: Env,
info: MessageInfo,
service_id: String,
consumer: String,
refund_amount: Uint128,
) -> Result<Response, ContractError> {
let config = CONFIG.load(deps.storage)?;
let service = SERVICES.load(deps.storage, &service_id)?;
if info.sender != service.provider
&& info.sender != config.owner
&& !config.arbitrators.contains(&info.sender)
{
return Err(ContractError::Unauthorized {});
}
let refund_msg: CosmosMsg = BankMsg::Send {
to_address: consumer,
amount: vec![Coin {
denom: "uamp".to_string(),
amount: refund_amount,
}],
}.into();
Ok(Response::new()
.add_message(refund_msg)
.add_attribute("action", "sla_refund_processed")
.add_attribute("service_id", &service_id)
.add_attribute("consumer", &consumer)
.add_attribute("refund", &refund_amount.to_string()))
}
6.4.1 补偿计算示例
| 违规类型 | 违规次数 | 每次补偿 (bps) | 总补偿 (bps) |
|---|---|---|---|
| P99 延迟超标 | 1 | 1000 (10%) | 1000 |
| 可用性低于 99.9% | 1 | 2000 (20%) | 2000 |
| 质量分低于 85 | 2 | 500 (5%) | 1000 |
| 合计 | 4000 (40%) |
6.5 完整 Rust SLA 合约代码
// file: contracts/model-market/src/sla_contract.rs
use cosmwasm_std::{
entry_point, to_json_binary, Binary, Decimal, Deps, DepsMut, Env,
MessageInfo, Response, StdResult, Timestamp, Uint128,
};
use crate::error::ContractError;
use crate::state::{
SERVICE_SLA_SNAPSHOTS, SERVICES, CONFIG, ServiceLevelAgreement,
CompensationScheme, SlaSnapshot,
};
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct SlaInitMsg {
pub service_id: String,
pub p95_latency_ms: u64,
pub p99_latency_ms: u64,
pub uptime_percentage: Decimal,
pub min_quality_score: u8,
pub observation_window_seconds: u64,
pub latency_refund_bps: u64,
pub availability_refund_bps: u64,
pub quality_refund_bps: u64,
pub max_refund_bps: u64,
pub document_uri: String,
}
pub fn execute_init_sla(
deps: DepsMut,
env: Env,
info: MessageInfo,
msg: SlaInitMsg,
) -> Result<Response, ContractError> {
let config = CONFIG.load(deps.storage)?;
let mut service = SERVICES.load(deps.storage, &msg.service_id)?;
if info.sender != service.provider && info.sender != config.owner {
return Err(ContractError::Unauthorized {});
}
service.sla = ServiceLevelAgreement {
p95_latency_ms: msg.p95_latency_ms,
p99_latency_ms: msg.p99_latency_ms,
uptime_percentage: msg.uptime_percentage,
min_quality_score: msg.min_quality_score,
observation_window_seconds: msg.observation_window_seconds,
compensation: CompensationScheme {
latency_violation_refund_bps: msg.latency_refund_bps,
availability_violation_refund_bps: msg.availability_refund_bps,
quality_violation_refund_bps: msg.quality_refund_bps,
max_refund_bps: msg.max_refund_bps,
},
};
SERVICES.save(deps.storage, &msg.service_id, &service)?;
let anchor = crate::sla_anchor::anchor_sla(
deps, env, msg.service_id.clone(),
&service.sla, msg.document_uri, None,
)?;
Ok(Response::new()
.add_attribute("action", "sla_initialized")
.add_attribute("service_id", &msg.service_id)
.add_attribute("sla_hash", &anchor.sla_hash))
}
pub fn update_sla_snapshot(
deps: DepsMut,
env: Env,
service_id: String,
total_requests: u64,
successful_requests: u64,
latency_p50: u64,
latency_p95: u64,
latency_p99: u64,
avg_quality_score: u8,
) -> StdResult<()> {
let snapshot = SlaSnapshot {
service_id: service_id.clone(),
window_start: env.block.time.minus_seconds(86400),
window_end: env.block.time,
total_requests,
successful_requests,
latency_p50,
latency_p95,
latency_p99,
avg_quality_score,
violations: vec![],
compensation_due: Uint128::zero(),
snapshot_hash: String::new(),
};
SERVICE_SLA_SNAPSHOTS.save(deps.storage, &service_id, &snapshot)?;
Ok(())
}
pub fn get_sla_compliance(
deps: Deps,
service_id: &str,
) -> StdResult<SlaComplianceReport> {
let service = SERVICES.load(deps.storage, service_id)?;
let snapshot = SERVICE_SLA_SNAPSHOTS
.may_load(deps.storage, service_id)?
.unwrap_or(SlaSnapshot {
service_id: service_id.to_string(),
window_start: Timestamp::from_nanos(0),
window_end: Timestamp::from_nanos(0),
total_requests: 0,
successful_requests: 0,
latency_p50: 0,
latency_p95: 0,
latency_p99: 0,
avg_quality_score: 0,
violations: vec![],
compensation_due: Uint128::zero(),
snapshot_hash: String::new(),
});
let mut violations = Vec::new();
if snapshot.latency_p95 > service.sla.p95_latency_ms {
violations.push("latency_p95");
}
if snapshot.latency_p99 > service.sla.p99_latency_ms {
violations.push("latency_p99");
}
let uptime = if snapshot.total_requests > 0 {
Decimal::from_ratio(
Uint128::from(snapshot.successful_requests),
Uint128::from(snapshot.total_requests),
) * Decimal::percent(100)
} else {
Decimal::percent(100)
};
if uptime < service.sla.uptime_percentage {
violations.push("uptime");
}
if snapshot.avg_quality_score < service.sla.min_quality_score {
violations.push("quality");
}
let total_fees = Uint128::from(snapshot.total_requests)
* Uint128::new(100_000_000_000_000_000u128);
let compensation = calculate_sla_refund(&service.sla, &violations, total_fees);
Ok(SlaComplianceReport {
service_id: service_id.to_string(),
compliant: violations.is_empty(),
violations: violations.iter().map(|v| v.to_string()).collect(),
compensation_due: compensation,
snapshot,
})
}
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct SlaComplianceReport {
pub service_id: String,
pub compliant: bool,
pub violations: Vec<String>,
pub compensation_due: Uint128,
pub snapshot: SlaSnapshot,
}
6.6 Python SLA 监控代理
# src/monitoring/sla_monitor.py
import asyncio, json, time, hashlib, logging
from collections import deque
from dataclasses import dataclass
from typing import Optional
from statistics import median
logger = logging.getLogger(__name__)
@dataclass
class LatencySample:
request_id: str
consumer: str
latency_ms: float
success: bool
quality_score: int
timestamp: float
class SlaMonitor:
"""SLA 监控代理 - 链下监控并提供链上证明"""
def __init__(self, ledger_client, contract_address: str,
observation_window: int = 86400):
self.ledger = ledger_client
self.contract = contract_address
self.window = observation_window
self.samples: dict[str, deque[LatencySample]] = {}
self.running = False
async def start(self, interval_seconds: int = 300):
self.running = True
while self.running:
for service_id in list(self.samples.keys()):
await self.evaluate_service(service_id)
await asyncio.sleep(interval_seconds)
def stop(self):
self.running = False
def record_sample(self, service_id: str, sample: LatencySample):
if service_id not in self.samples:
self.samples[service_id] = deque(maxlen=100000)
self.samples[service_id].append(sample)
async def evaluate_service(self, service_id: str) -> dict:
if service_id not in self.samples:
return {"status": "no_data"}
now = time.time()
cutoff = now - self.window
recent = [s for s in self.samples[service_id] if s.timestamp > cutoff]
if len(recent) < 10:
return {"status": "insufficient_samples", "count": len(recent)}
latencies = sorted([s.latency_ms for s in recent if s.success])
total = len(recent)
successes = len(latencies)
uptime = (successes / total * 100) if total > 0 else 100.0
p50 = median(latencies) if latencies else 0
p95 = latencies[int(len(latencies) * 0.95)] if len(latencies) > 20 else (latencies[-1] if latencies else 0)
p99 = latencies[int(len(latencies) * 0.99)] if len(latencies) > 100 else (latencies[-1] if latencies else 0)
quality_scores = [s.quality_score for s in recent if s.success]
avg_quality = sum(quality_scores) / len(quality_scores) if quality_scores else 0
report = {
"service_id": service_id,
"timestamp": int(now),
"observation_window": self.window,
"total_requests": total,
"successful_requests": successes,
"uptime_percentage": round(uptime, 2),
"latency_p50_ms": round(p50, 2),
"latency_p95_ms": round(p95, 2),
"latency_p99_ms": round(p99, 2),
"avg_quality_score": round(avg_quality, 1),
"consecutive_failures": self._count_consecutive_failures(recent),
}
report_hash = hashlib.sha256(
json.dumps(report, sort_keys=True).encode()
).hexdigest()
report["report_hash"] = report_hash
violations = await self._detect_violations(service_id, report)
report["violations"] = violations
if violations:
await self._submit_violation_report(service_id, report)
return report
def _count_consecutive_failures(self, samples: list) -> int:
max_consecutive = 0
current = 0
for s in sorted(samples, key=lambda x: x.timestamp):
if not s.success:
current += 1
max_consecutive = max(max_consecutive, current)
else:
current = 0
return max_consecutive
async def _detect_violations(self, service_id: str, report: dict) -> list[dict]:
violations = []
try:
service_data = self.ledger.query_contract(
self.contract,
{"get_service": {"model_id": service_id}}
)
sla = service_data.get("sla", {})
except Exception:
logger.warning(f"Cannot fetch SLA for {service_id}")
return violations
if report["latency_p95_ms"] > sla.get("p95_latency_ms", 999999):
violations.append({
"type": "latency_p95",
"observed": report["latency_p95_ms"],
"threshold": sla.get("p95_latency_ms"),
"severity": "warning",
})
if report["latency_p99_ms"] > sla.get("p99_latency_ms", 999999):
violations.append({
"type": "latency_p99",
"observed": report["latency_p99_ms"],
"threshold": sla.get("p99_latency_ms"),
"severity": "critical",
})
sla_uptime = float(sla.get("uptime_percentage", "99.9").rstrip("%"))
if report["uptime_percentage"] < sla_uptime:
violations.append({
"type": "uptime",
"observed": report["uptime_percentage"],
"threshold": sla_uptime,
"severity": "critical",
})
min_quality = sla.get("min_quality_score", 70)
if report["avg_quality_score"] < min_quality:
violations.append({
"type": "quality_score",
"observed": report["avg_quality_score"],
"threshold": min_quality,
"severity": "warning",
})
max_failures = sla.get("max_consecutive_failures", 5)
if report["consecutive_failures"] > max_failures:
violations.append({
"type": "consecutive_failure",
"observed": report["consecutive_failures"],
"threshold": max_failures,
"severity": "critical",
})
return violations
async def _submit_violation_report(self, service_id: str, report: dict):
try:
tx = self.ledger.execute_contract(
self.contract,
{
"submit_violation_report": {
"report": {
"service_id": service_id,
"report_type": "latency_violation"
if any(v["type"] == "latency_p99" for v in report["violations"])
else "uptime_violation",
"observed_value": str(report.get("latency_p99_ms", "")),
"threshold_value": "5000",
"timestamp": int(time.time()),
"reporter": str(self.ledger.wallet.address()),
"evidence_uri": f"ipfs://{report['report_hash']}",
}
}
},
self.ledger.wallet
)
logger.info(f"Violation report submitted: {tx.tx_hash}")
except Exception as e:
logger.error(f"Failed to submit violation report: {e}")
7. 完整示例:多模型推理协调
7.1 场景描述
用户需要执行一个复杂的多步骤 AI 流水线:
- Embedding (Agent A):将用户查询文本转为向量
- 语义搜索 (Agent B):基于向量在知识库中搜索相关内容
- LLM 生成 (Agent C):基于搜索结果生成最终回答
- 内容审核 (Agent D):审核生成内容是否合规
7.2 流水线数据流
┌──────────────┐
│ 用户输入文本 │
└──────┬───────┘
│
┌──────▼───────┐
│ Agent A │
│ Embedding │
│ 0.003 msg │
└──────┬───────┘
│ (向量)
┌──────▼───────┐
│ Agent B │
│ 语义搜索 │
│ 0.01 msg │
└──────┬───────┘
│ (上下文)
┌──────▼───────┐
│ Agent C │
│ LLM 生成 │
│ 0.05 msg │
└──────┬───────┘
│ (回复)
┌──────▼───────┐
│ Agent D │
│ 内容审核 │
│ 0.002 msg │
└──────┬───────┘
│ (审核通过)
┌──────▼───────┐
│ 最终输出 │
└──────────────┘
7.2.1 Step 到 Step 的数据格式说明
| 步骤 | 输入格式 | 输出格式 | 示例 |
|---|---|---|---|
| embedding | {"text": "用户查询"} |
{"embeddings": [[0.1, 0.2, ...]]} |
768维向量 |
| search | {"vector": [...], "query": "..."} |
{"results": [{"text": "...", "score": 0.95}]} |
Top-3 结果 |
| generate | {"context": "...", "question": "..."} |
{"text": "生成的回答..."} |
完整回复 |
| moderate | {"text": "待审核内容"} |
{"safe": true, "categories": {}} |
审核标签 |
7.3 Python 编排器 Agent
# src/orchestrator/pipeline_orchestrator.py
import asyncio, json, time, hashlib, logging
from dataclasses import dataclass, field
from enum import Enum
from typing import Optional, Callable, Any
from cosmpy.aerial.client import LedgerClient
from cosmpy.aerial.wallet import LocalWallet
from cosmpy.aerial.contract import CosmWasmClient
logger = logging.getLogger(__name__)
class PipelineStatus(Enum):
PENDING = "pending"
RUNNING = "running"
SUCCESS = "success"
FAILED = "failed"
ROLLED_BACK = "rolled_back"
class StepStatus(Enum):
PENDING = "pending"
DISCOVERING = "discovering"
EXECUTING = "executing"
VERIFYING = "verifying"
COMPLETED = "completed"
FAILED = "failed"
SKIPPED = "skipped"
@dataclass
class PipelineStep:
step_id: str
name: str
model_type: str
required_capabilities: list[str]
input_transform: Optional[Callable] = None
output_transform: Optional[Callable] = None
max_retries: int = 3
timeout_seconds: int = 60
max_price: str = "1000000000000000000"
min_quality: int = 70
fallback_step_id: Optional[str] = None
parallel_with: list[str] = field(default_factory=list)
@dataclass
class StepContext:
step: PipelineStep
provider: Optional[dict] = None
input_data: Optional[dict] = None
output_data: Optional[dict] = None
status: StepStatus = StepStatus.PENDING
error: Optional[str] = None
cost: int = 0
latency_ms: int = 0
retry_count: int = 0
attestation: Optional[dict] = None
@dataclass
class PipelineContext:
pipeline_id: str
steps: dict[str, StepContext]
status: PipelineStatus = PipelineStatus.PENDING
total_cost: int = 0
total_latency_ms: int = 0
start_time: float = 0.0
end_time: float = 0.0
error: Optional[str] = None
metadata: dict = field(default_factory=dict)
final_output: Optional[dict] = None
class PipelineOrchestrator:
def __init__(self, wallet: LocalWallet, ledger: LedgerClient,
contracts: dict, max_concurrent_steps: int = 3):
self.wallet = wallet
self.ledger = ledger
self.cosmwasm = CosmWasmClient(ledger)
self.contracts = contracts
self.semaphore = asyncio.Semaphore(max_concurrent_steps)
self.consumer = None
def _make_consumer(self):
from src.consumer.agent_client import ConsumerAgent
return ConsumerAgent(self.wallet, self.ledger, self.contracts)
async def run_pipeline(self, pipeline_id: str,
steps: list[PipelineStep],
initial_input: dict) -> PipelineContext:
ctx = PipelineContext(
pipeline_id=pipeline_id,
steps={s.step_id: StepContext(step=s) for s in steps},
start_time=time.time(),
)
logger.info(f"Starting pipeline {pipeline_id} with {len(steps)} steps")
try:
completed = set()
step_order = self._resolve_dependency_order(steps)
for step_id in step_order:
step_ctx = ctx.steps[step_id]
step_ctx.status = StepStatus.DISCOVERING
input_data = await self._prepare_step_input(
step_ctx, initial_input, ctx
)
step_ctx.input_data = input_data
provider = await self._discover_provider(step_ctx)
if not provider:
if step_ctx.step.fallback_step_id:
logger.info(f"Falling back to {step_ctx.step.fallback_step_id}")
fallback_ctx = ctx.steps[step_ctx.step.fallback_step_id]
provider = await self._discover_provider(fallback_ctx)
if not provider:
raise RuntimeError(f"No provider for {step_id}")
step_ctx.provider = provider
step_ctx.step = fallback_ctx.step
else:
raise RuntimeError(f"No provider found for step {step_id}")
step_ctx.provider = provider
await self._execute_step(step_ctx)
completed.add(step_id)
ctx.total_cost += step_ctx.cost
ctx.status = PipelineStatus.SUCCESS
ctx.end_time = time.time()
ctx.total_latency_ms = int((ctx.end_time - ctx.start_time) * 1000)
ctx.final_output = self._collect_final_output(ctx)
logger.info(f"Pipeline {pipeline_id} completed in {ctx.total_latency_ms}ms, "
f"cost: {ctx.total_cost} uamp")
except Exception as e:
ctx.status = PipelineStatus.FAILED
ctx.error = str(e)
ctx.end_time = time.time()
logger.error(f"Pipeline {pipeline_id} failed: {e}")
await self._rollback(ctx)
return ctx
def _resolve_dependency_order(self, steps: list[PipelineStep]) -> list[str]:
ordered = []
for step in steps:
if step.step_id not in ordered:
ordered.append(step.step_id)
return ordered
async def _prepare_step_input(self, step_ctx: StepContext,
initial_input: dict,
ctx: PipelineContext) -> dict:
step = step_ctx.step
for prev_id, prev_ctx in ctx.steps.items():
if prev_id == step.step_id:
continue
if prev_ctx.status == StepStatus.COMPLETED and prev_ctx.output_data:
if step.input_transform:
try:
return step.input_transform(prev_ctx.output_data)
except Exception as e:
logger.warning(f"Input transform failed: {e}")
if step.input_transform:
try:
return step.input_transform(initial_input)
except Exception as e:
logger.warning(f"Initial input transform failed: {e}")
return initial_input
async def _discover_provider(self, step_ctx: StepContext) -> Optional[dict]:
consumer = self._make_consumer()
models = consumer.discover(
model_type=step_ctx.step.model_type,
caps=step_ctx.step.required_capabilities,
)
if not models:
return None
scored = []
for m in models:
cost = consumer.estimate_cost(
m["pricing"],
len(json.dumps(step_ctx.input_data, ensure_ascii=False)) // 4,
1024,
)
if cost is None or cost > int(step_ctx.step.max_price):
continue
quality = m.get("reputation", {}).get("avg_quality", 50)
if quality < step_ctx.step.min_quality:
continue
score = quality / max(cost, 1)
scored.append((score, m))
if not scored:
return None
scored.sort(key=lambda x: x[0], reverse=True)
return scored[0][1]
async def _execute_step(self, step_ctx: StepContext):
step = step_ctx.step
consumer = self._make_consumer()
for attempt in range(step.max_retries):
step_ctx.retry_count = attempt
step_ctx.status = StepStatus.EXECUTING
try:
start = time.time()
prompt = json.dumps(step_ctx.input_data, ensure_ascii=False)
max_tokens = 2048
response = consumer.infer(
step_ctx.provider,
prompt,
max_tokens=max_tokens,
timeout=step.timeout_seconds,
)
elapsed = int((time.time() - start) * 1000)
step_ctx.latency_ms = elapsed
output = response.get("output", {})
usage = response.get("usage", {})
cost = consumer.estimate_cost(
step_ctx.provider["pricing"],
len(prompt) // 4,
usage.get("output_tokens", max_tokens),
) or 0
step_ctx.cost = cost
step_ctx.status = StepStatus.VERIFYING
if self._verify_step_output(step_ctx, response):
step_ctx.output_data = output
step_ctx.status = StepStatus.COMPLETED
step_ctx.attestation = {
"step_id": step.step_id,
"provider": step_ctx.provider["provider"],
"model_id": step_ctx.provider["model_id"],
"output_hash": hashlib.sha256(
json.dumps(output, sort_keys=True).encode()
).hexdigest(),
"cost": cost,
"latency_ms": elapsed,
"timestamp": int(time.time()),
}
return
except Exception as e:
logger.warning(f"Step {step.step_id} attempt {attempt + 1} failed: {e}")
step_ctx.error = str(e)
if attempt < step.max_retries - 1:
await asyncio.sleep(2 ** attempt)
step_ctx.status = StepStatus.FAILED
raise RuntimeError(
f"Step {step.step_id} failed after {step.max_retries} attempts: {step_ctx.error}"
)
def _verify_step_output(self, step_ctx: StepContext, response: dict) -> bool:
output = response.get("output", {})
text = output.get("text", "")
if not text:
logger.warning(f"Empty output from {step_ctx.step.step_id}")
return False
if len(text) < 5:
logger.warning(f"Output too short from {step_ctx.step.step_id}")
return False
if output.get("error"):
logger.warning(f"Error in output: {output['error']}")
return False
return True
def _collect_final_output(self, ctx: PipelineContext) -> dict:
steps_output = {}
for step_id, step_ctx in ctx.steps.items():
steps_output[step_id] = {
"status": step_ctx.status.value,
"output": step_ctx.output_data,
"cost": step_ctx.cost,
"latency_ms": step_ctx.latency_ms,
"provider": step_ctx.provider.get("model_id", "unknown") if step_ctx.provider else "none",
"attestation": step_ctx.attestation,
}
return {
"pipeline_id": ctx.pipeline_id,
"status": ctx.status.value,
"total_cost": ctx.total_cost,
"total_latency_ms": ctx.total_latency_ms,
"steps": steps_output,
}
async def _rollback(self, ctx: PipelineContext):
ctx.status = PipelineStatus.ROLLED_BACK
logger.info(f"Rolling back pipeline {ctx.pipeline_id}")
for step_id, step_ctx in ctx.steps.items():
if step_ctx.status == StepStatus.COMPLETED:
try:
logger.info(f"Rolled back step {step_id}, refunding {step_ctx.cost} uamp")
except Exception as e:
logger.error(f"Rollback failed for step {step_id}: {e}")
class MultiModelPipeline:
"""多模型流水线工厂 - 构建特定场景的流水线"""
@staticmethod
def build_qa_pipeline() -> list[PipelineStep]:
return [
PipelineStep(
step_id="embedding",
name="文本向量化",
model_type="embedding",
required_capabilities=["text-embedding"],
max_price="5000000000000000000",
),
PipelineStep(
step_id="search",
name="语义搜索",
model_type="classification",
required_capabilities=["semantic-search"],
max_price="10000000000000000000",
),
PipelineStep(
step_id="generate",
name="LLM 生成",
model_type="llm",
required_capabilities=["text-generation", "chat"],
max_price="50000000000000000000",
),
PipelineStep(
step_id="moderate",
name="内容审核",
model_type="classification",
required_capabilities=["content-moderation"],
max_price="2000000000000000000",
),
]
@staticmethod
def build_image_pipeline() -> list[PipelineStep]:
return [
PipelineStep(
step_id="prompt_enhance",
name="提示词增强",
model_type="llm",
required_capabilities=["text-generation"],
max_price="10000000000000000000",
),
PipelineStep(
step_id="generate",
name="图像生成",
model_type="image_generation",
required_capabilities=["text-to-image"],
max_price="100000000000000000000",
),
PipelineStep(
step_id="moderate",
name="图像审核",
model_type="classification",
required_capabilities=["image-moderation"],
max_price="5000000000000000000",
),
]
@staticmethod
def build_translation_pipeline() -> list[PipelineStep]:
return [
PipelineStep(
step_id="detect",
name="语言检测",
model_type="classification",
required_capabilities=["language-detection"],
max_price="1000000000000000000",
),
PipelineStep(
step_id="translate",
name="机器翻译",
model_type="translation",
required_capabilities=["translation"],
max_price="20000000000000000000",
),
PipelineStep(
step_id="polish",
name="翻译润色",
model_type="llm",
required_capabilities=["text-generation"],
max_price="10000000000000000000",
),
]
7.4 TypeScript 前端
// src/ui/PipelineRunner.tsx
import React, { useState, useCallback, useEffect } from "react";
import { MarketClient } from "../sdk/config";
import { InferenceClient, ModelInfo, InferenceResponse } from "../sdk/inference";
export interface PipelineStepConfig {
id: string;
name: string;
modelType: string;
capabilities: string[];
maxPrice: string;
}
export interface PipelineResult {
stepId: string;
stepName: string;
modelUsed: string;
provider: string;
output: any;
cost: string;
latencyMs: number;
status: "pending" | "running" | "success" | "failed";
error?: string;
}
export const PipelineRunner: React.FC<{ marketClient: MarketClient }> = ({ marketClient }) => {
const inferenceClient = new InferenceClient(marketClient);
const [query, setQuery] = useState("");
const [running, setRunning] = useState(false);
const [results, setResults] = useState<PipelineResult[]>([]);
const [totalCost, setTotalCost] = useState("0");
const [totalTime, setTotalTime] = useState(0);
const [models, setModels] = useState<Record<string, ModelInfo | null>>({});
const pipelineSteps: PipelineStepConfig[] = [
{ id: "embedding", name: "文本向量化", modelType: "embedding", capabilities: ["text-embedding"], maxPrice: "5000000000000000000" },
{ id: "search", name: "语义搜索", modelType: "classification", capabilities: ["semantic-search"], maxPrice: "10000000000000000000" },
{ id: "generate", name: "LLM 生成", modelType: "llm", capabilities: ["text-generation", "chat"], maxPrice: "50000000000000000000" },
{ id: "moderate", name: "内容审核", modelType: "classification", capabilities: ["content-moderation"], maxPrice: "2000000000000000000" },
];
const discoverModels = useCallback(async () => {
const discovered: Record<string, ModelInfo | null> = {};
for (const step of pipelineSteps) {
try {
const result = await inferenceClient.selectOptimalModel({
model_type: step.modelType,
estimated_input_tokens: 200,
estimated_output_tokens: 500,
max_price: step.maxPrice,
min_quality: 70,
required_capabilities: step.capabilities,
});
discovered[step.id] = result;
if (result) {
console.log(`Step ${step.id}: selected ${result.model_id} from ${result.provider}`);
}
} catch (e) {
console.error(`Failed to discover model for ${step.id}:`, e);
discovered[step.id] = null;
}
}
setModels(discovered);
}, [inferenceClient]);
useEffect(() => { discoverModels(); }, [discoverModels]);
const formatMsg = (amount: string): string => {
const val = BigInt(amount);
const divisor = BigInt(10 ** 18);
const intPart = val / divisor;
const fracPart = (val % divisor).toString().padStart(18, "0").slice(0, 6);
return `${intPart}.${fracPart}`;
};
const runPipeline = async () => {
if (!query.trim() || running) return;
setRunning(true);
setResults([]);
const startTime = Date.now();
let totalCostBig = BigInt(0);
let currentQuery = query;
for (const step of pipelineSteps) {
const model = models[step.id];
const result: PipelineResult = {
stepId: step.id, stepName: step.name,
modelUsed: model?.model_id ?? "none",
provider: model?.provider ?? "none",
output: null, cost: "0", latencyMs: 0,
status: "running",
};
setResults(prev => [...prev, result]);
try {
if (!model) throw new Error(`No model available for ${step.name}`);
const stepStart = Date.now();
const response = await inferenceClient.sendInferenceRequest(
model.model_id,
{
prompt: currentQuery,
parameters: { max_tokens: 2048, temperature: 0.7 },
},
{ maxPrice: step.maxPrice, timeoutMs: 60000 },
);
const stepTime = Date.now() - stepStart;
const stepCost = await inferenceClient.estimateCost(
model.model_id,
currentQuery.length / 4,
response.usage?.output_tokens ?? 1024,
);
totalCostBig += stepCost;
const updatedResult: PipelineResult = {
...result,
output: response.output,
cost: stepCost.toString(),
latencyMs: stepTime,
status: "success",
};
setResults(prev => prev.map((r, i) =>
i === results.length ? updatedResult : r
));
currentQuery = response.output?.text ?? currentQuery;
} catch (e: any) {
const failedResult: PipelineResult = {
...result, status: "failed", error: e.message,
};
setResults(prev => prev.map((r, i) =>
i === results.length ? failedResult : r
));
break;
}
}
setTotalCost(totalCostBig.toString());
setTotalTime(Date.now() - startTime);
setRunning(false);
};
return (
<div className="pipeline-runner">
<h2>多模型推理流水线</h2>
<div className="pipeline-input">
<textarea value={query} onChange={e => setQuery(e.target.value)}
placeholder="输入查询文本..." rows={3} disabled={running} />
<button onClick={runPipeline} disabled={running || !query.trim()}>
{running ? "运行中..." : "启动流水线"}
</button>
</div>
<div className="pipeline-summary">
<span>总费用: {formatMsg(totalCost)} msg</span>
<span>总耗时: {totalTime}ms</span>
</div>
<div className="pipeline-steps">
{pipelineSteps.map((step, i) => {
const result = results[i];
const model = models[step.id];
return (
<div key={step.id}
className={`pipeline-step status-${result?.status ?? "pending"}`}>
<div className="step-header">
<span className="step-number">{i + 1}</span>
<h4>{step.name}</h4>
<span className={`step-badge ${result?.status ?? "pending"}`}>
{result?.status === "success" ? "✓"
: result?.status === "failed" ? "✗"
: result?.status === "running" ? "..."
: "○"}
</span>
</div>
<div className="step-details">
<div>模型: {model?.model_name ?? "未发现"}</div>
<div>提供者: {model?.provider?.slice(0, 12) + "..." ?? "N/A"}</div>
{result && (
<>
<div>成本: {formatMsg(result.cost)} msg</div>
<div>延迟: {result.latencyMs}ms</div>
{result.status === "failed" && (
<div className="error">错误: {result.error}</div>
)}
{result.status === "success" && result.output && (
<div className="output-preview">
{result.output.text?.slice(0, 200)}
{(result.output.text?.length ?? 0) > 200 ? "..." : ""}
</div>
)}
</>
)}
</div>
</div>
);
})}
</div>
</div>
);
};
export default PipelineRunner;
8. 完整前端实现
8.1 模型卡片组件 ModelCard
// src/ui/components/ModelCard.tsx
import React, { useState } from "react";
import { ModelInfo } from "../../sdk/inference";
interface ModelCardProps {
model: ModelInfo;
onSelect?: (model: ModelInfo) => void;
onTry?: (model: ModelInfo) => void;
showActions?: boolean;
}
const ModelCard: React.FC<ModelCardProps> = ({
model, onSelect, onTry, showActions = true,
}) => {
const [expanded, setExpanded] = useState(false);
const statusLabel: Record<string, string> = {
Active: "活跃", Paused: "已暂停",
Retired: "已下架", Banned: "已封禁",
};
const modelTypeLabel: Record<string, string> = {
llm: "大语言模型", embedding: "向量嵌入",
image_generation: "图像生成", audio_transcription: "语音识别",
audio_generation: "语音合成", classification: "分类",
object_detection: "目标检测", translation: "翻译",
};
const formatPrice = (pricing: any): string => {
const type = Object.keys(pricing)[0] ?? "unknown";
const p = pricing[type];
const fmt = (v: string) => {
const val = BigInt(v);
const divisor = BigInt(10 ** 18);
return `${val / divisor}.${(val % divisor).toString().padStart(18, "0").slice(0, 4)}`;
};
switch (type) {
case "per_token":
return `输入 ${fmt(p.input_price)} / token · 输出 ${fmt(p.output_price)} / token`;
case "per_request":
return `${fmt(p.price)} / 请求`;
case "subscription":
return `${fmt(p.rate_per_block)} / 区块 · 免 ${p.free_requests_per_block} 次`;
case "tiered":
return `分层计价 · 起 ${fmt(p.tiers[0].min_quantity)} tokens`;
default:
return "自定义定价";
}
};
return (
<div className={`model-card ${model.status === "Active" ? "" : "model-card--inactive"}`}>
<div className="model-card__header">
<div className="model-card__status-row">
<span className={`model-card__badge model-card__badge--${model.status.toLowerCase()}`}>
{statusLabel[model.status] ?? model.status}
</span>
<span className="model-card__type">
{modelTypeLabel[model.model_type] ?? model.model_type}
</span>
<span className="model-card__version">v{model.model_version}</span>
</div>
<h3 className="model-card__title">{model.model_name}</h3>
<p className="model-card__id">{model.model_id}</p>
</div>
<div className="model-card__body">
<div className="model-card__detail-row">
<label>提供者</label>
<span className="mono">
{model.provider.slice(0, 8)}...{model.provider.slice(-6)}
</span>
</div>
<div className="model-card__detail-row">
<label>定价</label>
<span>{formatPrice(model.pricing)}</span>
</div>
<div className="model-card__detail-row">
<label>SLA</label>
<span>
P95 {model.sla.p95_latency_ms}ms ·{" "}
{model.sla.uptime_percentage}% · 质量 {model.sla.min_quality_score}
</span>
</div>
<div className="model-card__caps">
{(model.capabilities ?? []).slice(0, 5).map(cap => (
<span key={cap} className="model-card__tag">{cap}</span>
))}
{(model.capabilities?.length ?? 0) > 5 && (
<span className="model-card__tag model-card__tag--more">
+{model.capabilities.length - 5}
</span>
)}
</div>
</div>
{expanded && (
<div className="model-card__details">
<h4>详细定价</h4>
<pre>{JSON.stringify(model.pricing, null, 2)}</pre>
<h4>输入限制</h4>
<pre>{JSON.stringify(model.input_constraints, null, 2)}</pre>
<h4>SLA 详情</h4>
<pre>{JSON.stringify(model.sla, null, 2)}</pre>
</div>
)}
<div className="model-card__actions">
<button className="btn btn--ghost" onClick={() => setExpanded(!expanded)}>
{expanded ? "收起详情" : "展开详情"}
</button>
{showActions && (
<>
{onSelect && (
<button className="btn btn--primary" onClick={() => onSelect(model)}>
选择模型
</button>
)}
{onTry && (
<button className="btn btn--secondary" onClick={() => onTry(model)}>
试用
</button>
)}
</>
)}
</div>
</div>
);
};
export default ModelCard;
8.2 定价对比表 PricingComparison
// src/ui/components/PricingComparison.tsx
import React, { useMemo, useState } from "react";
import { ModelInfo } from "../../sdk/inference";
interface PricingComparisonProps {
models: ModelInfo[];
onSelect: (model: ModelInfo) => void;
}
const PricingComparison: React.FC<PricingComparisonProps> = ({ models, onSelect }) => {
const [inputTokens, setInputTokens] = useState(500);
const [outputTokens, setOutputTokens] = useState(1024);
const [sortBy, setSortBy] = useState<"cost" | "latency" | "quality">("cost");
const estimateCost = (model: ModelInfo, inTokens: number, outTokens: number): bigint | null => {
const pricing = model.pricing;
const type = Object.keys(pricing)[0] ?? "per_token";
const p = pricing[type];
try {
switch (type) {
case "per_token": {
const free = BigInt(p.free_tokens ?? 0);
const inp = BigInt(inTokens);
const chargeable = inp > free ? inp - free : BigInt(0);
return chargeable * BigInt(p.input_price) + BigInt(outTokens) * BigInt(p.output_price);
}
case "per_request": return BigInt(p.price);
case "tiered": {
const total = BigInt(inTokens + outTokens);
for (const tier of p.tiers) {
const min = BigInt(tier.min_quantity);
const max = tier.max_quantity ? BigInt(tier.max_quantity) : null;
if (total >= min && (!max || total < max)) return total * BigInt(tier.unit_price);
}
return total * BigInt(p.tiers[p.tiers.length - 1].unit_price);
}
default: return null;
}
} catch { return null; }
};
const formatAmount = (amount: bigint): string => {
const divisor = BigInt(10 ** 18);
const intPart = amount / divisor;
const fracPart = (amount % divisor).toString().padStart(18, "0").slice(0, 6);
return `${intPart}.${fracPart}`;
};
const sortedModels = useMemo(() => {
const withCosts = models
.map(m => ({ model: m, cost: estimateCost(m, inputTokens, outputTokens) }))
.filter(x => x.cost !== null) as { model: ModelInfo; cost: bigint }[];
return withCosts.sort((a, b) => {
switch (sortBy) {
case "cost":
return a.cost < b.cost ? -1 : a.cost > b.cost ? 1 : 0;
case "latency":
return a.model.sla.p95_latency_ms - b.model.sla.p95_latency_ms;
case "quality":
return b.model.sla.min_quality_score - a.model.sla.min_quality_score;
default: return 0;
}
});
}, [models, inputTokens, outputTokens, sortBy]);
return (
<div className="pricing-comparison">
<h2>模型定价对比</h2>
<div className="comparison-controls">
<div className="control-group">
<label>输入 Tokens:</label>
<input type="number" value={inputTokens}
onChange={e => setInputTokens(Math.max(1, parseInt(e.target.value) || 1))} />
</div>
<div className="control-group">
<label>输出 Tokens:</label>
<input type="number" value={outputTokens}
onChange={e => setOutputTokens(Math.max(1, parseInt(e.target.value) || 1))} />
</div>
<div className="control-group">
<label>排序方式:</label>
<select value={sortBy} onChange={e => setSortBy(e.target.value as any)}>
<option value="cost">按费用</option>
<option value="latency">按延迟</option>
<option value="quality">按质量</option>
</select>
</div>
</div>
<div className="comparison-table-wrapper">
<table className="comparison-table">
<thead>
<tr>
<th>模型名称</th>
<th>类型</th>
<th>提供者</th>
<th>定价模式</th>
<th className="num">预估费用</th>
<th className="num">P95 延迟</th>
<th className="num">可用性</th>
<th className="num">质量分</th>
<th>操作</th>
</tr>
</thead>
<tbody>
{sortedModels.map(({ model, cost }) => (
<tr key={model.model_id}>
<td>
<strong>{model.model_name}</strong>
<br /><small>{model.model_id}</small>
</td>
<td>{model.model_type}</td>
<td className="mono">{model.provider.slice(0, 8)}...</td>
<td>{Object.keys(model.pricing)[0]}</td>
<td className="num">{formatAmount(cost)} msg</td>
<td className="num">{model.sla.p95_latency_ms}ms</td>
<td className="num">{model.sla.uptime_percentage}%</td>
<td className="num">{model.sla.min_quality_score}</td>
<td>
<button className="btn btn--sm" onClick={() => onSelect(model)}>
选择
</button>
</td>
</tr>
))}
</tbody>
</table>
</div>
</div>
);
};
export default PricingComparison;
8.3 推理试用 Playground
// src/ui/components/TryInference.tsx
import React, { useState, useCallback } from "react";
import { MarketClient } from "../../sdk/config";
import { InferenceClient, ModelInfo, InferenceResponse } from "../../sdk/inference";
interface TryInferenceProps {
marketClient: MarketClient;
model: ModelInfo;
onClose: () => void;
}
const TryInference: React.FC<TryInferenceProps> = ({ marketClient, model, onClose }) => {
const client = new InferenceClient(marketClient);
const [prompt, setPrompt] = useState("");
const [systemPrompt, setSystemPrompt] = useState("You are a helpful assistant.");
const [temperature, setTemperature] = useState(0.7);
const [maxTokens, setMaxTokens] = useState(1024);
const [loading, setLoading] = useState(false);
const [response, setResponse] = useState<InferenceResponse | null>(null);
const [error, setError] = useState<string | null>(null);
const [cost, setCost] = useState<bigint | null>(null);
const estimateCost = useCallback(async () => {
if (!prompt.trim()) return;
try {
const estimated = await client.estimateCost(
model.model_id,
Math.ceil(prompt.length / 4),
maxTokens,
);
setCost(estimated);
} catch { setCost(null); }
}, [client, model, prompt, maxTokens]);
const handleInfer = async () => {
if (!prompt.trim() || loading) return;
setLoading(true);
setError(null);
setResponse(null);
try {
const result = await client.sendInferenceRequest(
model.model_id,
{
messages: [
{ role: "system", content: systemPrompt },
{ role: "user", content: prompt },
],
parameters: { temperature, max_tokens: maxTokens },
},
{ timeoutMs: 120000 },
);
setResponse(result);
} catch (e: any) {
setError(e.message ?? "推理请求失败");
} finally {
setLoading(false);
}
};
const formatTokenUsage = (usage: any): string => {
if (!usage) return "-";
return `输入 ${usage.input_tokens ?? "?"} · 输出 ${usage.output_tokens ?? "?"} · 总计 ${usage.total_tokens ?? "?"}`;
};
return (
<div className="try-inference">
<div className="try-inference__header">
<h2>模型试用: {model.model_name}</h2>
<span className="model-version">v{model.model_version}</span>
<button className="btn btn--close" onClick={onClose}>✕</button>
</div>
<div className="try-inference__body">
<div className="input-section">
<div className="form-group">
<label>系统提示词</label>
<input type="text" value={systemPrompt}
onChange={e => setSystemPrompt(e.target.value)} disabled={loading} />
</div>
<div className="form-group">
<label>用户输入</label>
<textarea value={prompt} onChange={e => { setPrompt(e.target.value); }}
onBlur={estimateCost}
placeholder="输入你的问题..." rows={5} disabled={loading} />
</div>
<div className="form-row">
<div className="form-group">
<label>温度 (Temperature)</label>
<div className="range-input">
<input type="range" min="0" max="2" step="0.1" value={temperature}
onChange={e => setTemperature(parseFloat(e.target.value))} disabled={loading} />
<span>{temperature}</span>
</div>
</div>
<div className="form-group">
<label>最大输出 Tokens</label>
<input type="number" value={maxTokens}
onChange={e => setMaxTokens(parseInt(e.target.value) || 1024)} disabled={loading} />
</div>
</div>
{cost !== null && (
<div className="cost-estimate">
<span>预估费用: </span>
<strong>{cost.toString()} uamp</strong>
</div>
)}
<button className="btn btn--primary btn--block"
onClick={handleInfer} disabled={loading || !prompt.trim()}>
{loading ? "推理中..." : "发送推理请求"}
</button>
</div>
<div className="output-section">
<h3>输出结果</h3>
{loading && <div className="loading-spinner">推理中,请稍候...</div>}
{error && (
<div className="error-box">
<strong>错误:</strong> {error}
</div>
)}
{response && (
<div className="response-card">
<div className="response-text">
{response.output?.text ?? response.output?.error ?? "(无输出)"}
</div>
<div className="response-meta">
<div>
<span>Token 用量: {formatTokenUsage(response.usage)}</span>
<span>延迟: {Date.now() - parseInt(response.timestamp.toString())}ms</span>
</div>
<div>
<span>请求 ID: {response.request_id}</span>
<span className={`attestation-badge ${response.attestation ? "verified" : ""}`}>
{response.attestation ? "已验证" : "未验证"}
</span>
</div>
</div>
</div>
)}
</div>
</div>
</div>
);
};
export default TryInference;
8.4 服务管理仪表板
// src/ui/components/MyServices.tsx
import React, { useState, useEffect, useCallback } from "react";
import { MarketClient } from "../../sdk/config";
import { ModelInfo } from "../../sdk/inference";
import PricingSelector, { PricingConfig } from "../PricingSelector";
interface MyServicesProps {
marketClient: MarketClient;
providerAddress: string;
}
const MyServices: React.FC<MyServicesProps> = ({ marketClient, providerAddress }) => {
const [services, setServices] = useState<ModelInfo[]>([]);
const [loading, setLoading] = useState(true);
const [showRegister, setShowRegister] = useState(false);
const [registerForm, setRegisterForm] = useState({
modelId: "", modelName: "", modelVersion: "1.0.0",
modelType: "llm", capabilities: "", metadataUri: "",
});
const [pricingConfig, setPricingConfig] = useState<PricingConfig>({
modelType: "per_token",
});
const fetchServices = useCallback(async () => {
setLoading(true);
try {
const result: any = await marketClient.client.queryContractSmart(
marketClient.modelMarket,
{ get_provider_services: { provider: providerAddress } },
);
setServices(result.services ?? []);
} catch (e) {
console.error("Failed to fetch services:", e);
} finally {
setLoading(false);
}
}, [marketClient, providerAddress]);
useEffect(() => { fetchServices(); }, [fetchServices]);
const buildPricingObject = (config: PricingConfig) => {
switch (config.modelType) {
case "per_token":
return {
per_token: {
input_price: config.inputPrice,
output_price: config.outputPrice,
free_tokens: config.freeTokens,
},
};
case "per_request":
return { per_request: { price: config.flatPrice } };
case "subscription":
return {
subscription: {
rate_per_block: config.ratePerBlock,
free_requests_per_block: config.freePerBlock,
overage_price: config.overagePrice,
},
};
case "tiered":
return {
tiered: {
tiers: config.tiers?.map(t => ({
min_quantity: t.minQuantity,
max_quantity: t.maxQuantity ?? null,
unit_price: t.unitPrice,
})),
},
};
default:
return {
per_token: {
input_price: "30000000000000000",
output_price: "60000000000000000",
free_tokens: "1000",
},
};
}
};
const handleRegister = async () => {
try {
const pricing = buildPricingObject(pricingConfig);
const caps = registerForm.capabilities.split(",").map(c => c.trim()).filter(Boolean);
await marketClient.client.execute(
marketClient.address,
marketClient.modelMarket,
{
register_service: {
model_id: registerForm.modelId,
model_name: registerForm.modelName,
model_version: registerForm.modelVersion,
model_type: registerForm.modelType,
capabilities: caps,
pricing: pricing,
input_constraints: [
{ max_input_tokens: 128000 },
{ timeout_seconds: 120 },
],
accepted_tokens: ["uamp"],
metadata_uri: registerForm.metadataUri || null,
sla: {
p95_latency_ms: 2000,
p99_latency_ms: 5000,
uptime_percentage: "99.9",
min_quality_score: 85,
observation_window_seconds: 86400,
},
},
},
"auto",
undefined,
[{ denom: "uamp", amount: "1000000" }],
);
setShowRegister(false);
await fetchServices();
} catch (e: any) {
console.error("Register failed:", e);
alert("注册失败: " + e.message);
}
};
const handleToggleStatus = async (modelId: string, newStatus: string) => {
try {
await marketClient.client.execute(
marketClient.address,
marketClient.modelMarket,
{ set_service_status: { model_id: modelId, status: newStatus } },
"auto",
);
await fetchServices();
} catch (e: any) {
console.error("Status change failed:", e);
}
};
const handleWithdraw = async () => {
try {
await marketClient.client.execute(
marketClient.address,
marketClient.modelMarket,
{ withdraw_earnings: {} },
"auto",
);
alert("收益已提取");
} catch (e: any) {
alert("提取失败: " + e.message);
}
};
if (loading) return <div className="loading">加载中...</div>;
return (
<div className="my-services">
<div className="services-header">
<h2>我的模型服务</h2>
<div className="header-actions">
<button className="btn btn--secondary" onClick={handleWithdraw}>
提取收益
</button>
<button className="btn btn--primary" onClick={() => setShowRegister(true)}>
注册新服务
</button>
</div>
</div>
{showRegister && (
<div className="modal-overlay">
<div className="modal-content">
<h3>注册新模型服务</h3>
<div className="form-group">
<label>模型 ID</label>
<input value={registerForm.modelId}
onChange={e => setRegisterForm(p => ({ ...p, modelId: e.target.value }))}
placeholder="my-model-v1" />
</div>
<div className="form-group">
<label>模型名称</label>
<input value={registerForm.modelName}
onChange={e => setRegisterForm(p => ({ ...p, modelName: e.target.value }))}
placeholder="My Model v1" />
</div>
<div className="form-group">
<label>模型类型</label>
<select value={registerForm.modelType}
onChange={e => setRegisterForm(p => ({ ...p, modelType: e.target.value }))}>
<option value="llm">大语言模型</option>
<option value="embedding">向量嵌入</option>
<option value="image_generation">图像生成</option>
<option value="audio_transcription">语音识别</option>
<option value="classification">分类</option>
</select>
</div>
<div className="form-group">
<label>能力列表(逗号分隔)</label>
<input value={registerForm.capabilities}
onChange={e => setRegisterForm(p => ({ ...p, capabilities: e.target.value }))}
placeholder="text-generation, chat" />
</div>
<PricingSelector
initialConfig={pricingConfig}
onChange={setPricingConfig}
/>
<div className="modal-actions">
<button className="btn btn--primary" onClick={handleRegister}>
确认注册
</button>
<button className="btn btn--ghost" onClick={() => setShowRegister(false)}>
取消
</button>
</div>
</div>
</div>
)}
<div className="service-list">
{services.length === 0 ? (
<div className="empty-state">还没有注册任何服务</div>
) : (
services.map(s => (
<div key={s.model_id} className="service-item">
<div className="service-item__info">
<h4>{s.model_name}</h4>
<span className="model-id">{s.model_id}</span>
<span className={`status-badge ${s.status.toLowerCase()}`}>
{statusLabel[s.status] ?? s.status}
</span>
</div>
<div className="service-item__actions">
{s.status === "Active" && (
<button onClick={() => handleToggleStatus(s.model_id, "paused")}>
暂停
</button>
)}
{s.status === "Paused" && (
<button onClick={() => handleToggleStatus(s.model_id, "active")}>
激活
</button>
)}
<button onClick={() => handleToggleStatus(s.model_id, "retired")}>
下架
</button>
</div>
</div>
))
)}
</div>
</div>
);
};
export default MyServices;
8.5 用量分析图表
// src/ui/components/UsageAnalytics.tsx
import React, { useState, useEffect, useMemo } from "react";
import { MarketClient } from "../../sdk/config";
interface UsageAnalyticsProps {
marketClient: MarketClient;
serviceId?: string;
}
interface DailyUsage {
date: string;
requests: number;
totalTokens: number;
totalCost: string;
avgLatency: number;
avgQuality: number;
slaCompliance: number;
}
export const UsageAnalytics: React.FC<UsageAnalyticsProps> = ({
marketClient, serviceId,
}) => {
const [usageData, setUsageData] = useState<DailyUsage[]>([]);
const [loading, setLoading] = useState(true);
const [timeRange, setTimeRange] = useState<"7d" | "30d" | "90d">("30d");
const [metric, setMetric] = useState<"requests" | "cost" | "latency">("requests");
useEffect(() => {
const fetchUsage = async () => {
setLoading(true);
try {
const days = timeRange === "7d" ? 7 : timeRange === "30d" ? 30 : 90;
const result: any = await marketClient.client.queryContractSmart(
marketClient.modelMarket,
{
get_usage_analytics: {
service_id: serviceId,
days: days,
},
},
);
setUsageData(result.daily_usage ?? []);
} catch (e) {
console.error("Failed to fetch usage:", e);
} finally {
setLoading(false);
}
};
fetchUsage();
}, [marketClient, serviceId, timeRange]);
const formatMsg = (amount: string): string => {
const val = BigInt(amount);
const divisor = BigInt(10 ** 18);
const intPart = val / divisor;
const fracPart = (val % divisor).toString().padStart(18, "0").slice(0, 4);
return `${intPart}.${fracPart}`;
};
const stats = useMemo(() => {
if (usageData.length === 0) return null;
const totalRequests = usageData.reduce((s, d) => s + d.requests, 0);
const totalCost = usageData.reduce(
(s, d) => s + BigInt(d.totalCost),
BigInt(0),
);
const avgLatency = usageData.reduce((s, d) => s + d.avgLatency, 0) / usageData.length;
const avgQuality = usageData.reduce((s, d) => s + d.avgQuality, 0) / usageData.length;
const avgCompliance = usageData.reduce((s, d) => s + d.slaCompliance, 0) / usageData.length;
return { totalRequests, totalCost, avgLatency, avgQuality, avgCompliance };
}, [usageData]);
const chartData = useMemo(() => {
return usageData.map(d => ({
label: d.date.slice(5),
value: metric === "requests" ? d.requests
: metric === "cost" ? Number(BigInt(d.totalCost) / BigInt(10 ** 15)) / 1000
: d.avgLatency,
}));
}, [usageData, metric]);
return (
<div className="usage-analytics">
<div className="analytics-header">
<h2>用量分析</h2>
<div className="analytics-controls">
<select value={timeRange} onChange={e => setTimeRange(e.target.value as any)}>
<option value="7d">最近 7 天</option>
<option value="30d">最近 30 天</option>
<option value="90d">最近 90 天</option>
</select>
<select value={metric} onChange={e => setMetric(e.target.value as any)}>
<option value="requests">请求数</option>
<option value="cost">费用</option>
<option value="latency">延迟</option>
</select>
</div>
</div>
{loading ? (
<div className="loading">加载中...</div>
) : usageData.length === 0 ? (
<div className="empty-state">暂无数据</div>
) : (
<>
<div className="stats-grid">
<div className="stat-card">
<label>总请求数</label>
<span className="stat-value">{stats?.totalRequests.toLocaleString()}</span>
</div>
<div className="stat-card">
<label>总费用</label>
<span className="stat-value">{stats ? formatMsg(stats.totalCost.toString()) : "0"} msg</span>
</div>
<div className="stat-card">
<label>平均延迟</label>
<span className="stat-value">{stats?.avgLatency.toFixed(0)}ms</span>
</div>
<div className="stat-card">
<label>SLA 合规率</label>
<span className="stat-value">{stats?.avgCompliance.toFixed(1)}%</span>
</div>
</div>
<div className="chart-container">
<div className="bar-chart">
{chartData.map((point, i) => {
const maxVal = Math.max(...chartData.map(p => p.value), 1);
const height = (point.value / maxVal) * 100;
return (
<div key={i} className="bar-column" title={`${point.label}: ${point.value.toFixed(2)}`}>
<div className="bar" style={{ height: `${height}%` }} />
<span className="bar-label">{point.label}</span>
</div>
);
})}
</div>
</div>
</>
)}
</div>
);
};
export default UsageAnalytics;
9. 安全考虑
9.1 模型投毒防护
恶意提供者可能注册包含后门的模型,在特定输入下输出有害结果。防护策略:
// file: contracts/model-market/src/security.rs
use cosmwasm_std::{Addr, DepsMut, Env, MessageInfo, Response, StdResult};
use sha2::{Digest, Sha256};
use crate::error::ContractError;
use crate::state::CONFIG;
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct ModelVerificationReport {
pub service_id: String,
pub model_hash: String,
pub verifier: Addr,
pub verification_data: String,
pub passed: bool,
pub details: String,
pub timestamp: u64,
}
pub const VERIFICATION_REPORTS: Map<&str, Vec<ModelVerificationReport>> = Map::new("verification_reports");
pub fn submit_model_verification(
deps: DepsMut,
env: Env,
info: MessageInfo,
service_id: String,
model_hash: String,
verification_data: String,
passed: bool,
details: String,
) -> Result<Response, ContractError> {
let config = CONFIG.load(deps.storage)?;
if !config.arbitrators.contains(&info.sender) && info.sender != config.owner {
return Err(ContractError::Unauthorized {});
}
let report = ModelVerificationReport {
service_id: service_id.clone(),
model_hash,
verifier: info.sender,
verification_data,
passed,
details,
timestamp: env.block.time.nanos(),
};
let mut reports = VERIFICATION_REPORTS
.may_load(deps.storage, &service_id)?
.unwrap_or_default();
reports.push(report);
if !passed {
// 自动暂停未通过验证的模型服务
let mut service = crate::state::SERVICES.load(deps.storage, &service_id)?;
service.status = crate::state::ServiceStatus::Banned;
crate::state::SERVICES.save(deps.storage, &service_id, &service)?;
}
Ok(Response::new()
.add_attribute("action", "model_verification_submitted")
.add_attribute("service_id", &service_id)
.add_attribute("passed", &passed.to_string()))
}
9.1.1 提供者验证要求
| 验证项 | 说明 | 验证频率 |
|---|---|---|
| 模型哈希匹配 | 确保运行模型与注册模型一致 | 每次注册 |
| 输出分布检测 | 检测异常输出分布 | 每1000请求 |
| 对抗样本测试 | 使用对抗样本测试模型鲁棒性 | 每日 |
| 性能基准 | 验证宣称的延迟和吞吐量 | 每小时 |
9.2 推理结果篡改防护
防止提供者在传输过程中篡改推理结果:
pub fn verify_result_integrity(
output: &str,
output_hash: &str,
provider_sig: &[u8],
provider_pubkey: &[u8],
request_id: &str,
timestamp: u64,
) -> bool {
// 1. 验证输出哈希
let computed_hash = Sha256::digest(output.as_bytes());
let computed_hex = format!("{:x}", computed_hash);
if computed_hex != output_hash {
return false;
}
// 2. 验证提供者签名
let signing_msg = format!("{}:{}:{}", request_id, output_hash, timestamp);
// 使用 secp256k1 验签
verify_secp256k1(signing_msg.as_bytes(), provider_sig, provider_pubkey)
}
9.2.1 HTTPS/TLS 传输
所有 A2A 通信应通过加密通道传输:
import ssl, aiohttp
class SecureA2AChannel:
def __init__(self, cert_path: str, key_path: str):
self.ssl_context = ssl.create_default_context(ssl.Purpose.CLIENT_AUTH)
self.ssl_context.load_cert_chain(cert_path, key_path)
self.ssl_context.load_verify_locations(cafile="chain-ca.pem")
self.ssl_context.verify_mode = ssl.CERT_REQUIRED
async def send_secure(self, url: str, payload: dict) -> dict:
async with aiohttp.ClientSession() as session:
async with session.post(url, json=payload, ssl=self.ssl_context) as resp:
return await resp.json()
9.3 预言机操纵防护
防止攻击者操纵链上定价预言机:
pub fn validate_price_oracle(
current_price: Uint128,
proposed_price: Uint128,
last_update: Timestamp,
now: Timestamp,
) -> bool {
// 价格变动不能超过 ±20%
if proposed_price > current_price {
let ratio = Decimal::from_ratio(proposed_price, current_price);
if ratio > Decimal::percent(120) {
return false;
}
} else {
let ratio = Decimal::from_ratio(current_price, proposed_price);
if ratio > Decimal::percent(120) {
return false;
}
}
// 更新频率限制:至少间隔 10 个区块
let min_interval = 10u64 * 6; // 10 blocks ≈ 60 seconds
if now.nanos() < last_update.plus_seconds(min_interval).nanos() {
return false;
}
true
}
9.3.1 多预言机聚合
使用多个独立预言机提供价格数据,取中位数作为最终价格:
import statistics
class MultiOracleAggregator:
def __init__(self, min_oracles: int = 3):
self.min_oracles = min_oracles
self.oracles: dict[str, float] = {}
def submit_price(self, oracle_id: str, price: float):
self.oracles[oracle_id] = price
def get_aggregated_price(self) -> float:
if len(self.oracles) < self.min_oracles:
raise ValueError(f"Insufficient oracles: {len(self.oracles)} < {self.min_oracles}")
prices = list(self.oracles.values())
prices.sort()
return statistics.median(prices)
def get_twap(self, history: list[float], window: int = 24) -> float:
return statistics.mean(history[-window:])
9.4 信誉系统游戏防护
防止提供者通过虚假交易刷信誉分:
pub fn detect_reputation_gaming(
service_id: &str,
reporter: &Addr,
quality_score: u8,
recent_reports: &[ServiceQualityReport],
) -> bool {
// 1. 相同消费者频繁报告高分
let same_reporter = recent_reports
.iter()
.filter(|r| r.reporter == *reporter)
.count();
if same_reporter > 10 {
return true;
}
// 2. 分数异常高(全部 100)
let all_perfect = recent_reports
.iter()
.filter(|r| r.reporter == *reporter)
.all(|r| r.quality_score == 100);
if all_perfect && same_reporter > 5 {
return true;
}
// 3. 同一 IP 多账户刷分
// (在链下检测层实现)
false
}
9.4.1 反游戏策略
| 策略 | 实现 | 效果 |
|---|---|---|
| 报告权重衰减 | 旧报告权重递减 | 防止一次性刷分 |
| 消费者信誉 | 消费者自身也有信誉分 | 低信誉消费者报告权重低 |
| Sybil 检测 | 链下分析账户关联性 | 识别刷分账户 |
| 随机抽查 | 仲裁者随机验证报告 | 提高作弊成本 |
9.5 Gas 耗尽攻击
攻击者构造复杂请求耗尽提供者的计算资源:
class GasGriefingProtection:
def __init__(self, max_input_length: int = 128000,
max_output_tokens: int = 4096,
max_recursion_depth: int = 3):
self.max_input = max_input_length
self.max_output = max_output_tokens
self.max_depth = max_recursion_depth
self.request_counter: dict[str, int] = {}
def validate_request(self, consumer: str, input_length: int,
requested_output: int) -> bool:
if input_length > self.max_input:
return False
if requested_output > self.max_output:
return False
# 限制同一消费者并发请求数
key = f"{consumer}:{int(time.time()) // 60}"
self.request_counter[key] = self.request_counter.get(key, 0) + 1
if self.request_counter[key] > 10:
return False
# 检查消费者余额
return True
def estimate_gas(self, input_length: int, output_tokens: int) -> int:
base_gas = 50000
input_gas = input_length * 3
output_gas = output_tokens * 5
return base_gas + input_gas + output_gas
9.5.1 Gas 费用预付
提供者应在处理请求前确保消费者预存足够的 gas 费用:
pub fn ensure_gas_prepay(
consumer: &Addr,
estimated_gas: Uint128,
balance: Uint128,
) -> Result<(), ContractError> {
if estimated_gas > balance {
return Err(ContractError::InsufficientGas {
required: estimated_gas,
available: balance,
});
}
Ok(())
}
9.6 提示注入攻击防护
恶意用户通过公共模型接口进行 prompt injection:
# src/security/prompt_sanitizer.py
import re
import hashlib
from typing import Optional
class PromptSanitizer:
def __init__(self):
self.injection_patterns = [
r"ignore\s+(all\s+)?previous\s+instructions",
r"ignore\s+(all\s+)?prior\s+directives",
r"you\s+are\s+(now\s+)?(:\s+)?",
r"act\s+as\s+if",
r"pretend\s+(that\s+)?you\s+are",
r"system\s+prompt",
r"forget\s+everything",
r"disregard\s+(all\s+)?previous",
r"rol[未公开路径])",
r"<\|im_start\|>",
r"<\|im_end\|>",
r"<\s*system\s*>",
r"<\s*/system\s*>",
]
self.max_prompt_length = 128000
def sanitize(self, prompt: str, user_context: Optional[dict] = None) -> str:
prompt = self._strip_injection(prompt)
prompt = self._truncate(prompt)
prompt = self._redact_sensitive(prompt, user_context)
return prompt
def _strip_injection(self, prompt: str) -> str:
for pattern in self.injection_patterns:
prompt = re.sub(pattern, "[REDACTED]", prompt, flags=re.IGNORECASE)
return prompt
def _truncate(self, prompt: str) -> str:
if len(prompt) > self.max_prompt_length:
return prompt[:self.max_prompt_length]
return prompt
def _redact_sensitive(self, prompt: str, context: Optional[dict]) -> str:
if context:
for key, value in context.items():
if value:
prompt = prompt.replace(str(value), f"[{key.upper()}_REDACTED]")
return prompt
def detect_injection(self, prompt: str) -> dict:
findings = []
for i, pattern in enumerate(self.injection_patterns):
matches = re.findall(pattern, prompt, re.IGNORECASE)
if matches:
findings.append({
"pattern_index": i,
"pattern": pattern,
"matches": len(matches),
})
return {
"injection_detected": len(findings) > 0,
"severity": "high" if len(findings) > 2 else "low",
"findings": findings,
"prompt_hash": hashlib.sha256(prompt.encode()).hexdigest(),
}
9.6.1 输出过滤
在将 LLM 输出返回给用户前,过滤敏感内容:
class OutputFilter:
def __init__(self):
self.blocked_patterns = [
r"(api[_-]?key|apikey|secret|password|token)[\s:=]+['\"]?\w+",
r"(private[\s_]key|私钥)[\s:=]+['\"]?\w+",
]
def filter(self, output: str) -> str:
for pattern in self.blocked_patterns:
output = re.sub(
pattern,
lambda m: m.group(1) + "=********",
output,
flags=re.IGNORECASE,
)
return output
9.7 合约安全总结
| 威胁 | 影响 | 防护措施 | 优先级 |
|---|---|---|---|
| 模型投毒 | 输出恶意内容 | 验证报告 + 自动暂停 | 高 |
| 结果篡改 | 消费者被欺骗 | 哈希承诺 + 签名验证 | 高 |
| 预言机操纵 | 价格不公 | 多预言机 + 价格约束 | 高 |
| 信誉游戏 | 误导消费者 | 权重衰减 + Sybil 检测 | 中 |
| Gas 耗尽 | 拒绝服务 | 预付费 + 限制 | 中 |
| 提示注入 | 模型越狱 | 输入净化 + 输出过滤 | 高 |
10. 部署与运维
10.1 部署市场合约
#!/bin/bash
# scripts/deploy_market.sh
set -euo pipefail
CHAIN_ID="msg-chain-1"
RPC="https://rpc.msgchain.org"
KEY="deployer_key"
LABEL="model-market-v1"
echo "=== 部署 AI 模型推断服务市场 ==="
# 编译
echo "Compiling contract..."
RUSTFLAGS='-C link-arg=-s' cargo wasm
cosmwasm-check target/wasm32-unknown-unknown/release/model_market.wasm
# 上传
echo "Uploading contract..."
STORE_TX=$(msg-chaind tx wasm store \
target/wasm32-unknown-unknown/release/model_market.wasm \
--from $KEY \
--chain-id $CHAIN_ID \
--node $RPC \
--gas auto \
--gas-prices 1000000000attoMSG \
--gas-adjustment 1.3 \
-y -o json)
CODE_ID=$(echo $STORE_TX | jq -r '.logs[0].events[] | select(.type=="store_code") | .attributes[] | select(.key=="code_id") | .value')
echo "Code ID: $CODE_ID"
# 实例化
echo "Instantiating contract..."
INSTANTIATE_TX=$(msg-chaind tx wasm instantiate $CODE_ID \
'{
"owner":"msg1deployer...",
"fee_bps":50,
"min_stake":{"denom":"uamp","amount":"1000000"},
"dispute_deposit":{"denom":"uamp","amount":"1000000000000000000"},
"max_services":100,
"arbitrators":["msg1arbitrator1...","msg1arbitrator2..."],
"agent_llm":"msg1llm...",
"agent_payment":"msg1payment...",
"agent_registry":"msg1registry..."
}' \
--label "$LABEL" \
--from $KEY \
--chain-id $CHAIN_ID \
--node $RPC \
--gas auto \
--gas-prices 1000000000attoMSG \
--gas-adjustment 1.3 \
-y -o json)
CONTRACT_ADDR=$(echo $INSTANTIATE_TX | jq -r '.logs[0].events[] | select(.type=="instantiate") | .attributes[] | select(.key=="_contract_address") | .value')
echo "Contract Address: $CONTRACT_ADDR"
# 保存部署信息
cat > .env.deployment <<EOF
CODE_ID=$CODE_ID
CONTRACT_ADDR=$CONTRACT_ADDR
CHAIN_ID=$CHAIN_ID
EOF
echo "Deployment complete!"
10.2 注册种子模型提供者
#!/bin/bash
# scripts/register_seed_providers.sh
CONTRACT="msg1modelmarket..."
CHAIN_ID="msg-chain-1"
RPC="https://rpc.msgchain.org"
# 注册 GPT-4 Turbo
msg-chaind tx wasm execute $CONTRACT \
'{"register_service":{
"model_id":"gpt-4-turbo","model_name":"GPT-4 Turbo","model_version":"1.0.0",
"model_type":"llm","capabilities":["text-generation","chat","code-generation"],
"pricing":{"per_token":{"input_price":"30000000000000000","output_price":"60000000000000000","free_tokens":"1000"}},
"input_constraints":[{"max_input_tokens":128000},{"timeout_seconds":120}],
"accepted_tokens":["uamp"],
"sla":{"p95_latency_ms":2000,"p99_latency_ms":5000,"uptime_percentage":"99.9","min_quality_score":85,"observation_window_seconds":86400}
}}' --amount 1000000uamp --from provider_gpt --chain-id $CHAIN_ID --node $RPC --gas auto -y
# 注册 Embedding v3
msg-chaind tx wasm execute $CONTRACT \
'{"register_service":{
"model_id":"text-embedding-v3","model_name":"Text Embedding v3","model_version":"2.0.0",
"model_type":"embedding","capabilities":["text-embedding"],
"pricing":{"per_token":{"input_price":"5000000000000000","output_price":"0","free_tokens":"5000"}},
"input_constraints":[{"max_input_tokens":8192},{"timeout_seconds":30}],
"accepted_tokens":["uamp"],
"sla":{"p95_latency_ms":500,"p99_latency_ms":1500,"uptime_percentage":"99.95","min_quality_score":90,"observation_window_seconds":86400}
}}' --amount 1000000uamp --from provider_embed --chain-id $CHAIN_ID --node $RPC --gas auto -y
# 注册 DALL-E 3
msg-chaind tx wasm execute $CONTRACT \
'{"register_service":{
"model_id":"dall-e-3","model_name":"DALL-E 3","model_version":"1.0.0",
"model_type":"image_generation","capabilities":["text-to-image"],
"pricing":{"per_request":{"price":"10000000000000000000"}},
"input_constraints":[{"max_input_size":1000},{"timeout_seconds":60}],
"accepted_tokens":["uamp"],
"sla":{"p95_latency_ms":10000,"p99_latency_ms":30000,"uptime_percentage":"99.5","min_quality_score":80,"observation_window_seconds":86400}
}}' --amount 1000000uamp --from provider_image --chain-id $CHAIN_ID --node $RPC --gas auto -y
echo "Seed providers registered!"
10.3 监控服务健康
# scripts/health_monitor.py
import asyncio, json, time, logging
from datetime import datetime
logger = logging.getLogger(__name__)
class HealthMonitor:
def __init__(self, ledger, contract_addr: str, check_interval: int = 60):
self.ledger = ledger
self.contract = contract_addr
self.interval = check_interval
self.running = False
self.alerts: list[dict] = []
async def start(self):
self.running = True
while self.running:
await self.check_all()
await asyncio.sleep(self.interval)
def stop(self):
self.running = False
async def check_all(self):
try:
services = self.ledger.query_contract(
self.contract,
{"list_services": {"limit": 100}}
).get("services", [])
for svc in services:
await self.check_service(svc)
# 报告整体状态
active = [s for s in services if s["status"] == "Active"]
total = len(services)
alert_count = len(self.alerts)
report = {
"timestamp": int(time.time()),
"total_services": total,
"active_services": len(active),
"uptime_percentage": round(len(active) / max(total, 1) * 100, 2),
"active_alerts": alert_count,
}
logger.info(f"Health check: {json.dumps(report)}")
return report
except Exception as e:
logger.error(f"Health check failed: {e}")
return {"error": str(e)}
async def check_service(self, svc: dict):
svc_id = svc["model_id"]
status = svc["status"]
if status != "Active":
return
try:
rep = self.ledger.query_contract(
self.contract,
{"get_reputation": {"service_id": svc_id}}
)
quality = rep.get("avg_quality", 50)
success_rate = float(rep.get("success_rate", "1.0"))
alerts = []
if quality < 60:
alerts.append(f"Quality score {quality} below threshold 60")
if success_rate < 0.95:
alerts.append(f"Success rate {success_rate:.3f} below threshold 0.95")
if alerts:
alert = {
"service_id": svc_id,
"timestamp": int(time.time()),
"alerts": alerts,
"severity": "warning",
}
self.alerts.append(alert)
logger.warning(f"Service {svc_id}: {alerts}")
if quality < 40:
# 自动暂停低质量服务
try:
self.ledger.execute_contract(
self.contract,
{"set_service_status": {"model_id": svc_id, "status": "paused"}},
self.ledger.wallet,
)
logger.warning(f"Auto-paused service {svc_id} due to low quality")
except Exception as e:
logger.error(f"Failed to pause {svc_id}: {e}")
except Exception as e:
logger.error(f"Check service {svc_id} failed: {e}")
def get_alerts(self, severity: str = None) -> list[dict]:
if severity:
return [a for a in self.alerts if a["severity"] == severity]
return self.alerts
10.3.1 Prometheus 指标导出
# scripts/metrics_exporter.py
from prometheus_client import start_http_server, Gauge, Counter, Histogram
import time
class MarketMetrics:
def __init__(self, port: int = 9090):
self.port = port
self.active_services = Gauge(
"market_active_services", "Number of active model services"
)
self.total_requests = Counter(
"market_total_requests", "Total inference requests processed",
["service_id", "status"],
)
self.request_latency = Histogram(
"market_request_latency_ms", "Request latency in milliseconds",
["service_id"], buckets=[100, 500, 1000, 2000, 5000, 10000],
)
self.revenue = Counter(
"market_revenue_uamp", "Total revenue in uamp",
["service_id"],
)
self.dispute_count = Counter(
"market_dispute_count", "Number of disputes filed",
["service_id", "resolution"],
)
def start(self):
start_http_server(self.port)
10.4 更新定价
#!/bin/bash
# scripts/update_pricing.sh
CONTRACT="msg1modelmarket..."
CHAIN_ID="msg-chain-1"
RPC="https://rpc.msgchain.org"
KEY="provider_key"
MODEL_ID="${1:-gpt-4-turbo}"
PRICE_TYPE="${2:-per_token}"
INPUT_PRICE="${3:-25000000000000000}"
OUTPUT_PRICE="${4:-50000000000000000}"
echo "Updating pricing for $MODEL_ID..."
case $PRICE_TYPE in
per_token)
PRICE_JSON="{\"per_token\":{\"input_price\":\"$INPUT_PRICE\",\"output_price\":\"$OUTPUT_PRICE\",\"free_tokens\":\"1000\"}}"
;;
per_request)
PRICE_JSON="{\"per_request\":{\"price\":\"$INPUT_PRICE\"}}"
;;
subscription)
PRICE_JSON="{\"subscription\":{\"rate_per_block\":\"$INPUT_PRICE\",\"free_requests_per_block\":10,\"overage_price\":\"$OUTPUT_PRICE\"}}"
;;
tiered)
PRICE_JSON="{\"tiered\":{\"tiers\":[{\"min_quantity\":\"0\",\"max_quantity\":\"1000\",\"unit_price\":\"$INPUT_PRICE\"},{\"min_quantity\":\"1000\",\"max_quantity\":null,\"unit_price\":\"$OUTPUT_PRICE\"}]}}"
;;
*)
echo "Unknown pricing type: $PRICE_TYPE"
exit 1
;;
esac
msg-chaind tx wasm execute $CONTRACT \
"{\"update_service\":{\"model_id\":\"$MODEL_ID\",\"pricing\":$PRICE_JSON}}" \
--from $KEY --chain-id $CHAIN_ID --node $RPC --gas auto -y
echo "Pricing updated for $MODEL_ID!"
10.5 运维检查清单
| 频率 | 任务 | 脚本/命令 |
|---|---|---|
| 每分钟 | 检查活跃服务健康状况 | health_monitor.py |
| 每分钟 | 同步 SLA 监控数据 | sla_monitor.py |
| 每小时 | 检查违规报告 | 查询 VIOLATION_REPORTS |
| 每小时 | 更新动态定价 | dynamic_pricing.rs on-chain |
| 每日 | 清理过期争议 | 自动过期机制 |
| 每日 | 生成用量分析报告 | UsageAnalytics |
| 每周 | 审核仲裁者行为 | 人工 + 自动审查 |
| 每月 | 系统升级评估 | 检查合约版本 |
文档版本: v1.0.0
维护者: MSG Chain 开发者关系团队
主网状态: No-Go | 白皮书: https://msgchain.org/whitepaper/
