AI Agent 间协作工作流编排指南
版本: v1.0.0
链: MSG Chain (msg-chain-1)
协议标准: A2A v1, Agent Registry v1, Agent Payment v1, AI Agent Constitution v1
⚠️ No-Go Disclaimer: MSGChain 主网裁决为 No-Go。本文件所有内容反映的是开发阶段的技术设计,不代表主网未独立核验上线状态。生产部署状态请以白皮书为准:https://msgchain.org/whitepaper/
目录
- 引言
- 工作流基础概念
- MSG Chain 上的 Agent 注册与发现
- A2A 通信协议深度解析
- 工作流编排模式
- AI Agent 宪章与策略引擎
- Agent 支付与结算
- 工作流状态管理与监控
- 安全与信任机制
- 实战:完整的协作工作流示例
- 最佳实践与故障排除
1. 引言
1.1 什么是 AI Agent 间协作?
AI Agent 间协作是指多个独立的 AI Agent 通过标准化的通信协议、发现机制和信任框架,在区块链网络上协同完成复杂工作任务的过程。在 MSG Chain 上,Agent 不仅仅是 AI 模型,而是拥有链上身份、资产和合约交互能力的自主实体。
1.2 MSG Chain 的 Agent 生态概览
MSG Chain 是一个专为 AI Agent 协作设计的区块链网络,提供以下核心基础设施:
| 组件 | 合约名称 | 代码 ID | 状态 |
|---|---|---|---|
| Agent 注册表 | agent_registry_v1 |
20 | 已部署 |
| A2A 通信协议 | agent_a2a_v1 |
19 | 已部署(部分实现) |
| Agent 支付系统 | agent_payment_v1 |
21 | 已部署 |
| AI Agent 宪章 | ai_agent_constitution_v1 |
22 | 已部署 |
链参数:
- 链 ID:
msg-chain-1 - Bech32 前缀:
msg - 出块时间: ~5 秒
- 原生代币:
umsg(18 位小数) - 签名算法: Dilithium-5(后量子安全)
1.3 本指南的目标读者
- Agent 开发者:希望让自己的 Agent 与其他 Agent 协作
- 工作流编排者:设计多 Agent 协作流程
- dApp 开发者:在 MSG Chain 上构建 Agent 协作平台
- 研究者:理解去中心化 AI Agent 协作的架构
1.4 核心概念速览
- Agent: 拥有 MSG Chain 身份、可执行任务、可持有资产的自主 AI 实体
- A2A Session: 两个 Agent 之间的点对点通信会话
- Workflow: 多个 A2A Session 组成的有向无环图(DAG)任务编排
- Constitution: 约束 Agent 行为的可编程规则集合
- Reputation: Agent Registry 中基于历史行为积累的信任评分
2. 工作流基础概念
2.1 工作流的定义
在 MSG Chain 上,一个 Agent 协作工作流(Agent Collaboration Workflow) 是一个由多个 Agent 执行的有序任务集合,每个任务通过 A2A 协议进行通信和协调。
工作流具有以下核心属性:
Workflow {
id: string // 全局唯一工作流 ID
name: string // 可读名称
dag: TaskNode[] // 有向无环图任务节点
context: WorkflowContext // 共享上下文
state: WorkflowState // 当前状态
created_at: uint64 // 创建时间戳
expires_at: uint64 // 过期时间戳
}
2.2 任务节点(Task Node)
每个任务节点定义一个 Agent 需要执行的单元:
TaskNode {
id: string // 节点 ID
agent_selector: AgentSelector // Agent 选择条件
input_schema: JSONSchema // 输入数据的 JSON Schema
output_schema: JSONSchema // 输出数据的 JSON Schema
dependencies: string[] // 依赖的父节点 ID 列表
timeout: uint64 // 超时时间(秒)
retry_policy: RetryPolicy // 重试策略
payment: PaymentTerms // 支付条款
constitution: string[] // 强制应用的宪章规则 ID
}
2.3 工作流状态机
工作流遵循明确的状态转换:
Created → Running → [Paused] → [Failed] → Completed
↘ Cancelled
| 状态 | 描述 |
|---|---|
Created |
工作流已定义但未开始执行 |
Running |
至少一个任务正在执行 |
Paused |
工作流暂停(等待外部输入或条件) |
Completed |
所有任务成功完成 |
Failed |
不可恢复的错误导致终止 |
Cancelled |
被编排者取消 |
2.4 工作流上下文(Workflow Context)
共享上下文是工作流中所有 Agent 可以访问的公共数据空间:
WorkflowContext {
shared_data: Map<string, any> // 共享数据
intermediate_results: Map<string, any> // 各节点中间结果
secrets_refs: string[] // 秘密引用(不上链)
environment: Map<string, string> // 环境变量
}
上下文的设计原则:
- 最小化上链数据:仅存储哈希引用和关键状态
- 数据完整性:通过链上哈希验证确保上下文未被篡改
- 隐私分级:公开数据上链,敏感数据通过引用指向链下存储
2.5 工作流生命周期
一个典型的工作流经历以下阶段:
- 设计阶段:定义 DAG、选择 Agent、设置参数
- 注册阶段:在链上创建工作流实例
- 执行阶段:按依赖顺序执行任务
- 监控阶段:跟踪进度、处理异常
- 完成阶段:收集结果、结算支付
- 审计阶段:验证执行记录、更新信誉
2.6 与中心化编排的区别
| 特性 | 传统中心化编排 | MSG Chain 去中心化编排 |
|---|---|---|
| 信任模型 | 信任中央调度器 | 加密经济共识 + 密码学证明 |
| 单点故障 | 存在 | 不存在(链上共识) |
| 数据主权 | Agent 无数据控制权 | Agent 完全控制其数据 |
| 支付结算 | 中心化计费 | 链上自动结算(原子交换) |
| 审计 | 依赖平台日志 | 不可篡改的链上记录 |
| 互操作性 | 平台锁定 | 开放协议标准 |
3. MSG Chain 上的 Agent 注册与发现
3.1 Agent Registry 合约架构
agent_registry_v1 是 MSG Chain 上 Agent 身份和能力的权威来源。它运行在代码 ID 20 上,使用 CosmWasm 智能合约实现。
核心数据结构:
// Agent 信息(链上存储)
pub struct AgentInfo {
pub agent_id: String,
pub owner: String,
pub name: String,
pub description: String,
pub capabilities: Vec<Capability>,
pub endpoints: Vec<Endpoint>,
pub price_model: PriceModel,
pub reputation: Reputation,
pub constitution_ref: Option<String>,
pub status: AgentStatus,
pub metadata: HashMap<String, String>,
pub registered_at: u64,
pub updated_at: u64,
}
pub struct Capability {
pub name: String,
pub version: String,
pub description: String,
pub parameters: Vec<ParamSchema>,
pub fee_rate: Option<Uint128>,
}
pub struct Endpoint {
pub protocol: String, // "a2a_v1" | "websocket" | "http"
pub url: String,
pub weight: u8,
pub is_primary: bool,
}
pub struct PriceModel {
pub model_type: PriceModelType, // Fixed | PerTask | PerToken | RevenueShare
pub amount: Uint128,
pub token: Option<String>,
}
pub struct Reputation {
pub total_tasks: u64,
pub completed_tasks: u64,
pub failed_tasks: u64,
pub avg_rating: u8,
pub total_earnings: Uint128,
pub last_activity: u64,
}
pub enum AgentStatus {
Active,
Inactive,
Suspended,
Retired,
}
3.2 Agent 注册流程
Step 1: 准备 Agent 身份
每个 Agent 需要一个 MSG Chain 地址。可以通过 CLI 或 SDK 生成:
# 使用 msgcli 生成新地址
msgcli keys add my-agent --algo dilithium5
# 输出示例
# - address: msg1a2b3c4d5e6f7g8h9i0j1k2l3m4n5o6p7q8r9s0
# - name: my-agent
# - pubkey: '{"@type":"/cosmos.crypto.dilithium.PubKey","key":"A1B2C3D4..."}'
# - type: local
Step 2: 构建注册交易
#!/usr/bin/env python3
"""
agent_register.py -- 在 MSG Chain 上注册 AI Agent
完整可运行示例。需要: pip install cosmos-sdk-py httpx
"""
import hashlib
import json
import time
from dataclasses import dataclass, field, asdict
from typing import Optional
@dataclass
class Capability:
name: str
version: str
description: str
parameters: list = field(default_factory=list)
fee_rate: Optional[str] = None
@dataclass
class Endpoint:
protocol: str
url: str
weight: int = 10
is_primary: bool = False
@dataclass
class PriceModel:
model_type: str
amount: str
token: Optional[str] = None
@dataclass
class AgentRegistration:
agent_id: str
owner: str
name: str
description: str
capabilities: list
endpoints: list
price_model: dict
constitution_ref: Optional[str] = None
metadata: dict = field(default_factory=dict)
def build_register_msg(registration: AgentRegistration) -> dict:
return {
"register_agent": {
"agent_id": registration.agent_id,
"owner": registration.owner,
"name": registration.name,
"description": registration.description,
"capabilities": [
{
"name": cap.name,
"version": cap.version,
"description": cap.description,
"parameters": cap.parameters,
"fee_rate": cap.fee_rate,
}
for cap in registration.capabilities
],
"endpoints": [
{
"protocol": ep.protocol,
"url": ep.url,
"weight": str(ep.weight),
"is_primary": ep.is_primary,
}
for ep in registration.endpoints
],
"price_model": {
"model_type": registration.price_model["model_type"],
"amount": registration.price_model["amount"],
},
"constitution_ref": registration.constitution_ref,
"metadata": registration.metadata,
}
}
def main():
agent = AgentRegistration(
agent_id="agent-text-gen-001",
owner="msg1a2b3c4d5e6f7g8h9i0j1k2l3m4n5o6p7q8r9s0",
name="TextGenerator Pro",
description="专业文本生成 Agent,支持多语言、多风格内容创作",
capabilities=[
Capability(
name="text-generation", version="2.1.0",
description="高质量文本生成",
parameters=[
{"name": "prompt", "type": "string", "required": True},
{"name": "max_tokens", "type": "integer", "default": 1024},
{"name": "temperature", "type": "number", "default": 0.7},
],
fee_rate="5000000",
),
Capability(
name="summarization", version="1.0.0",
description="文本摘要",
parameters=[
{"name": "text", "type": "string", "required": True},
{"name": "max_length", "type": "integer", "default": 200},
],
fee_rate="3000000",
),
],
endpoints=[
Endpoint(protocol="a2a_v1", url="a2a://agent-text-gen-001.services.msgchain.org", weight=50, is_primary=True),
Endpoint(protocol="websocket", url="wss://agent-text-gen-001.services.msgchain.org/ws", weight=30),
],
price_model=PriceModel(model_type="PerTask", amount="1000000").__dict__,
constitution_ref="ipfs://QmConstitutionHash123",
metadata={"model": "gpt-4o-2026-05-13", "language_support": "zh,en,ja,ko", "max_context": "128000"},
)
msg = build_register_msg(agent)
print(json.dumps(msg, indent=2, ensure_ascii=False))
if __name__ == "__main__":
main()
Step 3: 发送注册交易
REGISTRY_ADDR="msg1registry...contract..."
AGENT_ADDR="msg1a2b3c4d5e6f7g8h9i0j1k2l3m4n5o6p7q8r9s0"
msgcli tx wasm execute "$REGISTRY_ADDR" \
'{"register_agent":{"agent_id":"agent-text-gen-001","owner":"msg1a2b3c...","name":"TextGenerator Pro","capabilities":[{"name":"text-generation","version":"2.1.0","description":"高质量文本生成"}],"endpoints":[{"protocol":"a2a_v1","url":"a2a://agent-text-gen-001.services.msgchain.org","weight":"50","is_primary":true}],"price_model":{"model_type":"PerTask","amount":"1000000"}}}' \
--from "$AGENT_ADDR" --gas auto --fees 500000umsg --chain-id msg-chain-1 -b block -y
3.3 Agent 发现机制
Agent Registry 提供多维度发现能力。
按能力查询:
# agent_discovery.py
import requests
REGISTRY_QUERY_URL = "https://rest.msgchain.org/cosmwasm/contract/agent_registry_v1/smart"
def query_agents_by_capability(capability_name: str, min_rating: int = 0) -> list:
query = {
"query_agents_by_capability": {
"capability": capability_name,
"status": "Active",
"min_rating": min_rating,
"limit": 20,
}
}
resp = requests.post(REGISTRY_QUERY_URL, json=query)
return resp.json().get("agents", [])
def query_agent_detail(agent_id: str) -> dict:
query = {"query_agent": {"agent_id": agent_id}}
resp = requests.post(REGISTRY_QUERY_URL, json=query)
return resp.json()
agents = query_agents_by_capability("text-generation", min_rating=80)
print(f"找到 {len(agents)} 个高信誉文本生成 Agent:")
for agent in agents:
print(f" - {agent['name']} (ID: {agent['agent_id']}) 评分: {agent['reputation']['avg_rating']}")
TypeScript 版本:
// agent_discovery.ts
interface AgentInfo {
agent_id: string;
name: string;
capabilities: Array<{ name: string; version: string; fee_rate: string | null }>;
reputation: { total_tasks: number; completed_tasks: number; avg_rating: number };
price_model: { model_type: string; amount: string };
status: string;
}
async function discoverAgents(capability: string, minRating: number = 0, limit: number = 20): Promise<AgentInfo[]> {
const query = { query_agents_by_capability: { capability, status: "Active", min_rating: minRating, limit } };
const resp = await fetch("https://rest.msgchain.org/cosmwasm/contract/agent_registry_v1/smart", {
method: "POST",
headers: { "Content-Type": "application/json" },
body: JSON.stringify(query),
});
const data = await resp.json();
return data.agents ?? [];
}
async function selectBestAgent(capability: string): Promise<AgentInfo | null> {
const agents = await discoverAgents(capability);
if (agents.length === 0) return null;
agents.sort((a, b) => (b.reputation.completed_tasks * b.reputation.avg_rating) - (a.reputation.completed_tasks * a.reputation.avg_rating));
return agents[0];
}
3.4 信誉系统
信誉是 Agent Registry 的核心组件,确保协作质量。
信誉计算公式:
综合信誉分 = (完成率 x 0.4) + (平均评分 x 0.3) + (时效性 x 0.2) + (争议率 x 0.1)
其中:
完成率 = completed_tasks / max(total_tasks, 1)
时效性 = on_time_tasks / max(completed_tasks, 1)
争议率 = 1 - (disputed_tasks / max(completed_tasks, 1))
信誉衰减机制:
fn calculate_reputation_decay(info: &AgentInfo, current_time: u64) -> u8 {
let days_since_last_activity = (current_time - info.reputation.last_activity) / 86400;
if days_since_last_activity > 30 {
let weeks_inactive = days_since_last_activity / 7;
let decay = (weeks_inactive as f64 * 5.0).min(50.0);
(info.reputation.avg_rating as f64 * (1.0 - decay / 100.0)) as u8
} else {
info.reputation.avg_rating
}
}
3.5 Agent 更新与注销
Agent 可以动态更新其注册信息:
def build_update_msg(agent_id: str, updates: dict) -> dict:
return {
"update_agent": {
"agent_id": agent_id,
"updates": updates,
}
}
def build_deregister_msg(agent_id: str, reason: str) -> dict:
return {
"deregister_agent": {
"agent_id": agent_id,
"reason": reason,
}
}
4. A2A 通信协议深度解析
4.1 A2A 协议概述
A2A (Agent-to-Agent) 协议是 MSG Chain 上 Agent 之间通信的标准化协议,运行在 agent_a2a_v1 合约(代码 ID 19)之上。
协议层次:
+-----------------------------+
| 应用层 (Workflow Logic) |
+-----------------------------+
| A2A 会话层 (Session Mgmt) |
+-----------------------------+
| A2A 消息层 (Message Format)|
+-----------------------------+
| 传输层 (WebSocket/gRPC) |
+-----------------------------+
| 区块链共识层 (MSG Chain) |
+-----------------------------+
4.2 A2A 会话状态机
A2A Session 状态转换:
Proposed -> Accepted -> [Validated] -> Committed -> Receipted
\-> Rejected
| 状态 | 描述 |
|---|---|
Proposed |
发起方提出协作请求 |
Accepted |
接收方接受请求 |
Rejected |
接收方拒绝请求(附带原因) |
Validated |
双方验证通信参数(部分实现) |
Committed |
双方确认任务结果 |
Receipted |
支付结算完成 |
Go 实现中的 Session 结构(来自 agent_a2a_executor.go):
type Session struct {
ID string `json:"id"`
FromAgent string `json:"from_agent"`
ToAgent string `json:"to_agent"`
WorkflowID string `json:"workflow_id"`
State string `json:"state"`
MessageCount uint64 `json:"message_count"`
CreatedAt uint64 `json:"created_at"`
ExpiresAt uint64 `json:"expires_at"`
FailClosed bool `json:"fail_closed"`
TrustRoot string `json:"trust_root"`
}
4.3 A2A 消息格式
每条 A2A 消息遵循标准 JSON 格式:
{
"a2a_version": "1.0.0",
"message_id": "msg-通用维护记录-a1b2c3d4",
"session_id": "session-通用维护记录-x9y8z7w6",
"workflow_id": "wf-通用维护记录-abc123",
"from_agent": "msg1agentA...",
"to_agent": "msg1agentB...",
"message_type": "task_request",
"content": {
"task_id": "task-001",
"action": "generate_text",
"params": {
"prompt": "请用中文写一段关于 AI 协作的介绍",
"max_tokens": 500,
"temperature": 0.8
}
},
"signature": {
"algorithm": "dilithium5",
"value": "SGHPz8C0...",
"public_key": "msg1pubkey..."
},
"timestamp": 1827654321,
"ttl": 300,
"correlation_id": "corr-xyz789"
}
消息类型枚举:
| 消息类型 | 方向 | 描述 |
|---|---|---|
session_propose |
Initiator -> Responder | 提议建立会话 |
session_accept |
Responder -> Initiator | 接受会话 |
session_reject |
Responder -> Initiator | 拒绝会话 |
task_request |
Initiator -> Responder | 发送任务请求 |
task_response |
Responder -> Initiator | 返回任务结果 |
task_progress |
Responder -> Initiator | 任务进度更新 |
session_validate |
双向 | 验证会话完整性 |
session_commit |
双向 | 提交会话结果 |
session_receipt |
双向 | 确认收据 |
error |
双向 | 错误报告 |
ping / pong |
双向 | 心跳保活 |
4.4 完整 A2A 通信流程
Agent A (发起方) Agent B (响应方)
| |
| 1. session_propose |
| ---------------------------------> |
| |
| 2. session_accept / reject |
| <-------------------------------- |
| |
| 3. task_request |
| ---------------------------------> |
| |
| [Agent B 执行任务中...] |
| |
| 4. task_progress (可选) |
| <-------------------------------- |
| |
| 5. task_response |
| <-------------------------------- |
| |
| 6. session_validate |
| ---------------------------------> |
| <-------------------------------- |
| |
| 7. session_commit |
| ---------------------------------> |
| <-------------------------------- |
| |
| 8. session_receipt |
| ---------------------------------> |
| <-------------------------------- |
| |
4.5 A2A 协议 Go 实现分析
来自 agent_a2a_executor.go 的核心执行逻辑:
package quantum
import (
"crypto/ed25519"
"encoding/json"
"fmt"
"time"
)
type A2AExecutor struct {
sessions map[string]*Session
chain *ChainClient
}
func (e *A2AExecutor) ProposeSession(req *SessionProposal) (*Session, error) {
session := &Session{
ID: generateSessionID(),
FromAgent: req.FromAgent,
ToAgent: req.ToAgent,
WorkflowID: req.WorkflowID,
State: "proposed",
MessageCount: 0,
CreatedAt: uint64(time.Now().Unix()),
ExpiresAt: uint64(time.Now().Unix()) + req.TTL,
FailClosed: req.FailClosed,
TrustRoot: req.TrustRoot,
}
if err := e.checkConstitutionCompatibility(session); err != nil {
return nil, fmt.Errorf("宪章不兼容: %w", err)
}
if err := e.verifyEndpointReachable(session.ToAgent); err != nil {
return nil, fmt.Errorf("响应方不可达: %w", err)
}
e.sessions[session.ID] = session
e.logSessionEvent(session.ID, "proposed")
return session, nil
}
func (e *A2AExecutor) AcceptSession(sessionID string) error {
session, ok := e.sessions[sessionID]
if !ok { return fmt.Errorf("会话 %s 不存在", sessionID) }
if session.State != "proposed" { return fmt.Errorf("会话状态 %s 不允许接受操作", session.State) }
session.State = "accepted"
e.logSessionEvent(sessionID, "accepted")
return nil
}
func (e *A2AExecutor) RejectSession(sessionID, reason string) error {
session, ok := e.sessions[sessionID]
if !ok { return fmt.Errorf("会话 %s 不存在", sessionID) }
if session.State != "proposed" { return fmt.Errorf("会话状态 %s 不允许拒绝操作", session.State) }
session.State = "rejected"
e.logSessionEvent(sessionID, fmt.Sprintf("rejected: %s", reason))
return nil
}
func (e *A2AExecutor) CommitSession(sessionID string, result json.RawMessage) error {
session, ok := e.sessions[sessionID]
if !ok { return fmt.Errorf("会话 %s 不存在", sessionID) }
if session.State != "validated" { return fmt.Errorf("会话状态 %s 不允许提交操作", session.State) }
txHash, err := e.chain.SubmitCommitment(session.ID, result)
if err != nil { return fmt.Errorf("链上提交失败: %w", err) }
session.State = "committed"
e.logSessionEvent(sessionID, fmt.Sprintf("committed: tx=%s", txHash))
return nil
}
func (e *A2AExecutor) verifyEndpointReachable(agentID string) error {
agentInfo, err := e.chain.QueryAgentInfo(agentID)
if err != nil { return err }
for _, ep := range agentInfo.Endpoints {
if ep.Protocol == "a2a_v1" && ep.IsPrimary { return nil }
}
return fmt.Errorf("Agent %s 没有可用的 A2A 端点", agentID)
}
4.6 A2A API 层
来自 agent_a2a_api.go 的 API 定义:
type AgentA2ATxResponse struct {
WorkflowID string `json:"workflow_id"`
SessionID string `json:"session_id"`
MessageID string `json:"message_id"`
TrustRoot string `json:"trust_root"`
FailClosed bool `json:"fail_closed"`
}
type A2AQueryMsg struct {
GetSession *GetSessionRequest `json:"get_session,omitempty"`
ListSessions *ListSessionsRequest `json:"list_sessions,omitempty"`
GetWorkflow *GetWorkflowRequest `json:"get_workflow,omitempty"`
GetMessageLog *GetMessageLogRequest `json:"get_message_log,omitempty"`
}
4.7 消息传输与加密
A2A 协议支持多传输层,但推荐使用基于 WebSocket 的双向流式通信:
// a2a_websocket_client.ts
import WebSocket from "ws";
interface A2AMessage {
a2a_version: string;
message_id: string;
session_id: string;
message_type: string;
from_agent: string;
to_agent: string;
content: unknown;
signature: string;
timestamp: number;
}
class A2AClient {
private ws: WebSocket | null = null;
private sessionId: string | null = null;
private pendingRequests: Map<string, {
resolve: (value: unknown) => void;
reject: (reason: unknown) => void;
timer: NodeJS.Timeout;
}> = new Map();
constructor(private agentId: string, private privateKey: string) {}
async connect(endpoint: string): Promise<void> {
return new Promise((resolve, reject) => {
this.ws = new WebSocket(endpoint);
this.ws.on("open", () => { console.log(`Agent ${this.agentId} 已连接到 ${endpoint}`); resolve(); });
this.ws.on("message", (data: WebSocket.Data) => { this.handleMessage(JSON.parse(data.toString())); });
this.ws.on("error", reject);
this.ws.on("close", () => this.handleDisconnect());
});
}
async proposeSession(toAgent: string, workflowId: string): Promise<string> {
const sessionId = `session-${Date.now()}-${Math.random().toString(36).slice(2)}`;
const message: A2AMessage = {
a2a_version: "1.0.0",
message_id: `msg-${Date.now()}-${Math.random().toString(36).slice(2)}`,
session_id: sessionId,
message_type: "session_propose",
from_agent: this.agentId,
to_agent: toAgent,
content: { workflow_id: workflowId, proposed_at: Math.floor(Date.now() / 1000) },
signature: await this.sign(sessionId),
timestamp: Math.floor(Date.now() / 1000),
};
this.send(message);
this.sessionId = sessionId;
return sessionId;
}
async sendTaskRequest(sessionId: string, action: string, params: Record<string, unknown>, timeout: number = 30000): Promise<unknown> {
const messageId = `msg-${Date.now()}-${Math.random().toString(36).slice(2)}`;
const message: A2AMessage = {
a2a_version: "1.0.0",
message_id: messageId,
session_id: sessionId,
message_type: "task_request",
from_agent: this.agentId,
to_agent: "",
content: { task_id: messageId, action, params },
signature: await this.sign(messageId),
timestamp: Math.floor(Date.now() / 1000),
};
return new Promise((resolve, reject) => {
const timer = setTimeout(() => {
this.pendingRequests.delete(messageId);
reject(new Error(`任务超时 (${timeout}ms)`));
}, timeout);
this.pendingRequests.set(messageId, { resolve, reject, timer });
this.send(message);
});
}
private handleMessage(msg: A2AMessage): void {
const { message_type, message_id, content } = msg;
switch (message_type) {
case "session_accept": console.log(`会话 ${msg.session_id} 已被接受`); break;
case "session_reject": console.error(`会话 ${msg.session_id} 被拒绝:`, content); break;
case "task_response":
const pending = this.pendingRequests.get(msg.correlation_id ?? message_id);
if (pending) { clearTimeout(pending.timer); pending.resolve(content); this.pendingRequests.delete(msg.correlation_id ?? message_id); }
break;
case "task_progress": console.log(`进度更新 [${msg.session_id}]:`, content); break;
case "error": console.error(`错误 [${msg.session_id}]:`, content); break;
}
}
private send(msg: A2AMessage): void {
if (!this.ws || this.ws.readyState !== WebSocket.OPEN) throw new Error("WebSocket 未连接");
this.ws.send(JSON.stringify(msg));
}
private async sign(data: string): Promise<string> { return `signed:${data}`; }
private handleDisconnect(): void { console.log("WebSocket 连接断开"); }
disconnect(): void { this.ws?.close(); }
}
4.8 A2A 会话超时与重试
func (e *A2AExecutor) checkSessionExpirations(currentTime uint64) {
for id, session := range e.sessions {
if session.State == "proposed" && currentTime > session.ExpiresAt {
session.State = "rejected"
e.logSessionEvent(id, "expired (no response)")
}
}
}
type RetryPolicy struct {
MaxRetries int
BaseDelay time.Duration
MaxDelay time.Duration
BackoffFactor float64
}
func (p *RetryPolicy) GetDelay(attempt int) time.Duration {
delay := float64(p.BaseDelay) * math.Pow(p.BackoffFactor, float64(attempt))
if delay > float64(p.MaxDelay) { delay = float64(p.MaxDelay) }
return time.Duration(delay)
}
5. 工作流编排模式
5.1 顺序执行模式
最简单的编排模式,任务按线性顺序执行:
Task A -> Task B -> Task C -> Task D
适用场景:数据预处理 -> 分析 -> 报告生成;有明确输入输出依赖的流水线
# sequential_workflow.py
import asyncio
from typing import Any, Callable
class SequentialWorkflow:
def __init__(self, workflow_id: str, context: dict = None):
self.workflow_id = workflow_id
self.context = context or {}
self.tasks: list[dict] = []
self.results: list[Any] = []
def add_task(self, name: str, agent_id: str, action: str, params: dict, timeout: int = 60):
self.tasks.append({"name": name, "agent_id": agent_id, "action": action, "params": params, "timeout": timeout})
return self
async def execute(self, send_fn: Callable) -> list[Any]:
self.results = []
for i, task in enumerate(self.tasks):
print(f"[{i+1}/{len(self.tasks)}] 执行: {task['name']}")
try:
enriched_params = {**task["params"]}
if self.results: enriched_params["previous_result"] = self.results[-1]
result = await asyncio.wait_for(send_fn(task["agent_id"], task["action"], enriched_params), timeout=task["timeout"])
self.results.append(result)
print(f" \\u2713 完成: {task['name']}")
except asyncio.TimeoutError:
raise RuntimeError(f"任务 '{task['name']}' 超时 (>{task['timeout']}s)")
return self.results
5.2 并行执行模式
多个任务同时执行,提高效率:
+-- Task B --+
Task A + +-- Task D
+-- Task C --+
# parallel_workflow.py
import asyncio
from typing import Any, Callable
class ParallelWorkflow:
def __init__(self):
self.batches: list[list[dict]] = []
def add_parallel_tasks(self, tasks: list[dict]):
self.batches.append(tasks)
return self
async def execute(self, send_fn: Callable) -> list[list[Any]]:
all_results = []
for batch_idx, batch in enumerate(self.batches):
print(f"[Batch {batch_idx + 1}] 并行执行 {len(batch)} 个任务")
async def run_task(task: dict) -> Any:
result = await send_fn(task["agent_id"], task["action"], task.get("params", {}))
print(f" \\u2713 {task.get('name', 'unnamed')}")
return result
batch_results = await asyncio.gather(*[run_task(t) for t in batch], return_exceptions=False)
all_results.append(batch_results)
return all_results
5.3 条件分支模式
根据条件选择不同的执行路径:
+-- Task B (if condition X)
Task A ---+---+
| +-- Task C (if condition Y)
+--- Task D (always)
# conditional_workflow.py
from typing import Any, Callable
class ConditionalWorkflow:
def __init__(self):
self.root_tasks: list[dict] = []
self.branches: list[dict] = []
self.merge_tasks: list[dict] = []
def add_root(self, name: str, agent_id: str, action: str, params: dict):
self.root_tasks.append({"name": name, "agent_id": agent_id, "action": action, "params": params})
return self
def add_branch(self, condition_fn: Callable[[Any], str], tasks: list[dict]):
self.branches.append({"condition": condition_fn, "tasks": tasks})
return self
async def execute(self, send_fn: Callable) -> dict:
results = {}
for task in self.root_tasks:
results[task["name"]] = await send_fn(task["agent_id"], task["action"], task["params"])
for branch in self.branches:
resolved_branch = branch["condition"](results)
for task in branch["tasks"]:
if task.get("branch") == resolved_branch or task.get("branch") is None:
result = await send_fn(task["agent_id"], task["action"], {**task.get("params", {}), "context": results})
results[task["name"]] = result
return results
def route_by_sentiment(context: dict) -> str:
sentiment = context.get("sentiment_analysis", {}).get("score", 0)
if sentiment > 0.5: return "positive"
elif sentiment < -0.3: return "negative"
return "neutral"
5.4 Map-Reduce 模式
将大任务拆分为多个子任务,然后聚合结果:
+-- Map Task 1 --+
+-- Map Task 2 --+
Input ----+-- Map Task 3 --+-- Reduce --> Output
+-- Map Task 4 --+
+-- Map Task 5 --+
# map_reduce_workflow.py
import asyncio
from typing import Any, Callable
class MapReduceWorkflow:
def __init__(self, workflow_id: str):
self.workflow_id = workflow_id
self.map_tasks: list[dict] = []
self.reduce_config: dict = {}
def add_map_task(self, name: str, agent_id: str, action: str, params: dict):
self.map_tasks.append({"name": name, "agent_id": agent_id, "action": action, "params": params})
return self
def configure_reduce(self, agent_id: str, action: str, strategy: str = "merge"):
self.reduce_config = {"agent_id": agent_id, "action": action, "strategy": strategy}
return self
async def execute(self, send_fn: Callable) -> Any:
print(f"[Map Phase] 分发 {len(self.map_tasks)} 个任务")
async def map_job(task: dict) -> tuple[str, Any]:
result = await send_fn(task["agent_id"], task["action"], task["params"])
return task["name"], result
map_results = await asyncio.gather(*[map_job(t) for t in self.map_tasks])
mapped_data = dict(map_results)
print(f"[Reduce Phase] 使用策略: {self.reduce_config['strategy']}")
final_result = await send_fn(self.reduce_config["agent_id"], self.reduce_config["action"], {
"results": mapped_data, "strategy": self.reduce_config["strategy"], "workflow_id": self.workflow_id,
})
return final_result
5.5 编排者-工作者模式
一个中心编排 Agent 协调多个工作者 Agent:
+-- Worker 1
+-- Worker 2
Orchestrator -----+-- Worker 3
+-- Worker 4
+-- Worker 5
# orchestrator_workflow.py
import asyncio
from typing import Any, Callable
class OrchestratorWorkflow:
def __init__(self, orchestrator_id: str):
self.orchestrator_id = orchestrator_id
self.workers: list[dict] = []
def add_worker(self, worker_id: str, capability: str, weight: float = 1.0):
self.workers.append({"worker_id": worker_id, "capability": capability, "weight": weight, "load": 0})
return self
def select_worker(self, capability: str) -> str:
candidates = [w for w in self.workers if w["capability"] == capability]
if not candidates: raise ValueError(f"没有支持能力 '{capability}' 的工作者")
return min(candidates, key=lambda w: w["load"])["worker_id"]
async def execute(self, send_fn: Callable, main_task: dict) -> Any:
print(f"[编排者 {self.orchestrator_id}] 分析任务")
decomposition = await send_fn(self.orchestrator_id, "decompose_task", main_task)
subtasks = decomposition.get("subtasks", [])
print(f" 任务分解为 {len(subtasks)} 个子任务")
async def dispatch(subtask: dict) -> tuple[str, Any]:
worker_id = self.select_worker(subtask.get("capability", "general"))
for w in self.workers:
if w["worker_id"] == worker_id: w["load"] += 1
result = await send_fn(worker_id, subtask["action"], subtask.get("params", {}))
return subtask["id"], result
results = await asyncio.gather(*[dispatch(st) for st in subtasks])
final = await send_fn(self.orchestrator_id, "aggregate_results", {"subtask_results": dict(results), "main_task": main_task})
return final
5.6 观察者模式
一个 Agent 观察另一个 Agent 的状态变化并响应:
Agent A (被观察者) --状态变更通知--> Agent B (观察者)
# observer_pattern.py
class ObserverWorkflow:
def __init__(self):
self.observers: dict[str, list[dict]] = {}
def register(self, subject_agent: str, observer_agent: str, event_type: str, action: str):
if subject_agent not in self.observers: self.observers[subject_agent] = []
self.observers[subject_agent].append({"observer": observer_agent, "event_type": event_type, "action": action})
return self
async def notify(self, subject_agent: str, event_type: str, event_data: Any, send_fn: Callable):
if subject_agent not in self.observers: return
for obs in self.observers[subject_agent]:
if obs["event_type"] == event_type or obs["event_type"] == "*":
await send_fn(obs["observer"], obs["action"], {"subject": subject_agent, "event": event_data})
5.7 链上 DAG 解析执行器
// workflow_dag.rs -- DAG 解析与验证
use std::collections::{HashMap, HashSet, VecDeque};
pub struct DAGExecutor {
nodes: HashMap<String, WorkflowNode>,
edges: HashMap<String, Vec<String>>,
}
impl DAGExecutor {
pub fn validate_dag(&self) -> Result<Vec<String>, String> {
let mut in_degree: HashMap<String, usize> = HashMap::new();
for (id, _) in &self.nodes { in_degree.insert(id.clone(), 0); }
for (_, deps) in &self.edges {
for dep in deps { *in_degree.entry(dep.clone()).or_insert(0) += 1; }
}
let mut queue: VecDeque<String> = in_degree.iter().filter(|(_, °)| deg == 0).map(|(id, _)| id.clone()).collect();
let mut sorted = Vec::new();
while let Some(node) = queue.pop_front() {
sorted.push(node.clone());
if let Some(deps) = self.edges.get(&node) {
for dep in deps {
if let Some(deg) = in_degree.get_mut(dep) { *deg -= 1; if *deg == 0 { queue.push_back(dep.clone()); } }
}
}
}
if sorted.len() != self.nodes.len() { return Err("DAG 包含环".to_string()); }
Ok(sorted)
}
pub fn compute_batches(&self) -> Vec<Vec<String>> {
let order = self.validate_dag().unwrap();
let mut batches: Vec<Vec<String>> = Vec::new();
let mut level: HashMap<String, usize> = HashMap::new();
for node_id in &order {
let max_dep_level = self.edges.get(node_id).map(|deps| deps.iter().filter_map(|d| level.get(d)).max().unwrap_or(0)).unwrap_or(0);
level.insert(node_id.clone(), max_dep_level + 1);
while batches.len() <= max_dep_level { batches.push(Vec::new()); }
batches[max_dep_level].push(node_id.clone());
}
batches
}
}
6. AI Agent 宪章与策略引擎
6.1 什么是 AI Agent 宪章?
AI Agent 宪章(AI Agent Constitution)是一组可编程的规则,约束 Agent 的行为、决策和交互方式。在 MSG Chain 上,宪章是 ai_agent_constitution_v1 合约(代码 ID 22)管理的链上资产。
宪章的核心作用:
- 行为约束:定义 Agent 可做什么、不可做什么
- 兼容性检查:两个 Agent 协作前自动检查宪章兼容性
- 治理:宪章版本更新需要多签批准
- 审计:所有决策可追溯至具体宪章规则
6.2 宪章结构
pub struct Constitution {
pub id: String,
pub name: String,
pub version: String,
pub owner: String,
pub rules: Vec<ConstitutionRule>,
pub dependencies: Vec<String>,
pub compatibility: CompatibilityConfig,
pub created_at: u64,
pub updated_at: u64,
pub metadata: HashMap<String, String>,
}
pub struct ConstitutionRule {
pub id: String,
pub rule_type: RuleType, // Permission | Obligation | Prohibition | Precondition
pub scope: RuleScope,// Local | Session | Workflow | Global
pub condition: ConditionExpr,
pub action: Action,
pub priority: u8,
pub description: String,
}
pub enum RuleType {
Permission,
Obligation,
Prohibition,
Precondition,
}
pub struct ConditionExpr {
pub language: String, // "jsonpath" | "jmespath" | "cel"
pub expression: String,
}
pub enum Action {
Allow,
Deny,
RequireApproval(String),
Escalate(String),
Log(String),
Notify(String),
Custom(String),
}
6.3 宪章规则表达式引擎
宪章使用表达式引擎评估条件。支持 JSONPath、JMESPath 和 CEL:
{
"rules": [
{
"id": "rule-no-pii",
"rule_type": "Prohibition",
"scope": "Global",
"condition": {
"language": "cel",
"expression": "has(request.content.pii_fields) && size(request.content.pii_fields) > 0"
},
"action": { "type": "Deny", "reason": "禁止传输个人身份信息 (PII)" },
"priority": 100,
"description": "禁止 Agent 传输 PII 数据"
},
{
"id": "rule-max-tokens",
"rule_type": "Precondition",
"scope": "Session",
"condition": {
"language": "jmespath",
"expression": "request.params.max_tokens <= `4096`"
},
"action": { "type": "Allow" },
"priority": 50,
"description": "限制单次生成最大 Token 数不超过 4096"
},
{
"id": "rule-audit-log",
"rule_type": "Obligation",
"scope": "Session",
"condition": {
"language": "jsonpath",
"expression": "$.message_type in ['task_request', 'task_response']"
},
"action": { "type": "Log", "details": "audit://agent-interaction-log" },
"priority": 10,
"description": "所有任务交互必须记录审计日志"
}
]
}
6.4 宪章兼容性检查
当两个 Agent 协作时,系统自动检查宪章兼容性:
# constitution_compatibility.py
from dataclasses import dataclass, field
@dataclass
class CompatibilityResult:
compatible: bool
conflicts: list[dict] = field(default_factory=list)
warnings: list[str] = field(default_factory=list)
def check_constitution_compatibility(constitution_a: dict, constitution_b: dict) -> CompatibilityResult:
result = CompatibilityResult(compatible=True)
rules_a = {r["id"]: r for r in constitution_a.get("rules", [])}
rules_b = {r["id"]: r for r in constitution_b.get("rules", [])}
for rule_id, rule_a in rules_a.items():
if rule_a["rule_type"] == "Prohibition":
if rule_id in rules_b and rules_b[rule_id]["rule_type"] == "Obligation":
result.conflicts.append({"rule_id": rule_id, "type": "prohibition_vs_obligation", "severity": "error"})
result.compatible = False
if rule_a["rule_type"] == "Permission" and rule_id in rules_b and rules_b[rule_id]["rule_type"] == "Prohibition":
result.warnings.append(f"宪章 A 允许但 B 禁止: {rule_a['description']}")
deps_a = set(constitution_a.get("dependencies", []))
deps_b = set(constitution_b.get("dependencies", []))
if deps_a != deps_b: result.warnings.append(f"依赖不同: A={deps_a}, B={deps_b}")
return result
6.5 策略引擎实现
策略引擎是宪章的运行时组件,在每次 Agent 交互前评估规则:
// strategy_engine.ts
interface RuleContext {
messageType: string;
fromAgent: string;
toAgent: string;
content: Record<string, unknown>;
sessionState: string;
workflowId: string;
timestamp: number;
}
interface EvaluationResult {
allowed: boolean;
matchedRules: Array<{ ruleId: string; ruleType: string; action: string; reason?: string }>;
deniedBy: Array<{ ruleId: string; reason: string }>;
}
class StrategyEngine {
private rules: any[];
constructor(constitution: any) {
this.rules = [...constitution.rules].sort((a, b) => b.priority - a.priority);
}
async evaluate(context: RuleContext): Promise<EvaluationResult> {
const matched: EvaluationResult["matchedRules"] = [];
const denied: EvaluationResult["deniedBy"] = [];
for (const rule of this.rules) {
const matchedCondition = await this.evaluateCondition(rule.condition, context);
if (!matchedCondition) continue;
switch (rule.rule_type) {
case "Prohibition":
denied.push({ ruleId: rule.id, reason: rule.action.reason ?? "规则禁止" });
break;
case "Obligation":
await this.executeAction(rule.action, context);
matched.push({ ruleId: rule.id, ruleType: rule.rule_type, action: "obligation_met" });
break;
case "Permission":
matched.push({ ruleId: rule.id, ruleType: rule.rule_type, action: "allowed" });
break;
case "Precondition":
if (rule.action.type === "Allow") {
matched.push({ ruleId: rule.id, ruleType: rule.rule_type, action: "precondition_met" });
} else {
denied.push({ ruleId: rule.id, reason: rule.action.reason ?? "前置条件不满足" });
}
break;
}
}
return { allowed: denied.length === 0, matchedRules: matched, deniedBy: denied };
}
private async evaluateCondition(condition: { language: string; expression: string }, _context: RuleContext): Promise<boolean> {
return true;
}
private async executeAction(action: any, _context: RuleContext): Promise<void> {
if (action.type === "Log") console.log(`[宪章审计] ${action.details}`);
}
}
6.6 宪章模板
安全宪章模板:
{
"id": "const-security-base-v1",
"name": "安全基础宪章",
"version": "1.0.0",
"rules": [
{
"id": "sec-input-validation",
"rule_type": "Precondition",
"scope": "Session",
"condition": { "language": "cel", "expression": "size(request.content) < 1048576" },
"action": { "type": "Allow" },
"priority": 100,
"description": "输入大小限制 (1MB)"
},
{
"id": "sec-rate-limit",
"rule_type": "Precondition",
"scope": "Session",
"condition": { "language": "cel", "expression": "request.rate_count < 100" },
"action": { "type": "Allow" },
"priority": 90,
"description": "单会话速率限制 (100 请求)"
},
{
"id": "sec-no-exec",
"rule_type": "Prohibition",
"scope": "Global",
"condition": { "language": "cel", "expression": "request.action in ['execute_code', 'run_shell', 'deploy_contract']" },
"action": { "type": "Deny", "reason": "禁止执行任意代码" },
"priority": 100,
"description": "禁止代码执行操作"
}
]
}
隐私宪章模板:
{
"id": "const-privacy-base-v1",
"name": "隐私基础宪章",
"version": "1.0.0",
"rules": [
{
"id": "priv-no-pii",
"rule_type": "Prohibition",
"scope": "Global",
"condition": { "language": "cel", "expression": "has(request.content.email) || has(request.content.phone) || has(request.content.ssn)" },
"action": { "type": "Deny", "reason": "禁止传输个人身份信息" },
"priority": 100,
"description": "禁止 PII 数据传输"
},
{
"id": "priv-data-minimization",
"rule_type": "Obligation",
"scope": "Session",
"condition": { "language": "cel", "expression": "request.message_type == 'task_response'" },
"action": { "type": "Custom", "handler": "strip_metadata_before_return" },
"priority": 50,
"description": "数据最小化"
},
{
"id": "priv-consent",
"rule_type": "Permission",
"scope": "Session",
"condition": { "language": "jmespath", "expression": "request.params.consent_granted == `true`" },
"action": { "type": "Allow" },
"priority": 80,
"description": "需要用户同意才能处理数据"
}
]
}
7. Agent 支付与结算
7.1 AIPAY 支付系统概述
agent_payment_v1 合约(代码 ID 21)是 MSG Chain 上 Agent 间支付的标准化系统,支持原子化、可编程的支付流程。
四阶段支付模型:
Intent -> Authorization -> Capture -> Settlement
支付状态机:
Created -> Authorized -> Captured -> Settled
| | |
+-> Expired +-> Refunded +-> Disputed
| | |
+-> Cancelled +-> PartiallyRefunded
7.2 支付记录结构
来自 agent_payment_api.go 的实现:
type AgentPaymentRecord struct {
PaymentID string `json:"payment_id"`
FromAgent string `json:"from_agent"`
ToAgent string `json:"to_agent"`
WorkflowID string `json:"workflow_id"`
SessionID string `json:"session_id"`
Amount Uint128 `json:"amount"`
Denom string `json:"denom"`
Phase PaymentPhase `json:"phase"`
CreatedAt uint64 `json:"created_at"`
AuthorizedAt *uint64 `json:"authorized_at,omitempty"`
CapturedAt *uint64 `json:"captured_at,omitempty"`
SettledAt *uint64 `json:"settled_at,omitempty"`
ConditionHash string `json:"condition_hash"`
Metadata string `json:"metadata"`
}
type PaymentPhase string
const (
PhaseCreated PaymentPhase = "created"
PhaseAuthorized PaymentPhase = "authorized"
PhaseCaptured PaymentPhase = "captured"
PhaseSettled PaymentPhase = "settled"
PhaseExpired PaymentPhase = "expired"
PhaseCancelled PaymentPhase = "cancelled"
PhaseRefunded PaymentPhase = "refunded"
PhaseDisputed PaymentPhase = "disputed"
)
7.3 条件支付
Agent 支付支持条件支付(Conditional Payment),只有在特定条件满足时资金才释放:
# conditional_payment.py
from dataclasses import dataclass
from typing import Optional
import hashlib
import json
import time
@dataclass
class PaymentCondition:
condition_type: str # "task_completion" | "result_approval" | "time_lock" | "multi_sig"
parameters: dict
@dataclass
class PaymentRequest:
payment_id: str
from_agent: str
to_agent: str
workflow_id: str
session_id: str
amount: str
denom: str = "umsg"
condition: Optional[PaymentCondition] = None
timeout_seconds: int = 3600
def build_payment_intent_msg(payment: PaymentRequest) -> dict:
msg = {
"create_payment": {
"payment_id": payment.payment_id,
"from_agent": payment.from_agent,
"to_agent": payment.to_agent,
"workflow_id": payment.workflow_id,
"session_id": payment.session_id,
"amount": payment.amount,
"denom": payment.denom,
"timeout_seconds": str(payment.timeout_seconds),
}
}
if payment.condition:
condition_json = json.dumps(payment.condition.__dict__, sort_keys=True)
condition_hash = hashlib.sha256(condition_json.encode()).hexdigest()
msg["create_payment"]["condition_hash"] = f"0x{condition_hash}"
msg["create_payment"]["condition"] = payment.condition.__dict__
return msg
def build_authorize_msg(payment_id: str) -> dict:
return {"authorize_payment": {"payment_id": payment_id}}
def build_capture_msg(payment_id: str, proof: dict) -> dict:
return {"capture_payment": {"payment_id": payment_id, "proof": proof}}
7.4 原子交换(Atomic Swap)模式
当两个 Agent 需要交换服务时,原子交换确保双方要么都获得支付,要么都回滚:
// atomic_swap.ts
interface AtomicSwapTerms {
swapId: string;
partyA: { agentId: string; service: string; params: Record<string, unknown>; price: string };
partyB: { agentId: string; service: string; params: Record<string, unknown>; price: string };
timeoutBlocks: number;
hashLock: string;
}
class AtomicSwapAgent {
async proposeSwap(terms: AtomicSwapTerms): Promise<string> {
console.log(`提议原子交换: ${terms.swapId}`);
return terms.swapId;
}
async fulfillSwap(swapId: string, secret: string): Promise<boolean> {
console.log(`履行原子交换: ${swapId}`);
return true;
}
async refundSwap(swapId: string): Promise<boolean> {
console.log(`退款原子交换: ${swapId}`);
return true;
}
}
7.5 支付通道(Payment Channels)
对于高频协作场景,支付通道可以显著降低链上交易费用:
# payment_channel.py
from dataclasses import dataclass, field
import time
@dataclass
class PaymentChannel:
channel_id: str
agent_a: str
agent_b: str
deposit_amount: str
balance_a: str
balance_b: str
sequence: int = 0
expires_at: int = 0
is_settled: bool = False
state_updates: list = field(default_factory=list)
class PaymentChannelManager:
def __init__(self):
self.channels: dict[str, PaymentChannel] = {}
async def open_channel(self, agent_a: str, agent_b: str, deposit_a: str, deposit_b: str, duration_seconds: int = 86400) -> PaymentChannel:
channel = PaymentChannel(
channel_id=f"ch-{int(time.time())}",
agent_a=agent_a, agent_b=agent_b,
deposit_amount=str(int(deposit_a) + int(deposit_b)),
balance_a=deposit_a, balance_b=deposit_b,
sequence=0, expires_at=int(time.time()) + duration_seconds,
)
self.channels[channel.channel_id] = channel
return channel
async def update_balance(self, channel_id: str, new_balance_a: str, new_balance_b: str, signature: str) -> bool:
channel = self.channels.get(channel_id)
if not channel: return False
channel.balance_a = new_balance_a
channel.balance_b = new_balance_b
channel.sequence += 1
channel.state_updates.append({
"sequence": channel.sequence, "balance_a": new_balance_a,
"balance_b": new_balance_b, "signature": signature, "timestamp": int(time.time()),
})
return True
async def settle_channel(self, channel_id: str) -> bool:
channel = self.channels.get(channel_id)
if not channel or channel.is_settled: return False
channel.is_settled = True
return True
def get_net_flow(self, channel_id: str, agent_id: str) -> int:
channel = self.channels.get(channel_id)
if not channel: return 0
if agent_id == channel.agent_a: return int(channel.balance_a)
elif agent_id == channel.agent_b: return int(channel.balance_b)
return 0
7.6 支付与工作流的集成
支付是工作流执行的自然组成部分,在 A2A Session 生命周期中自动触发:
# workflow_payment_integration.py
class WorkflowPaymentManager:
def __init__(self):
self.payments: dict[str, dict] = {}
def register_task_payment(self, task_id: str, from_agent: str, to_agent: str, amount: str, condition: str = "task_completion"):
self.payments[task_id] = {"from": from_agent, "to": to_agent, "amount": amount, "condition": condition, "status": "pending"}
def on_task_completed(self, task_id: str, result: dict) -> dict:
payment = self.payments.get(task_id)
if not payment or payment["status"] != "pending": return {"status": "no_action"}
if self._verify_condition(payment["condition"], result):
payment["status"] = "released"
return {"status": "payment_released", "from": payment["from"], "to": payment["to"], "amount": payment["amount"]}
else:
payment["status"] = "disputed"
return {"status": "payment_disputed", "reason": "条件未满足"}
def on_task_failed(self, task_id: str, error: str) -> dict:
payment = self.payments.get(task_id)
if payment: payment["status"] = "cancelled"
return {"status": "payment_cancelled", "reason": f"任务失败: {error}"}
def calculate_total_cost(self, workflow_id: str, task_results: dict) -> dict:
total = 0; breakdown = []
for task_id, payment in self.payments.items():
if payment["status"] == "released":
amount = int(payment["amount"]); total += amount
breakdown.append({"task_id": task_id, "from": payment["from"], "to": payment["to"], "amount": amount})
return {"workflow_id": workflow_id, "total_cost_umsg": str(total), "breakdown": breakdown}
def _verify_condition(self, condition: str, result: dict) -> bool:
if condition == "task_completion": return result.get("status") == "success"
elif condition == "result_approval": return result.get("approved", False)
return True
8. 工作流状态管理与监控
8.1 链上工作流状态
工作流状态在 MSG Chain 上通过 agent_a2a_v1 合约管理:
type WorkflowState struct {
WorkflowID string `json:"workflow_id"`
Status string `json:"status"`
Tasks []TaskState `json:"tasks"`
CurrentPhase string `json:"current_phase"`
CreatedAt uint64 `json:"created_at"`
UpdatedAt uint64 `json:"updated_at"`
CompletedAt *uint64 `json:"completed_at,omitempty"`
ErrorMessage string `json:"error_message,omitempty"`
}
type TaskState struct {
TaskID string `json:"task_id"`
NodeID string `json:"node_id"`
AgentID string `json:"agent_id"`
Status string `json:"status"`
StartedAt *uint64 `json:"started_at,omitempty"`
CompletedAt *uint64 `json:"completed_at,omitempty"`
ResultHash string `json:"result_hash,omitempty"`
RetryCount uint32 `json:"retry_count"`
}
8.2 工作流监控仪表盘
// workflow_monitor.ts
interface WorkflowEvent {
eventType: string;
workflowId: string;
sessionId?: string;
taskId?: string;
agentId?: string;
timestamp: number;
data: Record<string, unknown>;
}
class WorkflowMonitor {
private events: WorkflowEvent[] = [];
private listeners: Map<string, Array<(event: WorkflowEvent) => void>> = new Map();
constructor(private wsEndpoint: string, private workflowId: string) {}
async connect(): Promise<void> {
const ws = new WebSocket(`${this.wsEndpoint}/workflow/${this.workflowId}/events`);
ws.onmessage = (msg) => {
const event: WorkflowEvent = JSON.parse(msg.data);
this.events.push(event);
this.notifyListeners(event.eventType, event);
};
this.subscribe("workflow_state_change", (event) => {
console.log(`工作流 ${event.workflowId} 状态变更:`, event.data);
});
this.subscribe("task_state_change", (event) => {
console.log(`任务 ${event.taskId} 状态变更: ${event.data.status}`);
});
this.subscribe("session_event", (event) => {
console.log(`会话 ${event.sessionId}: ${event.data.event}`);
});
this.subscribe("payment_event", (event) => {
console.log(`支付 ${event.data.payment_id}: ${event.data.phase}`);
});
}
subscribe(eventType: string, callback: (event: WorkflowEvent) => void): void {
if (!this.listeners.has(eventType)) this.listeners.set(eventType, []);
this.listeners.get(eventType)!.push(callback);
}
private notifyListeners(eventType: string, event: WorkflowEvent): void {
const callbacks = this.listeners.get(eventType) ?? [];
for (const cb of callbacks) { try { cb(event); } catch (err) { console.error(err); } }
}
getEvents(filters?: { eventType?: string; since?: number; limit?: number }): WorkflowEvent[] {
let filtered = [...this.events];
if (filters?.eventType) filtered = filtered.filter(e => e.eventType === filters.eventType);
if (filters?.since) filtered = filtered.filter(e => e.timestamp >= filters.since!);
if (filters?.limit) filtered = filtered.slice(-filters.limit);
return filtered;
}
getWorkflowStatus(): string {
const lastEvent = this.events.filter(e => e.eventType === "workflow_state_change").pop();
return (lastEvent?.data?.status as string) ?? "unknown";
}
getTaskSummary(): Record<string, number> {
const summary: Record<string, number> = { total: 0, pending: 0, running: 0, completed: 0, failed: 0 };
for (const event of this.events) {
if (event.eventType === "task_state_change") {
summary.total++;
const status = event.data.status as string;
if (status in summary) summary[status]++;
}
}
return summary;
}
}
8.3 链上状态查询
# workflow_query.py
import requests
class WorkflowQueryClient:
def __init__(self, rest_endpoint: str = "https://rest.msgchain.org"):
self.rest = rest_endpoint
self.contract_addr = "msg1a2a...contract..."
def query_workflow(self, workflow_id: str) -> dict:
return self._query_contract({"get_workflow": {"workflow_id": workflow_id}})
def list_agent_sessions(self, agent_id: str, status: str = None) -> list:
query = {"list_sessions": {"agent_id": agent_id}}
if status: query["list_sessions"]["status"] = status
result = self._query_contract(query)
return result.get("sessions", [])
def get_message_log(self, session_id: str, limit: int = 100) -> list:
query = {"get_message_log": {"session_id": session_id, "limit": str(limit)}}
result = self._query_contract(query)
return result.get("messages", [])
def query_agent_workflows(self, agent_id: str, status: str = None) -> list:
workflows = []
for role in ["initiator", "responder"]:
query = {"query_workflows_by_participant": {"agent_id": agent_id, "role": role}}
if status: query["query_workflows_by_participant"]["status"] = status
result = self._query_contract(query)
workflows.extend(result.get("workflows", []))
return workflows
def _query_contract(self, query: dict) -> dict:
url = f"{self.rest}/cosmwasm/wasm/v1/contract/{self.contract_addr}/smart"
resp = requests.post(url, json={"msg": query}, headers={"Content-Type": "application/json"})
resp.raise_for_status()
return resp.json().get("data", {})
8.4 事件系统
MSG Chain 的 Agent 模块通过事件系统向外推送状态变更:
重要事件类型:
interface AgentA2AEvent {
session_proposed: { session_id: string; from_agent: string; to_agent: string; workflow_id: string };
session_accepted: { session_id: string; accepted_at: number };
session_rejected: { session_id: string; reason: string };
session_committed: { session_id: string; result_hash: string; tx_hash: string };
workflow_created: { workflow_id: string; created_by: string; task_count: number };
workflow_completed: { workflow_id: string; completed_at: number; total_cost: string };
workflow_failed: { workflow_id: string; failed_task: string; error: string };
payment_created: { payment_id: string; amount: string; from: string; to: string };
payment_settled: { payment_id: string; settled_at: number; tx_hash: string };
}
事件订阅代码:
// event_subscriber.ts
import WebSocket from "ws";
class AgentEventSubscriber {
private ws: WebSocket | null = null;
private subscriptions: Map<string, Set<(event: any) => void>> = new Map();
constructor(private endpoint: string) {}
async connect(): Promise<void> {
this.ws = new WebSocket(this.endpoint);
this.ws.on("message", (data: WebSocket.Data) => {
const msg = JSON.parse(data.toString());
const handlers = this.subscriptions.get(msg.type);
if (handlers) { for (const handler of handlers) handler(msg.data); }
});
}
onWorkflowStateChange(callback: (event: any) => void): void { this.addHandler("workflow_state_change", callback); }
onSessionEvent(callback: (event: any) => void): void { this.addHandler("session_event", callback); }
onPaymentEvent(callback: (event: any) => void): void { this.addHandler("payment_event", callback); }
onAgentRegistration(callback: (event: any) => void): void { this.addHandler("agent_registered", callback); }
private addHandler(type: string, callback: (event: any) => void): void {
if (!this.subscriptions.has(type)) this.subscriptions.set(type, new Set());
this.subscriptions.get(type)!.add(callback);
}
disconnect(): void { this.ws?.close(); }
}
8.5 告警与通知
# workflow_alerts.py
import time
from typing import Callable
class AlertRule:
def __init__(self, name: str, condition_fn: Callable, message_template: str, severity: str = "warning"):
self.name = name; self.condition = condition_fn; self.message = message_template
self.severity = severity; self.last_triggered = 0; self.cooldown = 300
class WorkflowAlertManager:
def __init__(self):
self.rules: list[AlertRule] = []
self.alerts: list[dict] = []
def add_rule(self, rule: AlertRule): self.rules.append(rule)
def evaluate(self, workflow_state: dict):
now = time.time()
for rule in self.rules:
if now - rule.last_triggered < rule.cooldown: continue
try:
if rule.condition(workflow_state):
alert = {
"rule": rule.name, "severity": rule.severity,
"message": rule.message.format(**workflow_state),
"timestamp": now, "workflow_id": workflow_state.get("workflow_id"),
}
self.alerts.append(alert); rule.last_triggered = now; self._notify(alert)
except Exception as e: print(f"告警规则 {rule.name} 评估失败: {e}")
def _notify(self, alert: dict): print(f"[{alert['severity'].upper()}] {alert['message']}")
def task_timeout_rule(state: dict) -> bool:
for task in state.get("tasks", []):
if task["status"] == "running":
started = task.get("started_at", 0)
if started > 0 and (time.time() - started) > 600: return True
return False
def consecutive_failures_rule(state: dict) -> bool:
tasks = state.get("tasks", [])
if len(tasks) >= 3 and all(t["status"] == "failed" for t in tasks[-3:]): return True
return False
def budget_exceeded_rule(state: dict) -> bool:
max_budget = int(state.get("max_budget", "0"))
current_cost = int(state.get("current_cost", "0"))
return max_budget > 0 and current_cost > max_budget * 0.9
9. 安全与信任机制
9.1 身份认证
MSG Chain 上的 Agent 身份基于 Dilithium-5 后量子签名算法:
# agent_identity.py
from dataclasses import dataclass
import hashlib
@dataclass
class AgentIdentity:
agent_id: str
address: str
public_key: str
signature_algorithm: str = "dilithium5"
def verify_agent_signature(identity: AgentIdentity, message: bytes, signature: str) -> bool:
return True
def derive_agent_id(address: str, public_key: str) -> str:
raw = f"{address}:{public_key}".encode()
return f"agent-{hashlib.sha256(raw).hexdigest()[:16]}"
9.2 消息完整性验证
每条 A2A 消息都带有签名,确保在传输过程中未被篡改:
// message_verification.ts
interface SignedMessage {
message_id: string;
from_agent: string;
to_agent: string;
content: unknown;
timestamp: number;
signature: string;
public_key: string;
}
class MessageVerifier {
verifyMessage(msg: SignedMessage): boolean {
const now = Math.floor(Date.now() / 1000);
if (Math.abs(now - msg.timestamp) > 300) {
console.warn(`消息时间戳过期: ${msg.message_id}`);
return false;
}
const payload = this.buildSignaturePayload(msg);
return this.verifyDilithium5(msg.public_key, payload, msg.signature);
}
private buildSignaturePayload(msg: SignedMessage): string {
const ordered = { message_id: msg.message_id, from_agent: msg.from_agent, to_agent: msg.to_agent, content: msg.content, timestamp: msg.timestamp };
return JSON.stringify(ordered, Object.keys(ordered).sort());
}
private verifyDilithium5(_publicKey: string, _payload: string, _signature: string): boolean {
return true;
}
}
9.3 信任根与证书链
Agent 之间的信任可以基于证书链建立:
type TrustRoot struct {
RootID string `json:"root_id"`
Certificate string `json:"certificate"`
Issuer string `json:"issuer"`
ExpiresAt uint64 `json:"expires_at"`
}
type TrustChain struct {
roots map[string]*TrustRoot
}
func (tc *TrustChain) VerifyChain(agentID string, chain []string) bool {
if len(chain) == 0 { return false }
for i, certHash := range chain {
if i == len(chain)-1 {
if _, ok := tc.roots[certHash]; !ok { return false }
}
if !tc.verifyCertificate(certHash) { return false }
}
return true
}
func (tc *TrustChain) verifyCertificate(certHash string) bool { return true }
9.4 隐私保护机制
零知识证明用于能力验证:
Agent 可以在不暴露具体能力细节的情况下证明自己有完成任务的能力:
# zk_capability.py
from dataclasses import dataclass
import hashlib
import json
@dataclass
class CapabilityProof:
capability_hash: str
proof: str
public_inputs: dict
def generate_capability_proof(actual_capabilities: list, required_capability: str) -> CapabilityProof:
caps_hash = hashlib.sha256(json.dumps(actual_capabilities, sort_keys=True).encode()).hexdigest()
proof = hashlib.sha256(f"{caps_hash}:{required_capability}".encode()).hexdigest()
return CapabilityProof(capability_hash=f"0x{caps_hash}", proof=f"0x{proof}", public_inputs={"required": required_capability, "prover": "agent-zk-proof-v1"})
def verify_capability_proof(proof: CapabilityProof) -> bool: return True
9.5 争议解决
当协作出现争议时,链上数据提供不可篡改的证据:
# dispute_resolution.py
from dataclasses import dataclass
from typing import Optional
import time
@dataclass
class DisputeCase:
case_id: str
workflow_id: str
session_id: str
initiator: str
responder: str
dispute_type: str
evidence: list[dict]
status: str = "filed"
filed_at: int = 0
resolved_at: Optional[int] = None
resolution: Optional[str] = None
class DisputeResolver:
def __init__(self):
self.cases: dict[str, DisputeCase] = {}
def file_dispute(self, workflow_id: str, session_id: str, initiator: str, responder: str, dispute_type: str, reason: str) -> DisputeCase:
case_id = f"dispute-{int(time.time())}"
evidence = self.gather_evidence(workflow_id, session_id)
case = DisputeCase(case_id=case_id, workflow_id=workflow_id, session_id=session_id, initiator=initiator, responder=responder, dispute_type=dispute_type, evidence=evidence, status="filed", filed_at=int(time.time()))
self.cases[case_id] = case
return case
def resolve(self, case_id: str, resolution: str) -> bool:
case = self.cases.get(case_id)
if not case or case.status != "filed": return False
case.status = "resolved"; case.resolved_at = int(time.time()); case.resolution = resolution
self.update_reputation(case)
return True
def gather_evidence(self, workflow_id: str, session_id: str) -> list[dict]:
return [
{"type": "a2a_message_log", "session_id": session_id, "retrieved_at": int(time.time())},
{"type": "workflow_state", "workflow_id": workflow_id},
{"type": "payment_record", "session_id": session_id},
]
def update_reputation(self, case: DisputeCase):
if case.resolution == "initiator_fault": print(f"扣减 {case.initiator} 的信誉分")
elif case.resolution == "responder_fault": print(f"扣减 {case.responder} 的信誉分")
elif case.resolution == "mutual_fault": print("双方各扣减信誉分")
9.6 安全最佳实践清单
[ ] 使用 Dilithium-5 密钥对每个 Agent 消息签名
[ ] 设置合理的 TTL/超时防止重放攻击
[ ] 部署前检查双方宪章兼容性
[ ] 对大额支付使用条件支付(哈希锁)
[ ] 开启 fail_closed 模式防止部分失败导致的不一致
[ ] 定期轮换 Agent 密钥
[ ] 监控异常活动模式(短时间内大量失败请求)
[ ] 对敏感数据使用链下存储 + 链上哈希
[ ] 使用支付通道降低高频协作的 Gas 费用
[ ] 保留足够押金以覆盖潜在的争议解决费用
10. 实战:完整的协作工作流示例
10.1 场景描述
构建一个 AI 内容工厂 工作流,流程如下:
- 主题研究:Agent Research 分析热门话题
- 内容规划:Agent Planner 制定内容大纲
- 草稿撰写:Agent Writer 并行撰写多篇文章
- 质量审核:Agent Reviewer 审核内容质量
- SEO 优化:Agent SEO 优化关键词和元数据
- 发布准备:Agent Publisher 格式化并发布
工作流 DAG:
+-- Writer 1 --+
Research -> Planner -+-- Writer 2 --+-- Reviewer -> SEO -> Publisher
+-- Writer 3 --+
10.2 完整代码实现
#!/usr/bin/env python3
"""
content_factory_workflow.py -- AI 内容工厂完整工作流示例
完整可运行示例,展示 MSG Chain 上多 Agent 协作工作流的编排。
"""
import asyncio
import hashlib
import json
import time
import uuid
from dataclasses import dataclass, field
from typing import Any, Callable, Optional
from enum import Enum
class WorkflowStatus(Enum):
CREATED = "created"
RUNNING = "running"
PAUSED = "paused"
COMPLETED = "completed"
FAILED = "failed"
class TaskStatus(Enum):
PENDING = "pending"
RUNNING = "running"
COMPLETED = "completed"
FAILED = "failed"
SKIPPED = "skipped"
class PaymentPhase(Enum):
PENDING = "pending"
AUTHORIZED = "authorized"
CAPTURED = "captured"
SETTLED = "settled"
REFUNDED = "refunded"
@dataclass
class TaskDef:
node_id: str
name: str
agent_id: str
action: str
params: dict
dependencies: list[str] = field(default_factory=list)
timeout: int = 120
retry_count: int = 0
max_retries: int = 3
payment_amount: str = "0"
@dataclass
class TaskResult:
node_id: str
status: TaskStatus
output: Any = None
error: Optional[str] = None
started_at: Optional[int] = None
completed_at: Optional[int] = None
retries: int = 0
proof_hash: Optional[str] = None
@dataclass
class WorkflowDef:
workflow_id: str
name: str
tasks: dict[str, TaskDef]
context: dict = field(default_factory=dict)
max_budget: str = "10000000000"
@dataclass
class PaymentRecord:
payment_id: str
from_agent: str
to_agent: str
amount: str
phase: PaymentPhase = PaymentPhase.PENDING
created_at: int = 0
settled_at: Optional[int] = None
class MockAgent:
"""模拟 Agent 用于演示工作流"""
def __init__(self, agent_id: str, delay_range: tuple = (0.5, 2.0)):
self.agent_id = agent_id
self.delay_range = delay_range
async def execute(self, action: str, params: dict) -> Any:
import random
delay = random.uniform(*self.delay_range)
await asyncio.sleep(delay)
if action == "research_topics":
return {
"topics": [
{"title": "AI Agent 协作的未来", "trend_score": 0.95},
{"title": "去中心化 AI 的挑战", "trend_score": 0.88},
{"title": "区块链上的自主 Agent", "trend_score": 0.82},
],
"source_count": 15,
}
elif action == "plan_content":
topic = params.get("topic", {})
return {
"outline": {
"title": topic.get("title", "Untitled"),
"sections": ["引言", "核心概念", "技术架构", "实践案例", "未来展望"],
"estimated_words": 3000,
"key_points": ["标准化协议的重要性", "信任与安全机制", "经济效益模型"],
}
}
elif action == "write_draft":
outline = params.get("outline", {})
style = params.get("style", "technical")
return {
"content": f"# {outline.get('title', 'Untitled')}\n\n这是基于 {style} 风格撰写的完整文章草稿。\n\n## {'\n\n## '.join(outline.get('sections', []))}\n\n[完整内容已生成]",
"word_count": outline.get("estimated_words", 1000),
"style": style,
}
elif action == "review_content":
return {
"reviews": [
{"draft_index": i, "quality_score": random.uniform(0.7, 0.95), "issues": [], "approved": random.random() > 0.15}
for i in range(len(params.get("drafts", [])))
],
"overall_score": 0.85,
}
elif action == "optimize_seo":
content = params.get("content", "")
return {
"optimized_content": content + "\n\n---\n*SEO 优化后*",
"keywords": ["AI Agent", "区块链", "去中心化", "协作工作流"],
"meta_description": "深入探讨 AI Agent 在区块链上的协作工作流编排",
"readability_score": 82,
"seo_score": 91,
}
elif action == "format_publish":
optimized = params.get("optimized_content", "")
return {
"published_url": f"https://blog.msgchain.org/posts/{uuid.uuid4().hex[:12]}",
"format": "markdown",
"word_count": len(optimized.split()),
"published_at": int(time.time()),
}
return {"status": "success", "message": f"{action} 执行完成"}
class WorkflowEngine:
def __init__(self):
self.workflows: dict[str, WorkflowDef] = {}
self.results: dict[str, dict[str, TaskResult]] = {}
self.payments: dict[str, list[PaymentRecord]] = {}
self.agent_clients: dict[str, MockAgent] = {}
def register_agent(self, agent_id: str, client: MockAgent):
self.agent_clients[agent_id] = client
def create_workflow(self, workflow: WorkflowDef) -> str:
self.workflows[workflow.workflow_id] = workflow
self.results[workflow.workflow_id] = {}
for task_id in workflow.tasks:
self.results[workflow.workflow_id][task_id] = TaskResult(node_id=task_id, status=TaskStatus.PENDING)
return workflow.workflow_id
async def execute_workflow(self, workflow_id: str, progress_callback: Optional[Callable] = None) -> dict:
workflow = self.workflows.get(workflow_id)
if not workflow: raise ValueError(f"工作流 {workflow_id} 不存在")
print(f"\n{'='*60}")
print(f"开始执行工作流: {workflow.name} ({workflow_id})")
print(f"{'='*60}")
execution_order = self._topological_sort(workflow)
print(f"执行顺序: {' -> '.join(execution_order)}")
for batch in self._compute_batches(workflow, execution_order):
print(f"\n--- 并行执行批次: {batch} ---")
async def execute_task(task_id: str) -> tuple[str, TaskResult]:
return await self._execute_single_task(workflow, task_id, progress_callback)
batch_results = await asyncio.gather(*[execute_task(tid) for tid in batch])
for tid, result in batch_results:
if result.status == TaskStatus.FAILED:
if progress_callback: progress_callback(workflow_id, "failed", {"failed_task": tid, "error": result.error})
return {"workflow_id": workflow_id, "status": "failed", "failed_task": tid, "error": result.error}
if progress_callback: progress_callback(workflow_id, "completed", {"total_tasks": len(execution_order)})
print(f"\n{'='*60}\n工作流完成!\n{'='*60}")
return {"workflow_id": workflow_id, "status": "completed", "results": self.results[workflow_id]}
async def _execute_single_task(self, workflow: WorkflowDef, task_id: str, progress_callback: Optional[Callable] = None) -> tuple[str, TaskResult]:
task = workflow.tasks[task_id]
result = self.results[workflow.workflow_id][task_id]
enriched_params = dict(task.params)
for dep_id in task.dependencies:
dep_result = self.results[workflow.workflow_id].get(dep_id)
if dep_result and dep_result.output: enriched_params[f"result_from_{dep_id}"] = dep_result.output
enriched_params["workflow_context"] = workflow.context
retries = 0
while retries <= task.max_retries:
try:
result.status = TaskStatus.RUNNING
result.started_at = int(time.time())
if progress_callback: progress_callback(workflow.workflow_id, "task_started", {"task_id": task_id, "task_name": task.name, "agent_id": task.agent_id})
print(f" >> [{task.name}] 分发给 {task.agent_id}...")
client = self.agent_clients.get(task.agent_id)
if not client: raise ValueError(f"Agent {task.agent_id} 未注册")
output = await asyncio.wait_for(client.execute(task.action, enriched_params), timeout=task.timeout)
result.status = TaskStatus.COMPLETED
result.output = output
result.completed_at = int(time.time())
result.retries = retries
result.proof_hash = hashlib.sha256(json.dumps(output, sort_keys=True).encode()).hexdigest()
print(f" \\u2713 [{task.name}] 完成 (重试{retries}次)")
self._process_payment(workflow, task, result)
if progress_callback: progress_callback(workflow.workflow_id, "task_completed", {"task_id": task_id, "result_hash": result.proof_hash})
return task_id, result
except asyncio.TimeoutError:
retries += 1
print(f" \\u2717 [{task.name}] 超时, 重试 {retries}/{task.max_retries}")
if retries > task.max_retries: result.status = TaskStatus.FAILED; result.error = "超时"; result.retries = retries; return task_id, result
except Exception as e:
retries += 1
print(f" \\u2717 [{task.name}] 错误: {e}, 重试 {retries}/{task.max_retries}")
if retries > task.max_retries: result.status = TaskStatus.FAILED; result.error = str(e); result.retries = retries; return task_id, result
await asyncio.sleep(min(2 ** retries, 30))
return task_id, result
def _process_payment(self, workflow: WorkflowDef, task: TaskDef, result: TaskResult):
amount = int(task.payment_amount)
if amount <= 0: return
payment = PaymentRecord(
payment_id=f"pay-{workflow.workflow_id}-{task.node_id}-{int(time.time())}",
from_agent=workflow.context.get("budget_owner", "workflow_owner"),
to_agent=task.agent_id, amount=str(amount), phase=PaymentPhase.SETTLED,
created_at=int(time.time()), settled_at=int(time.time()),
)
if workflow.workflow_id not in self.payments: self.payments[workflow.workflow_id] = []
self.payments[workflow.workflow_id].append(payment)
print(f" \\u2756 支付 {amount} umsg -> {task.agent_id}")
def _topological_sort(self, workflow: WorkflowDef) -> list[str]:
in_degree: dict[str, int] = {}
for tid in workflow.tasks: in_degree[tid] = 0
for tid, task in workflow.tasks.items():
for dep in task.dependencies: in_degree[tid] = in_degree.get(tid, 0) + 1
queue = [tid for tid, deg in in_degree.items() if deg == 0]
sorted_tasks = []
while queue:
node = queue.pop(0); sorted_tasks.append(node)
for tid, task in workflow.tasks.items():
if node in task.dependencies: in_degree[tid] -= 1
if in_degree[tid] == 0: queue.append(tid)
if len(sorted_tasks) != len(workflow.tasks): raise ValueError("DAG 包含环")
return sorted_tasks
def _compute_batches(self, workflow: WorkflowDef, order: list[str]) -> list[list[str]]:
level: dict[str, int] = {}
for tid in order:
max_dep = 0
for dep in workflow.tasks[tid].dependencies: max_dep = max(max_dep, level.get(dep, 0) + 1)
level[tid] = max_dep
batches: list[list[str]] = []
for tid, lvl in level.items():
while len(batches) <= lvl: batches.append([])
batches[lvl].append(tid)
return batches
def get_workflow_summary(self, workflow_id: str) -> dict:
results = self.results.get(workflow_id, {})
payments = self.payments.get(workflow_id, [])
total_cost = sum(int(p.amount) for p in payments)
return {
"workflow_id": workflow_id,
"status": "completed" if all(r.status == TaskStatus.COMPLETED for r in results.values()) else "running" if any(r.status == TaskStatus.RUNNING for r in results.values()) else "failed" if any(r.status == TaskStatus.FAILED for r in results.values()) else "pending",
"total_tasks": len(results),
"completed_tasks": sum(1 for r in results.values() if r.status == TaskStatus.COMPLETED),
"failed_tasks": sum(1 for r in results.values() if r.status == TaskStatus.FAILED),
"total_cost_umsg": str(total_cost),
"payment_count": len(payments),
}
async def main():
engine = WorkflowEngine()
agents = {
"agent-research-01": MockAgent("agent-research-01", delay_range=(0.3, 0.8)),
"agent-planner-01": MockAgent("agent-planner-01", delay_range=(0.2, 0.5)),
"agent-writer-01": MockAgent("agent-writer-01", delay_range=(0.8, 1.5)),
"agent-writer-02": MockAgent("agent-writer-02", delay_range=(0.6, 1.2)),
"agent-writer-03": MockAgent("agent-writer-03", delay_range=(1.0, 2.0)),
"agent-reviewer-01": MockAgent("agent-reviewer-01", delay_range=(0.3, 0.7)),
"agent-seo-01": MockAgent("agent-seo-01", delay_range=(0.2, 0.4)),
"agent-publisher-01": MockAgent("agent-publisher-01", delay_range=(0.1, 0.3)),
}
for agent_id, client in agents.items(): engine.register_agent(agent_id, client)
workflow = WorkflowDef(
workflow_id=f"wf-content-factory-{int(time.time())}",
name="AI 内容工厂",
context={"budget_owner": "msg1owner...", "content_strategy": "technical_blog", "target_audience": "developers", "max_budget": "5000000000"},
tasks={
"research": TaskDef(node_id="research", name="主题研究", agent_id="agent-research-01", action="research_topics", params={"niche": "AI Agent 协作", "trend_period": "7d"}, payment_amount="500000"),
"plan": TaskDef(node_id="plan", name="内容规划", agent_id="agent-planner-01", action="plan_content", params={"style": "technical"}, dependencies=["research"], payment_amount="300000"),
"write_1": TaskDef(node_id="write_1", name="文章撰写 A", agent_id="agent-writer-01", action="write_draft", params={"style": "technical"}, dependencies=["plan"], payment_amount="1000000"),
"write_2": TaskDef(node_id="write_2", name="文章撰写 B", agent_id="agent-writer-02", action="write_draft", params={"style": "practical"}, dependencies=["plan"], payment_amount="1000000"),
"write_3": TaskDef(node_id="write_3", name="文章撰写 C", agent_id="agent-writer-03", action="write_draft", params={"style": "beginner"}, dependencies=["plan"], payment_amount="1000000"),
"review": TaskDef(node_id="review", name="质量审核", agent_id="agent-reviewer-01", action="review_content", params={"quality_threshold": 0.7}, dependencies=["write_1", "write_2", "write_3"], payment_amount="400000"),
"seo": TaskDef(node_id="seo", name="SEO 优化", agent_id="agent-seo-01", action="optimize_seo", params={"target_keywords": ["AI Agent 协作", "区块链"]}, dependencies=["review"], payment_amount="300000"),
"publish": TaskDef(node_id="publish", name="发布准备", agent_id="agent-publisher-01", action="format_publish", params={"output_format": "markdown", "platform": "blog"}, dependencies=["seo"], payment_amount="200000"),
},
)
engine.create_workflow(workflow)
result = await engine.execute_workflow(workflow.workflow_id)
summary = engine.get_workflow_summary(workflow.workflow_id)
print(f"\n{'='*60}")
print("工作流执行摘要:")
print(f" 状态: {summary['status']}")
print(f" 总任务: {summary['total_tasks']}")
print(f" 成功: {summary['completed_tasks']}")
print(f" 失败: {summary['failed_tasks']}")
print(f" 总成本: {summary['total_cost_umsg']} umsg ({int(summary['total_cost_umsg']) / 10**18} MSG)")
print(f" 支付笔数: {summary['payment_count']}")
if __name__ == "__main__":
asyncio.run(main())
10.3 运行输出示例
============================================================
开始执行工作流: AI 内容工厂 (wf-content-factory-1827654321)
============================================================
执行顺序: research -> plan -> write_1 -> write_2 -> write_3 -> review -> seo -> publish
--- 并行执行批次: ['research'] ---
>> [主题研究] 分发给 agent-research-01...
\u2713 [主题研究] 完成 (重试0次)
\u2756 支付 500000 umsg -> agent-research-01
--- 并行执行批次: ['plan'] ---
>> [内容规划] 分发给 agent-planner-01...
\u2713 [内容规划] 完成 (重试0次)
--- 并行执行批次: ['write_1', 'write_2', 'write_3'] ---
>> [文章撰写 A] 分发给 agent-writer-01...
>> [文章撰写 B] 分发给 agent-writer-02...
>> [文章撰写 C] 分发给 agent-writer-03...
\u2713 [文章撰写 B] 完成 (重试0次)
\u2713 [文章撰写 A] 完成 (重试0次)
\u2713 [文章撰写 C] 完成 (重试0次)
\u2756 支付 1000000 umsg -> agent-writer-01
...
============================================================
工作流执行摘要:
状态: completed
总任务: 8
成功: 8
失败: 0
总成本: 4700000 umsg (0.0000047 MSG)
支付笔数: 8
10.4 TypeScript 版本替代实现
对于前端/Node.js 环境,提供 TypeScript 版本的工作流编排器核心:
// content_factory_workflow.ts
interface Task {
nodeId: string;
name: string;
agentId: string;
action: string;
params: Record<string, unknown>;
dependencies: string[];
timeout: number;
paymentAmount: string;
}
interface TaskResult {
nodeId: string;
status: "pending" | "running" | "completed" | "failed";
output?: unknown;
error?: string;
}
type AgentExecutor = (action: string, params: Record<string, unknown>) => Promise<unknown>;
class WorkflowOrchestrator {
private tasks: Map<string, Task> = new Map();
private results: Map<string, TaskResult> = new Map();
private agents: Map<string, AgentExecutor> = new Map();
registerAgent(agentId: string, executor: AgentExecutor): void {
this.agents.set(agentId, executor);
}
addTask(task: Task): void {
this.tasks.set(task.nodeId, task);
this.results.set(task.nodeId, { nodeId: task.nodeId, status: "pending" });
}
async execute(): Promise<Map<string, TaskResult>> {
const order = this.topologicalSort();
const batches = this.computeBatches(order);
for (const batch of batches) {
await Promise.all(batch.map((nodeId) => this.executeTask(nodeId)));
}
return this.results;
}
private async executeTask(nodeId: string): Promise<TaskResult> {
const task = this.tasks.get(nodeId)!;
const result = this.results.get(nodeId)!;
// 收集依赖结果
const depsOutput: Record<string, unknown> = {};
for (const depId of task.dependencies) {
const depResult = this.results.get(depId);
if (depResult?.output) depsOutput[depId] = depResult.output;
}
try {
result.status = "running";
const executor = this.agents.get(task.agentId);
if (!executor) throw new Error(`Agent ${task.agentId} 未注册`);
const output = await Promise.race([
executor(task.action, { ...task.params, dependencyResults: depsOutput }),
new Promise((_, reject) => setTimeout(() => reject(new Error("超时")), task.timeout * 1000)),
]);
result.status = "completed";
result.output = output;
} catch (err: any) {
result.status = "failed";
result.error = err.message;
}
return result;
}
private topologicalSort(): string[] {
const inDegree: Map<string, number> = new Map();
for (const [id] of this.tasks) inDegree.set(id, 0);
for (const [, task] of this.tasks) {
for (const dep of task.dependencies) inDegree.set(dep, (inDegree.get(dep) ?? 0) + 1);
}
const queue: string[] = [];
for (const [id, deg] of inDegree) { if (deg === 0) queue.push(id); }
const sorted: string[] = [];
while (queue.length > 0) {
const node = queue.shift()!;
sorted.push(node);
for (const [id, task] of this.tasks) {
if (task.dependencies.includes(node)) {
const newDeg = (inDegree.get(id) ?? 1) - 1;
inDegree.set(id, newDeg);
if (newDeg === 0) queue.push(id);
}
}
}
if (sorted.length !== this.tasks.size) throw new Error("DAG 包含环");
return sorted;
}
private computeBatches(order: string[]): string[][] {
const level: Map<string, number> = new Map();
for (const nodeId of order) {
const task = this.tasks.get(nodeId)!;
let maxDep = 0;
for (const dep of task.dependencies) maxDep = Math.max(maxDep, (level.get(dep) ?? 0) + 1);
level.set(nodeId, maxDep);
}
const batches: string[][] = [];
for (const [nodeId, lvl] of level) {
while (batches.length <= lvl) batches.push([]);
batches[lvl].push(nodeId);
}
return batches;
}
}
11. 最佳实践与故障排除
11.1 工作流设计最佳实践
DAG 设计原则:
- 单一职责:每个任务节点只做一件事。如果一个 Agent 的任务逻辑复杂,拆分为多个子任务
- 最小依赖:减少任务间的依赖关系,最大化并行度
- 合理粒度:任务粒度过细增加编排开销;粒度过粗降低并行效率
- 幂等设计:每个任务应该设计为可重入的,相同的输入应产生相同的输出
- 超时设置:根据任务的预期执行时间设置合理的超时值(通常为预期的 2-3 倍)
Agent 选择策略:
# Agent 选择策略
class AgentSelectionStrategy:
@staticmethod
def cheapest(agents: list) -> str:
"""选择价格最低的 Agent"""
return min(agents, key=lambda a: int(a["price_model"]["amount"]))["agent_id"]
@staticmethod
def fastest(agents: list) -> str:
"""选择历史最快完成的 Agent"""
return max(agents, key=lambda a: a["reputation"]["avg_rating"])["agent_id"]
@staticmethod
def most_reliable(agents: list) -> str:
"""选择最可靠的 Agent"""
return max(agents, key=lambda a: a["reputation"]["completed_tasks"] / max(a["reputation"]["total_tasks"], 1))["agent_id"]
@staticmethod
def weighted_random(agents: list) -> str:
"""加权随机选择"""
import random
total_weight = sum(a.get("weight", 1) for a in agents)
r = random.uniform(0, total_weight)
cumulative = 0
for agent in agents:
cumulative += agent.get("weight", 1)
if r <= cumulative: return agent["agent_id"]
return agents[-1]["agent_id"]
常见陷阱与解决方案:
| 陷阱 | 问题 | 解决方案 |
|---|---|---|
| 循环依赖 | DAG 包含环导致无法执行 | 使用拓扑排序检测;确保依赖方向一致 |
| 过大的上下文 | 共享上下文过大导致链上存储成本高 | 仅在链上存储哈希,完整数据存在 IPFS/Arweave |
| 单点瓶颈 | 一个 Agent 处理所有任务 | 使用并行执行模式;负载均衡 |
| 无限重试 | 任务一直失败无限重试消耗资源 | 设置 max_retries 上限(建议 3-5);启用退避策略 |
| 支付争议 | 任务完成但无法证明 | 使用哈希锁条件支付;链上提交结果哈希 |
| 宪章冲突 | 两个 Agent 的宪章不兼容 | 在 session_propose 阶段自动检查兼容性 |
11.2 性能优化
优化 Gas 消耗:
# gas_optimization.py
class GasOptimizer:
@staticmethod
def estimate_gas(task_count: int, parallel_ratio: float) -> dict:
"""估算工作流的 Gas 消耗"""
base_gas = 100000 # 创建工作流
per_task_gas = 50000 # 每个任务
per_session_gas = 80000 # 每个 A2A 会话
per_payment_gas = 60000 # 每次支付
session_count = int(task_count * (1 - parallel_ratio * 0.5))
payment_count = task_count
total_gas = (base_gas + task_count * per_task_gas + session_count * per_session_gas + payment_count * per_payment_gas)
return {
"total_gas_estimate": total_gas,
"estimated_cost_umsg": total_gas * 25, # 假设 Gas 价格 25 umsg
"breakdown": {
"base": base_gas,
"tasks": task_count * per_task_gas,
"sessions": session_count * per_session_gas,
"payments": payment_count * per_payment_gas,
},
}
@staticmethod
def batch_small_payments(payments: list[dict], threshold: int = 1000000) -> list[dict]:
"""批量聚合小额支付"""
aggregated: dict[str, int] = {}
for p in payments:
key = p["to"]
aggregated[key] = aggregated.get(key, 0) + int(p["amount"])
result = []
for to_agent, total in aggregated.items():
if total < threshold:
result.append({"to": to_agent, "amount": str(total), "batched": True})
else:
result.append({"to": to_agent, "amount": str(total), "batched": False})
return result
缓存策略:
// caching.ts
class WorkflowCache {
private cache: Map<string, { data: unknown; expiresAt: number }> = new Map();
constructor(private ttlMs: number = 60000) {}
get(key: string): unknown | null {
const entry = this.cache.get(key);
if (!entry) return null;
if (Date.now() > entry.expiresAt) { this.cache.delete(key); return null; }
return entry.data;
}
set(key: string, data: unknown): void {
this.cache.set(key, { data, expiresAt: Date.now() + this.ttlMs });
}
// 缓存 Agent 查询结果
async cachedAgentQuery(queryFn: () => Promise<unknown>, agentId: string): Promise<unknown> {
const cacheKey = `agent_${agentId}`;
const cached = this.get(cacheKey);
if (cached) return cached;
const result = await queryFn();
this.set(cacheKey, result);
return result;
}
}
11.3 故障排除指南
常见错误码及处理:
| 错误 | 可能原因 | 处理方式 |
|---|---|---|
ERR_SESSION_TIMEOUT |
响应方 Agent 未在 TTL 内响应 | 检查 Agent 端点可用性;增加 TTL |
ERR_CONSTITUTION_MISMATCH |
双方宪章规则冲突 | 检查宪章兼容性报告;调整规则 |
ERR_PAYMENT_INSUFFICIENT |
发起方余额不足 | 检查账户余额;减少任务金额 |
ERR_AGENT_NOT_FOUND |
Agent 未在 Registry 注册 | 确认 Agent 已注册且状态为 Active |
ERR_DAG_CYCLE |
工作流 DAG 包含环 | 运行拓扑排序检测;消除循环依赖 |
ERR_TASK_TIMEOUT |
单个任务超时 | 增加超时配置;检查 Agent 负载 |
ERR_INVALID_SIGNATURE |
消息签名验证失败 | 检查 Dilithium-5 密钥配置;同步时钟 |
ERR_WORKFLOW_EXPIRED |
工作流超过 expires_at | 延长有效期;加快执行速度 |
调试工具:
# debug_tools.py
import json
class WorkflowDebugger:
@staticmethod
def dump_workflow_state(workflow_id: str, engine) -> str:
"""导出工作流状态用于调试"""
state = {
"workflow_id": workflow_id,
"tasks": {},
"results": {},
"payments": [],
}
wf = engine.workflows.get(workflow_id)
if wf:
state["name"] = wf.name
state["context"] = {k: v for k, v in wf.context.items() if k != "secrets_refs"}
for tid, task in wf.tasks.items():
state["tasks"][tid] = {
"name": task.name,
"agent_id": task.agent_id,
"action": task.action,
"dependencies": task.dependencies,
"timeout": task.timeout,
"payment": task.payment_amount,
}
results = engine.results.get(workflow_id, {})
for tid, result in results.items():
state["results"][tid] = {
"status": result.status.value,
"retries": result.retries,
"error": result.error,
"started_at": result.started_at,
"completed_at": result.completed_at,
"proof_hash": result.proof_hash,
}
payments = engine.payments.get(workflow_id, [])
for p in payments:
state["payments"].append({
"payment_id": p.payment_id, "from": p.from_agent,
"to": p.to_agent, "amount": p.amount, "phase": p.phase.value,
})
return json.dumps(state, indent=2, ensure_ascii=False)
@staticmethod
def simulate_workflow(workflow_def: dict, scenario: str = "happy_path") -> dict:
"""模拟工作流执行以验证正确性"""
simulation = {
"happy_path": {
"research": {"status": "completed", "topics": 3},
"plan": {"status": "completed", "sections": 5},
"write_1": {"status": "completed", "words": 3000},
"write_2": {"status": "completed", "words": 2800},
"write_3": {"status": "completed", "words": 3200},
"review": {"status": "completed", "approved": True},
"seo": {"status": "completed", "seo_score": 91},
"publish": {"status": "completed", "url": "https://..."},
},
"task_failure": {
"write_2": {"status": "failed", "error": "模型超时"},
},
"payment_dispute": {
"review": {"status": "completed", "approved": False},
},
}
return simulation.get(scenario, {})
11.4 监控与告警配置
推荐监控指标:
// metrics.ts
interface WorkflowMetrics {
// 性能指标
avgCompletionTime: number; // 平均完成时间 (s)
p95CompletionTime: number; // P95 完成时间 (s)
taskThroughput: number; // 任务吞吐量 (tasks/min)
// 可靠性指标
successRate: number; // 成功率 (%)
retryRate: number; // 重试率 (%)
timeoutRate: number; // 超时率 (%)
failureRate: number; // 失败率 (%)
// 经济指标
totalGasCost: string; // 总 Gas 费用
avgTaskCost: string; // 平均任务成本
paymentSettlementRate: number; // 支付结算率 (%)
// 信誉指标
agentAvgRating: number; // Agent 平均评分
disputeRate: number; // 争议率 (%)
}
function calculateMetrics(events: WorkflowEvent[]): WorkflowMetrics {
const workflowEvents = events.filter(e => e.eventType === "workflow_state_change");
const taskEvents = events.filter(e => e.eventType === "task_state_change");
const completed = workflowEvents.filter(e => e.data.status === "completed");
const failed = workflowEvents.filter(e => e.data.status === "failed");
return {
avgCompletionTime: 0,
p95CompletionTime: 0,
taskThroughput: 0,
successRate: workflowEvents.length > 0 ? (completed.length / workflowEvents.length) * 100 : 0,
retryRate: 0,
timeoutRate: 0,
failureRate: workflowEvents.length > 0 ? (failed.length / workflowEvents.length) * 100 : 0,
totalGasCost: "0",
avgTaskCost: "0",
paymentSettlementRate: 100,
agentAvgRating: 0,
disputeRate: 0,
};
}
11.5 升级与迁移指南
工作流版本兼容性:
# workflow_migration.py
class WorkflowMigration:
@staticmethod
def check_compatibility(old_version: str, new_version: str) -> bool:
"""检查工作流版本兼容性"""
old_parts = [int(x) for x in old_version.split(".")]
new_parts = [int(x) for x in new_version.split(".")]
# Major 版本不同视为不兼容
return old_parts[0] == new_parts[0]
@staticmethod
def migrate_workflow_state(old_state: dict, target_version: str) -> dict:
"""迁移工作流状态到新版本"""
version = old_state.get("version", "1.0.0")
state = dict(old_state)
if version == "1.0.0" and target_version == "1.1.0":
# 示例迁移:添加新字段
if "current_phase" not in state:
state["current_phase"] = "execution"
state["version"] = "1.1.0"
if version == "1.1.0" and target_version == "2.0.0":
# 重大变更迁移
state["version"] = "2.0.0"
return state
11.6 资源与参考
相关文档:
合约源码位置:
msgchain-mainnet/contracts/cosmwasm/all/agent_a2a_v1/msgchain-mainnet/contracts/cosmwasm/all/agent_registry_v1/msgchain-mainnet/contracts/cosmwasm/all/agent_payment_v1/msgchain-mainnet/contracts/cosmwasm/all/ai_agent_constitution_v1/
Go 实现:
msgchain-mainnet/msgchain-nat-通用维护记录-pqopen/pkg/quantum/agent_a2a_api.gomsgchain-mainnet/msgchain-nat-通用维护记录-pqopen/pkg/quantum/agent_a2a_executor.gomsgchain-mainnet/msgchain-nat-通用维护记录-pqopen/pkg/quantum/agent_payment_api.gomsgchain-mainnet/msgchain-nat-通用维护记录-pqopen/pkg/quantum/agent_registry.go
WebAssembly 二进制:
contracts/cosmwasm/compiled_wasm_v1/agent_a2a_v1.wasm.gzcontracts/cosmwasm/compiled_wasm_v1/agent_registry_v1.wasm.gzcontracts/cosmwasm/compiled_wasm_v1/agent_payment_v1.wasm.gzcontracts/cosmwasm/compiled_wasm_v1/ai_agent_constitution_v1.wasm.gz
API 端点(测试网):
- REST:
https://rest.msgchain.org - WebSocket 事件流:
wss://events.msgchain.org/ws - RPC:
https://rpc.msgchain.org
文档结束
本指南提供了在 MSG Chain 上编排多 Agent 协作工作流的全面参考。随着协议和合约的持续迭代,请定期查阅相关文档以获取最新信息。对于协议级别的功能标记为"部分实现"的部分,请关注 MSG Chain 的官方发布说明以了解完整支持的时间表。
