dApp Docs/AI Agent 间协作工作流编排指南
Development reference. Not independently verified for production.

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/


目录

  1. 引言
  2. 工作流基础概念
  3. MSG Chain 上的 Agent 注册与发现
  4. A2A 通信协议深度解析
  5. 工作流编排模式
  6. AI Agent 宪章与策略引擎
  7. Agent 支付与结算
  8. 工作流状态管理与监控
  9. 安全与信任机制
  10. 实战:完整的协作工作流示例
  11. 最佳实践与故障排除

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 已部署

链参数:

1.3 本指南的目标读者

1.4 核心概念速览


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 工作流生命周期

一个典型的工作流经历以下阶段:

  1. 设计阶段:定义 DAG、选择 Agent、设置参数
  2. 注册阶段:在链上创建工作流实例
  3. 执行阶段:按依赖顺序执行任务
  4. 监控阶段:跟踪进度、处理异常
  5. 完成阶段:收集结果、结算支付
  6. 审计阶段:验证执行记录、更新信誉

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)| 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)管理的链上资产。

宪章的核心作用:

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 内容工厂 工作流,流程如下:

  1. 主题研究:Agent Research 分析热门话题
  2. 内容规划:Agent Planner 制定内容大纲
  3. 草稿撰写:Agent Writer 并行撰写多篇文章
  4. 质量审核:Agent Reviewer 审核内容质量
  5. SEO 优化:Agent SEO 优化关键词和元数据
  6. 发布准备: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 设计原则:

  1. 单一职责:每个任务节点只做一件事。如果一个 Agent 的任务逻辑复杂,拆分为多个子任务
  2. 最小依赖:减少任务间的依赖关系,最大化并行度
  3. 合理粒度:任务粒度过细增加编排开销;粒度过粗降低并行效率
  4. 幂等设计:每个任务应该设计为可重入的,相同的输入应产生相同的输出
  5. 超时设置:根据任务的预期执行时间设置合理的超时值(通常为预期的 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 资源与参考

相关文档:

合约源码位置:

Go 实现:

WebAssembly 二进制:

API 端点(测试网):


文档结束

本指南提供了在 MSG Chain 上编排多 Agent 协作工作流的全面参考。随着协议和合约的持续迭代,请定期查阅相关文档以获取最新信息。对于协议级别的功能标记为"部分实现"的部分,请关注 MSG Chain 的官方发布说明以了解完整支持的时间表。