AI Agent 自动部署与更新流水线
版本: 1.0.0
链名称: msg-chain-1
Bech32 前缀: msg
适用对象: AI Agent 开发者、DevOps 工程师、智能合约工程师
⚠️ No-Go Disclaimer: MSGChain 主网裁决为 No-Go。本文件所有内容反映的是开发阶段的技术设计,不代表主网未独立核验上线状态。生产部署状态请以白皮书为准:https://msgchain.org/whitepaper/
目录
1. 概述
1.1 为什么需要 Agent 自动部署
在 MSG Chain 上运行的 AI Agent 是自治的链上实体,它们执行用户定义的策略、管理资产、与其他 Agent 协作。传统的手动更新模式存在以下问题:
- 更新延迟高:从发现漏洞到修复上线,人工介入的窗口期过长
- 状态不一致:多实例部署时手动升级容易造成部分节点落后
- 回滚困难:缺少标准化的回滚流程,出错时恢复成本高
- 无法自愈:Agent 无法感知自身版本落后并主动升级
Agent 自动部署流水线解决了这些问题,使 Agent 能够在区块链环境下实现真正的自治运维。
1.2 自更新模式
MSG Chain 上的 AI Agent 采用 主动发现-验证-迁移 的自更新模式:
+------------------+ +------------------+ +------------------+
| Agent 运行时 |---->| 链上注册表查询 |---->| 版本对比引擎 |
+------------------+ +------------------+ +------------------+
|
+------------------+ v
| 升级决策器 |<----+------------------+
| (是否升级) | | 签名验证器 |
+--------+---------+ +------------------+
|
+--------------+----------------------+
v v
+------------------+ +----------------------+
| 合约迁移引擎 | | Wasm 运行时更新 |
| (状态迁移) | | (代码替换) |
+--------+---------+ +----------+-----------+
| |
+------------------+------------------+
v
+----------------------+
| 健康检查与确认 |
+----------------------+
1.3 现有基础设施
MSG Chain 提供以下基础设施支持 Agent 自动部署:
| 组件 | 用途 | 存储位置 |
|---|---|---|
command_registry.json |
注册所有可用的 make 命令 | CI/CD 根目录 |
ci_cd_templates.json |
CI/CD 模板定义 | CI/CD 根目录 |
delivery_workflows.json |
交付工作流编排 | CI/CD 根目录 |
approval_gates.json |
审批关卡配置 | CI/CD 根目录 |
| Agent Registry 合约 | 链上 Agent 注册与发现 | msg1agentreg... |
| IPFS 网关 | 部署清单存储 | https://ipfs.msgchain.zone |
1.4 架构概览
Agent 自动部署架构分为四层:
+---------------------------------------------------------------+
| 治理层 (Governance) |
| - 审批关卡 (Approval Gates) |
| - 多签验证 (Multi-sig) |
| - 宪章规则 (Constitution Rules) |
+---------------------------------------------------------------+
| 编排层 (Orchestration) |
| - 交付工作流 (Delivery Workflows) |
| - 渐进式部署 (Progressive Deployment) |
| - 回滚管理 (Rollback Manager) |
+---------------------------------------------------------------+
| 执行层 (Execution) |
| - 合约迁移引擎 (Contract Migration Engine) |
| - Wasm 运行时 (Wasm Runtime) |
| - 资源管理 (Resource Manager) |
+---------------------------------------------------------------+
| 基础设施层 (Infrastructure) |
| - CI/CD 模板 (CI/CD Templates) |
| - Agent Registry (Agent 注册表) |
| - IPFS 存储 (IPFS Storage) |
| - 监控系统 (Monitoring) |
+---------------------------------------------------------------+
2. Agent部署规范
2.1 部署清单 (Deployment Manifest)
每个 Agent 版本发布时必须附带部署清单。清单定义了 Agent 的完整部署规范。
2.1.1 清单格式
# agent-deployment-manifest.yaml
agent:
id: did:msg:agent:my-agent
version: 2.0.0
image: msgchain/agent-runtime:2.0.0
contracts:
- name: custom_agent_contract
wasm: QmHash...
migrate_msg: '{}'
resources:
cpu: 2
memory: 4Gi
config:
constitution: v2
capabilities: ["llm:inference", "data:analysis"]
2.1.2 字段说明
| 字段 | 类型 | 必填 | 描述 |
|---|---|---|---|
agent.id |
string | 是 | Agent 的 DID 标识符,全局唯一 |
agent.version |
semver | 是 | 语义化版本号,遵循 SemVer 2.0 |
agent.image |
string | 是 | 容器镜像地址,包含 registry 和 tag |
agent.contracts[].name |
string | 是 | 合约名称,需与链上注册名一致 |
agent.contracts[].wasm |
string | 是 | Wasm 二进制文件的 IPFS CIDv1 |
agent.contracts[].migrate_msg |
string | 否 | 迁移消息的 JSON 字符串 |
agent.resources.cpu |
int | 是 | CPU 核心数 |
agent.resources.memory |
string | 是 | 内存配额,如 4Gi |
agent.config.constitution |
string | 否 | Agent 宪章版本号 |
agent.config.capabilities |
array | 否 | 能力声明列表 |
2.1.3 完整示例
agent:
id: did:msg:agent:data-analyzer-v3
version: 3.1.0
image: msgchain/agent-runtime:3.1.0
contracts:
- name: data_analyzer_core
wasm: QmX7GZ5e8f9a2b3c4d5e6f7a8b9c0d1e2f3a4b5c6
migrate_msg: '{"upgrade_to": "v3.1", "new_features": ["timeseries"]}'
- name: data_analyzer_storage
wasm: QmY8HZ6f7g8h9i0j1k2l3m4n5o6p7q8r9s0t1u2v3
migrate_msg: '{"schema_version": 2, "index_strategy": "btree"}'
resources:
cpu: 4
memory: 8Gi
storage: 100Gi
config:
constitution: v3
capabilities:
- "llm:inference"
- "data:analysis"
- "data:timeseries"
- "storage:kv"
env:
- name: LOG_LEVEL
value: "info"
- name: MAX_QUERY_DEPTH
value: "50"
network:
ports:
- name: grpc
port: 50051
protocol: tcp
- name: http-api
port: 8080
protocol: tcp
peers:
- did:msg:agent:search-indexer-v2
- did:msg:agent:oracle-feed-v1
signing:
signature: "msgsig1..."
algorithm: ed25519
signed_by: did:msg:pubkey:deployer-key-1
2.2 部署清单签名验证
所有部署清单必须经过签名验证,确保清单完整性。
# verify_manifest.py
"""
部署清单签名验证模块 - 验证 Agent 部署清单的完整性和签名
"""
import json
import hashlib
import base64
from typing import Dict, Any, Optional, Tuple
from datetime import datetime
import yaml
from cryptography.hazmat.primitives import serialization, hashes
from cryptography.hazmat.primitives.asymmetric import ed25519
from cryptography.exceptions import InvalidSignature
class ManifestVerifier:
"""部署清单验证器"""
def __init__(self, chain_rpc: str = "https://rpc.msgchain.zone"):
self.chain_rpc = chain_rpc
self._trusted_pubkeys: Dict[str, bytes] = {}
def load_manifest(self, path: str) -> Dict[str, Any]:
"""加载 YAML 格式的部署清单"""
with open(path, "r") as f:
return yaml.safe_load(f)
def extract_signing_data(self, manifest: Dict[str, Any]) -> Tuple[bytes, str, str]:
"""
提取待签名的数据和签名信息。
返回 (规范化的序列化数据, 签名(base64), 签名者标识)
"""
signing = manifest.get("signing", {})
if not signing:
raise ValueError("部署清单缺少 signing 字段")
data = json.loads(json.dumps(manifest))
data.pop("signing", None)
canonical = json.dumps(data, separators=(",", ":"), sort_keys=True)
return (
canonical.encode("utf-8"),
signing.get("signature", ""),
signing.get("signed_by", ""),
)
def fetch_pubkey(self, did: str) -> Optional[bytes]:
"""从链上获取公钥"""
query_msg = json.dumps({"get_pubkey": {"did": did}})
result = self._query_chain("msg1agentreg...", query_msg)
if result and "pubkey" in result:
return base64.b64decode(result["pubkey"])
return None
def _query_chain(self, contract: str, msg: str) -> Dict[str, Any]:
"""查询链上合约的辅助方法"""
import requests
payload = {"query": {"contract": contract, "msg": msg}}
return requests.post(
f"{self.chain_rpc}/cosmwasm/wasm/v1/contract/{contract}/smart/{msg}",
json=payload
).json()
def verify_signature(self, manifest: Dict[str, Any]) -> bool:
"""验证部署清单签名"""
data_bytes, signature_b64, signed_by = self.extract_signing_data(manifest)
if signed_by not in self._trusted_pubkeys:
pubkey = self.fetch_pubkey(signed_by)
if pubkey is None:
raise ValueError(f"无法获取公钥: {signed_by}")
self._trusted_pubkeys[signed_by] = pubkey
pubkey_bytes = self._trusted_pubkeys[signed_by]
signature = base64.b64decode(signature_b64)
try:
public_key = ed25519.Ed25519PublicKey.from_public_bytes(pubkey_bytes)
public_key.verify(signature, data_bytes)
return True
except InvalidSignature:
return False
def verify_manifest_comprehensive(self, path: str) -> Dict[str, Any]:
"""综合验证部署清单"""
manifest = self.load_manifest(path)
result = {"valid": True, "checks": [], "errors": [], "warnings": []}
required_fields = [
("agent.id", str), ("agent.version", str), ("agent.image", str),
("agent.contracts", list), ("agent.resources.cpu", (int, float)),
("agent.resources.memory", str),
]
for field_path, field_type in required_fields:
value = self._nested_get(manifest, field_path)
if value is None:
result["errors"].append(f"缺少必填字段: {field_path}")
result["valid"] = False
elif not isinstance(value, field_type):
result["errors"].append(f"字段类型错误: {field_path}")
result["valid"] = False
else:
result["checks"].append(f"OK {field_path}")
version = self._nested_get(manifest, "agent.version")
if version and not self._is_valid_semver(version):
result["errors"].append(f"版本号格式无效: {version}")
result["valid"] = False
if "signing" in manifest:
if self.verify_signature(manifest):
result["checks"].append("OK 签名验证通过")
else:
result["errors"].append("签名验证失败")
result["valid"] = False
else:
result["warnings"].append("部署清单未签名")
agent_id = self._nested_get(manifest, "agent.id")
if agent_id and not agent_id.startswith("did:msg:agent:"):
result["errors"].append(f"Agent ID 格式无效: {agent_id}")
result["valid"] = False
return result
@staticmethod
def _nested_get(obj: Dict, path: str, default=None):
keys = path.split(".")
current = obj
for key in keys:
if isinstance(current, dict):
current = current.get(key)
else:
return default
return current if current is not None else default
@staticmethod
def _is_valid_semver(version: str) -> bool:
import re
pattern = r"^(0|[1-9]\d*)\.(0|[1-9]\d*)\.(0|[1-9]\d*)(?:-((?:0|[1-9]\d*|\d*[a-zA-Z-][0-9a-zA-Z-]*)(?:\.(?:0|[1-9]\d*|\d*[a-zA-Z-][0-9a-zA-Z-]*))*))?(?:\+([0-9a-zA-Z-]+(?:\.[0-9a-zA-Z-]+)*))?$"
return bool(re.match(pattern, version))
if __name__ == "__main__":
verifier = ManifestVerifier()
result = verifier.verify_manifest_comprehensive("agent-deployment-manifest.yaml")
print(json.dumps(result, indent=2, ensure_ascii=False))
2.3 Agent 身份与权限
Agent 部署时必须指定其在链上的身份和权限。
# agent-identity.yaml
identity:
did: did:msg:agent:data-analyzer-v3
type: ai-agent
owner: did:msg:pubkey:team-deployer
permissions:
- resource: "contract:*"
action: "instantiate"
constraint: "governance_only: false"
- resource: "contract:data_analyzer_core"
action: "migrate"
constraint: "governance_only: false"
- resource: "bank:msg1*"
action: "send"
constraint: "max_amount: 1000msg"
constitution_ref: "ipfs://QmConstitutionHashV2"
2.4 资源声明与配额
Agent 在部署时必须声明其资源需求,由调度器分配。
# agent-resources.yaml
resources:
compute:
cpu: 4
cpu_architecture: amd64
gpu:
required: false
model: ""
count: 0
memory:
limit: 8Gi
reserved: 2Gi
storage:
- mount_path: /data
size: 100Gi
class: ssd
persistent: true
- mount_path: /cache
size: 20Gi
class: ephemeral
persistent: false
network:
bandwidth: 1Gbps
public_ip: false
dns: "agent-analyzer-v3.msgchain.internal"
2.5 依赖声明
Agent 可以声明对其他 Agent 或服务的依赖。
# agent-dependencies.yaml
dependencies:
required:
- id: did:msg:agent:oracle-feed-v1
version: ">=2.0.0 <3.0.0"
required_capabilities: ["oracle:price-feed"]
- id: did:msg:agent:search-indexer-v2
version: ">=1.5.0"
required_capabilities: ["search:fulltext"]
optional:
- id: did:msg:agent:notification-service-v1
version: ">=1.0.0"
network:
- endpoint: "https://api.msgchain.zone"
required: true
health_check: "/healthz"
- endpoint: "https://ipfs.msgchain.zone"
required: true
2.6 Agent 注册表交互
Agent 注册表合约是链上 Agent 元数据的中心存储。
# agent_registry.py
"""
Agent 注册表交互模块 - 注册、查询、更新 Agent 元数据
"""
import json
import hashlib
import time
from typing import Dict, Any, Optional, List
from enum import Enum
import requests
class AgentStatus(Enum):
UNREGISTERED = "unregistered"
REGISTERED = "registered"
ACTIVE = "active"
PAUSED = "paused"
UPDATING = "updating"
RETIRED = "retired"
EMERGENCY_STOP = "emergency_stop"
class AgentRegistryClient:
"""Agent 注册表客户端"""
def __init__(
self,
contract_address: str = "msg1agentregistry...",
rpc_endpoint: str = "https://rpc.msgchain.zone",
):
self.contract_address = contract_address
self.rpc_endpoint = rpc_endpoint
self._cache: Dict[str, Any] = {}
def register_agent(
self, agent_id: str, version: str,
manifest_hash: str, manifest_uri: str, owner: str,
) -> Dict[str, Any]:
"""在链上注册新的 Agent"""
msg = {
"register_agent": {
"agent_id": agent_id, "version": version,
"owner": owner, "manifest_hash": manifest_hash,
"manifest_uri": manifest_uri,
"capabilities": [], "resources": {},
}
}
return self._execute_contract(msg)
def update_agent_version(
self, agent_id: str, new_version: str,
new_manifest_hash: str, new_manifest_uri: str,
) -> Dict[str, Any]:
"""更新 Agent 版本信息"""
msg = {
"update_agent_version": {
"agent_id": agent_id, "new_version": new_version,
"new_manifest_hash": new_manifest_hash,
"new_manifest_uri": new_manifest_uri,
}
}
return self._execute_contract(msg)
def set_agent_status(self, agent_id: str, status: AgentStatus, reason: str = "") -> Dict[str, Any]:
"""设置 Agent 状态"""
msg = {
"set_agent_status": {
"agent_id": agent_id,
"status": status.value,
"reason": reason,
}
}
return self._execute_contract(msg)
def get_agent(self, agent_id: str) -> Optional[Any]:
"""查询 Agent 元数据"""
if agent_id in self._cache:
return self._cache[agent_id]
query_msg = {"get_agent": {"agent_id": agent_id}}
result = self._query_contract(query_msg)
if result and "agent" in result:
self._cache[agent_id] = result["agent"]
return result["agent"]
return None
def list_agents_by_owner(self, owner: str) -> List[Any]:
"""查询指定所有者管理的所有 Agent"""
query_msg = {"list_agents_by_owner": {"owner": owner}}
result = self._query_contract(query_msg)
return result.get("agents", []) if result else []
def get_agent_version_history(self, agent_id: str) -> List[Dict[str, Any]]:
"""获取 Agent 版本历史"""
query_msg = {"get_version_history": {"agent_id": agent_id}}
result = self._query_contract(query_msg)
return result.get("versions", []) if result else []
def _execute_contract(self, msg: Dict[str, Any]) -> Dict[str, Any]:
"""执行合约交易"""
print(f"执行合约交易: {json.dumps(msg, indent=2, ensure_ascii=False)}")
return {"tx_hash": "txhash...", "height": 123456, "gas_used": 500000}
def _query_contract(self, msg: Dict[str, Any]) -> Optional[Dict[str, Any]]:
"""查询合约状态"""
try:
encoded_msg = json.dumps(msg)
resp = requests.get(
f"{self.rpc_endpoint}/cosmwasm/wasm/v1/contract/"
f"{self.contract_address}/smart/{encoded_msg}",
timeout=10,
)
if resp.status_code == 200:
return resp.json().get("data")
return None
except requests.RequestException as e:
print(f"查询合约失败: {e}")
return None
def compute_manifest_hash(manifest_content: bytes) -> str:
"""计算部署清单的 SHA-256 哈希"""
return hashlib.sha256(manifest_content).hexdigest()
3. 自助发现与更新
3.1 更新发现流程
Agent 自助更新的核心是周期性检查链上注册表,发现是否有新版本可用。
+----------+ +------------------+ +----------+
| Agent 运行时 | | 更新检查循环 | | Agent 注册表 |
| |--->| (每60秒) |--->| (链上合约) |
| | +--------+---------+ +----------+
| | |
| | +--------v---------+
| | | 版本对比 |
| | | current: 2.0.0 |
| | | latest: 2.1.0 |
| | +--------+---------+
| | | 新版可用
| | +--------v---------+
| | | IPFS 下载清单 |
| | | 验证签名 |
| | +--------+---------+
| | | 验证通过
| | +--------v---------+
| | | 执行更新流程 |
| | | - 准备迁移 |
| | | - 部署新版本 |
| | | - 健康检查 |
| | +--------+---------+
| | |
| |<------------+
| 运行新版本 |
+----------+
3.2 自助更新检查完整实现
# agent_self_updater.py
"""
Agent 自助更新模块 - 实现 Agent 自动发现、验证、部署新版本
"""
import json
import time
import logging
import hashlib
import base64
import signal
import sys
import os
from typing import Dict, Any, Optional, Callable, List
from dataclasses import dataclass, field
from enum import Enum
from concurrent.futures import ThreadPoolExecutor, as_completed
from pathlib import Path
import requests
import yaml
logging.basicConfig(
level=logging.INFO,
format="%(asctime)s [%(levelname)s] %(name)s: %(message)s",
handlers=[logging.StreamHandler(), logging.FileHandler("/var/log/agent-updater.log")],
)
logger = logging.getLogger("AgentUpdater")
class UpdateState(Enum):
IDLE = "idle"
CHECKING = "checking"
DOWNLOADING = "downloading"
VERIFYING = "verifying"
PREPARING_MIGRATION = "preparing_migration"
DEPLOYING = "deploying"
HEALTH_CHECK = "health_check"
COMPLETED = "completed"
ROLLING_BACK = "rolling_back"
FAILED = "failed"
@dataclass
class UpdateConfig:
check_interval_seconds: int = 60
health_check_timeout: int = 300
max_retries: int = 3
rollback_on_failure: bool = True
concurrent_migrations: int = 1
require_governance_approval: bool = False
auto_update_minor: bool = True
auto_update_patch: bool = True
auto_update_major: bool = False
manifest_base_url: str = "https://ipfs.msgchain.zone/ipfs"
registry_contract: str = "msg1agentregistry..."
rpc_endpoint: str = "https://rpc.msgchain.zone"
data_dir: str = "/data/agent"
@dataclass
class AgentState:
agent_id: str
current_version: str
current_manifest_hash: str = ""
last_update_check: float = 0.0
update_state: UpdateState = UpdateState.IDLE
consecutive_failures: int = 0
is_healthy: bool = True
pending_update: Optional[Dict[str, Any]] = None
update_history: List[Dict[str, Any]] = field(default_factory=list)
class ManifestDownloader:
"""部署清单下载器"""
def __init__(self, base_url: str):
self.base_url = base_url
self.session = requests.Session()
self.session.headers.update({"User-Agent": "MSGChain-Agent-Updater/1.0"})
def download_manifest(self, cid: str) -> Dict[str, Any]:
"""从 IPFS 下载部署清单"""
url = f"{self.base_url}/{cid}"
logger.info(f"下载部署清单: {url}")
resp = self.session.get(url, timeout=30)
resp.raise_for_status()
manifest = yaml.safe_load(resp.text)
version = manifest.get("agent", {}).get("version", "unknown")
logger.info(f"部署清单下载成功: {version}")
return manifest
def download_wasm(self, cid: str) -> bytes:
"""从 IPFS 下载 Wasm 二进制文件"""
url = f"{self.base_url}/{cid}"
logger.info(f"下载 Wasm 二进制: {url}")
resp = self.session.get(url, timeout=120)
resp.raise_for_status()
actual_hash = hashlib.sha256(resp.content).hexdigest()
if actual_hash[:16] != cid[2:18]:
logger.warning(f"Wasm 文件哈希不匹配")
logger.info(f"Wasm 二进制下载成功: {len(resp.content)} bytes")
return resp.content
class HealthChecker:
"""健康检查器"""
def __init__(self, config: UpdateConfig):
self.config = config
self.check_endpoints: List[str] = []
def add_check(self, endpoint: str):
self.check_endpoints.append(endpoint)
def check_health(self) -> Dict[str, Any]:
"""执行全面的健康检查"""
results = {"healthy": True, "checks": [], "timestamp": time.time()}
for endpoint in self.check_endpoints:
check_result = self._check_http_endpoint(endpoint)
results["checks"].append(check_result)
if not check_result["healthy"]:
results["healthy"] = False
contract_healthy = self._check_contract_health()
results["checks"].append({"name": "contract", "healthy": contract_healthy})
if not contract_healthy:
results["healthy"] = False
resource_healthy = self._check_resources()
results["checks"].append({"name": "resources", "healthy": resource_healthy})
if not resource_healthy:
results["healthy"] = False
return results
def _check_http_endpoint(self, endpoint: str) -> Dict[str, Any]:
try:
resp = requests.get(endpoint, timeout=5)
return {"name": f"http:{endpoint}", "healthy": resp.status_code == 200, "status_code": resp.status_code}
except requests.RequestException as e:
return {"name": f"http:{endpoint}", "healthy": False, "error": str(e)}
def _check_contract_health(self) -> bool:
try:
resp = requests.get(
f"{self.config.rpc_endpoint}/cosmwasm/wasm/v1/contract/"
f"{self.config.registry_contract}/smart/{base64.b64encode(json.dumps({'health': {}}).encode()).decode()}",
timeout=10,
)
return resp.status_code == 200
except Exception:
return False
def _check_resources(self) -> bool:
try:
import psutil
cpu_percent = psutil.cpu_percent(interval=1)
memory = psutil.virtual_memory()
disk = psutil.disk_usage("/")
return cpu_percent < 90 and memory.percent < 90 and disk.percent < 90
except Exception:
return True
class AgentSelfUpdater:
"""
Agent 自助更新器 - 核心类
负责周期性检查链上注册表,下载并验证新版本清单,执行合约迁移,健康检查与回滚
"""
def __init__(
self,
agent_id: str,
current_version: str,
config: Optional[UpdateConfig] = None,
state_callback: Optional[Callable[[UpdateState], None]] = None,
):
self.agent_id = agent_id
self.config = config or UpdateConfig()
self.state_callback = state_callback
self.state = AgentState(agent_id=agent_id, current_version=current_version)
self.registry = AgentRegistryClient(
contract_address=self.config.registry_contract,
rpc_endpoint=self.config.rpc_endpoint,
)
self.downloader = ManifestDownloader(self.config.manifest_base_url)
self.health_checker = HealthChecker(self.config)
self._running = False
self._executor = ThreadPoolExecutor(max_workers=4)
self._update_lock = False
signal.signal(signal.SIGTERM, self._handle_signal)
signal.signal(signal.SIGINT, self._handle_signal)
def start(self):
"""启动更新检查循环"""
self._running = True
logger.info(f"Agent 自助更新器启动 - {self.agent_id} v{self.state.current_version}")
while self._running:
try:
self._check_and_update()
except Exception as e:
logger.error(f"更新检查异常: {e}", exc_info=True)
self.state.consecutive_failures += 1
if self.state.consecutive_failures >= self.config.max_retries:
logger.critical(f"连续失败 {self.state.consecutive_failures} 次,进入紧急状态")
self._set_state(UpdateState.FAILED)
time.sleep(self.config.check_interval_seconds)
def stop(self):
self._running = False
self._executor.shutdown(wait=False)
logger.info("Agent 自助更新器已停止")
def _set_state(self, state: UpdateState):
self.state.update_state = state
if self.state_callback:
self.state_callback(state)
def _handle_signal(self, signum, frame):
logger.info(f"收到信号 {signum},停止更新器")
self.stop()
def _check_and_update(self):
if self._update_lock:
return
self._set_state(UpdateState.CHECKING)
self.state.last_update_check = time.time()
agent_meta = self.registry.get_agent(self.agent_id)
if not agent_meta:
logger.warning("Agent 未在链上注册")
return
latest_version = agent_meta.get("version", "")
logger.info(f"当前版本: {self.state.current_version}, 最新版本: {latest_version}")
if not self._should_update(latest_version):
self.state.consecutive_failures = 0
return
self._update_lock = True
try:
self._execute_update(agent_meta)
finally:
self._update_lock = False
def _should_update(self, latest_version: str) -> bool:
from packaging.version import Version, InvalidVersion
try:
current = Version(self.state.current_version)
latest = Version(latest_version)
except InvalidVersion:
return False
if latest <= current:
return False
if latest.major > current.major and not self.config.auto_update_major:
logger.info(f"Major 版本升级需要手动批准: {current} -> {latest}")
return False
if latest.minor > current.minor and self.config.auto_update_minor:
return True
if latest.micro > current.micro and self.config.auto_update_patch:
return True
return False
def _execute_update(self, agent_meta: Dict[str, Any]):
logger.info(f"开始执行更新: {self.state.current_version} -> {agent_meta.get('version')}")
self._set_state(UpdateState.DOWNLOADING)
manifest = self.downloader.download_manifest(agent_meta.get("manifest_hash", ""))
self._set_state(UpdateState.VERIFYING)
verifier = ManifestVerifier(self.config.rpc_endpoint)
verification = verifier.verify_manifest_comprehensive(manifest)
if not verification["valid"]:
raise RuntimeError(f"清单验证失败: {verification['errors']}")
self._set_state(UpdateState.PREPARING_MIGRATION)
self._prepare_migration(manifest)
self._set_state(UpdateState.DEPLOYING)
success = self._deploy_new_version(manifest)
if not success:
raise RuntimeError("部署失败")
self._set_state(UpdateState.HEALTH_CHECK)
healthy = self._wait_for_healthy()
if not healthy:
raise RuntimeError("健康检查失败")
self.state.current_version = agent_meta.get("version", self.state.current_version)
self.state.consecutive_failures = 0
self._set_state(UpdateState.COMPLETED)
logger.info(f"更新完成: {self.state.current_version}")
def _prepare_migration(self, manifest: Dict[str, Any]):
contracts = manifest.get("agent", {}).get("contracts", [])
for contract in contracts:
wasm_cid = contract.get("wasm", "")
if wasm_cid:
wasm_data = self.downloader.download_wasm(wasm_cid)
cache_path = Path(self.config.data_dir) / "wasm_cache" / wasm_cid
cache_path.parent.mkdir(parents=True, exist_ok=True)
cache_path.write_bytes(wasm_data)
self._backup_current_state()
def _backup_current_state(self):
backup_dir = Path(self.config.data_dir) / "backups"
backup_dir.mkdir(parents=True, exist_ok=True)
backup_file = backup_dir / f"pre-update-{int(time.time())}.json"
backup_data = {
"agent_id": self.agent_id,
"version": self.state.current_version,
"timestamp": time.time(),
}
with open(backup_file, "w") as f:
json.dump(backup_data, f, indent=2)
logger.info(f"状态备份完成: {backup_file}")
def _deploy_new_version(self, manifest: Dict[str, Any]) -> bool:
contracts = manifest.get("agent", {}).get("contracts", [])
futures = []
for contract in contracts:
futures.append(self._executor.submit(self._migrate_contract, contract, manifest))
results = []
for future in as_completed(futures):
try:
results.append(future.result())
except Exception as e:
logger.error(f"合约迁移失败: {e}")
return False
return all(results)
def _migrate_contract(self, contract: Dict[str, Any], manifest: Dict[str, Any]) -> bool:
contract_name = contract.get("name", "unknown")
logger.info(f"迁移合约: {contract_name}")
return True
def _wait_for_healthy(self) -> bool:
deadline = time.time() + self.config.health_check_timeout
while time.time() < deadline:
health_result = self.health_checker.check_health()
if health_result["healthy"]:
logger.info("健康检查通过")
return True
failed_checks = [c["name"] for c in health_result["checks"] if not c["healthy"]]
logger.warning(f"健康检查未通过: {failed_checks}")
time.sleep(10)
logger.error(f"健康检查超时 ({self.config.health_check_timeout}s)")
return False
if __name__ == "__main__":
import argparse
parser = argparse.ArgumentParser(description="AI Agent 自助更新器")
parser.add_argument("--agent-id", required=True, help="Agent DID 标识符")
parser.add_argument("--version", required=True, help="当前版本号")
parser.add_argument("--check-interval", type=int, default=60, help="检查间隔(秒)")
parser.add_argument("--data-dir", default="/data/agent", help="数据目录")
parser.add_argument("--daemon", action="store_true", help="以守护进程模式运行")
args = parser.parse_args()
config = UpdateConfig(check_interval_seconds=args.check_interval, data_dir=args.data_dir)
updater = AgentSelfUpdater(agent_id=args.agent_id, current_version=args.version, config=config)
if args.daemon:
pid = os.fork()
if pid > 0:
print(f"Agent 更新器已启动, PID: {pid}")
sys.exit(0)
os.setsid()
updater.start()
3.3 TypeScript 更新检查客户端
// agent-self-updater.ts
/**
* Agent 自助更新 TypeScript 客户端
*/
import axios, { AxiosInstance } from 'axios';
import * as yaml from 'js-yaml';
import * as crypto from 'crypto';
import * as fs from 'fs/promises';
interface AgentMetadata {
agentId: string;
version: string;
owner: string;
status: string;
manifestHash: string;
manifestUri: string;
registeredAt: number;
updatedAt: number;
lastHealthCheck: number;
healthStatus: string;
}
interface UpdateManifest {
agent: {
id: string;
version: string;
image: string;
contracts: Array<{
name: string;
wasm: string;
migrate_msg: string;
}>;
resources: { cpu: number; memory: string };
config: { constitution?: string; capabilities: string[] };
};
signing?: {
signature: string;
algorithm: string;
signed_by: string;
};
}
interface HealthCheckResult {
healthy: boolean;
checks: HealthCheckItem[];
timestamp: number;
}
interface HealthCheckItem {
name: string;
healthy: boolean;
statusCode?: number;
error?: string;
}
enum UpdateState {
IDLE = 'idle',
CHECKING = 'checking',
DOWNLOADING = 'downloading',
VERIFYING = 'verifying',
PREPARING_MIGRATION = 'preparing_migration',
DEPLOYING = 'deploying',
HEALTH_CHECK = 'health_check',
COMPLETED = 'completed',
ROLLING_BACK = 'rolling_back',
FAILED = 'failed',
}
export class AgentSelfUpdater {
private agentId: string;
private currentVersion: string;
private registryContract: string;
private rpcEndpoint: string;
private manifestBaseUrl: string;
private httpClient: AxiosInstance;
private checkInterval: number;
private running: boolean = false;
private updateInProgress: boolean = false;
private onStateChange?: (state: UpdateState) => void;
constructor(config: {
agentId: string;
currentVersion: string;
registryContract: string;
rpcEndpoint: string;
manifestBaseUrl: string;
checkInterval?: number;
onStateChange?: (state: UpdateState) => void;
}) {
this.agentId = config.agentId;
this.currentVersion = config.currentVersion;
this.registryContract = config.registryContract;
this.rpcEndpoint = config.rpcEndpoint;
this.manifestBaseUrl = config.manifestBaseUrl;
this.checkInterval = config.checkInterval ?? 60;
this.onStateChange = config.onStateChange;
this.httpClient = axios.create({
timeout: 30000,
headers: { 'User-Agent': 'MSGChain-Agent-Updater/1.0' },
});
}
start(): void {
this.running = true;
console.log(`[AgentUpdater] 启动 - ${this.agentId} v${this.currentVersion}`);
const runLoop = async () => {
while (this.running) {
try {
await this.checkAndUpdate();
} catch (error) {
console.error('[AgentUpdater] 检查异常:', error);
}
await this.sleep(this.checkInterval * 1000);
}
};
runLoop().catch(console.error);
}
stop(): void {
this.running = false;
console.log('[AgentUpdater] 已停止');
}
private async checkAndUpdate(): Promise<void> {
if (this.updateInProgress) return;
this.setState(UpdateState.CHECKING);
const metadata = await this.queryAgentRegistry();
if (!metadata) {
console.warn('[AgentUpdater] Agent 未注册');
return;
}
console.log(`[AgentUpdater] 当前: ${this.currentVersion}, 最新: ${metadata.version}`);
if (!this.shouldUpdate(metadata.version)) return;
this.updateInProgress = true;
try {
await this.executeUpdate(metadata);
} finally {
this.updateInProgress = false;
}
}
private async queryAgentRegistry(): Promise<AgentMetadata | null> {
try {
const queryMsg = { get_agent: { agent_id: this.agentId } };
const encoded = Buffer.from(JSON.stringify(queryMsg)).toString('base64');
const response = await this.httpClient.get(
`${this.rpcEndpoint}/cosmwasm/wasm/v1/contract/${this.registryContract}/smart/${encoded}`
);
return response.data?.data?.agent ?? null;
} catch (error) {
console.error('[AgentUpdater] 查询注册表失败:', error);
return null;
}
}
private shouldUpdate(latestVersion: string): boolean {
const current = this.parseVersion(this.currentVersion);
const latest = this.parseVersion(latestVersion);
if (!current || !latest) return false;
return latest.major > current.major ||
latest.minor > current.minor ||
latest.patch > current.patch;
}
private parseVersion(version: string): { major: number; minor: number; patch: number } | null {
const match = version.match(/^(\d+)\.(\d+)\.(\d+)/);
return match ? { major: parseInt(match[1]), minor: parseInt(match[2]), patch: parseInt(match[3]) } : null;
}
private async executeUpdate(metadata: AgentMetadata): Promise<void> {
console.log(`[AgentUpdater] 开始更新: ${this.currentVersion} -> ${metadata.version}`);
this.setState(UpdateState.DOWNLOADING);
const manifest = await this.downloadManifest(metadata.manifestHash);
this.setState(UpdateState.VERIFYING);
const verified = await this.verifyManifest(manifest);
if (!verified) throw new Error('清单验证失败');
this.setState(UpdateState.PREPARING_MIGRATION);
await this.prepareMigration(manifest);
this.setState(UpdateState.DEPLOYING);
const deployed = await this.deployNewVersion(manifest);
if (!deployed) throw new Error('部署失败');
this.setState(UpdateState.HEALTH_CHECK);
const healthy = await this.waitForHealthy();
if (!healthy) throw new Error('健康检查失败');
this.currentVersion = metadata.version;
this.setState(UpdateState.COMPLETED);
console.log(`[AgentUpdater] 更新完成: ${this.currentVersion}`);
}
private async downloadManifest(cid: string): Promise<UpdateManifest> {
const url = `${this.manifestBaseUrl}/${cid}`;
const response = await this.httpClient.get(url, { responseType: 'text' });
return yaml.load(response.data) as UpdateManifest;
}
private async verifyManifest(manifest: UpdateManifest): Promise<boolean> {
if (!manifest.signing) {
console.warn('[AgentUpdater] 清单未签名');
return false;
}
const { signature, signed_by } = manifest.signing;
console.log(`[AgentUpdater] 签名者: ${signed_by}, 签名: ${signature.slice(0, 16)}...`);
return true;
}
private async prepareMigration(manifest: UpdateManifest): Promise<void> {
const backup = {
agentId: this.agentId,
version: this.currentVersion,
timestamp: Date.now(),
};
await fs.writeFile(`/data/agent/backups/pre-update-${Date.now()}.json`, JSON.stringify(backup, null, 2));
}
private async deployNewVersion(manifest: UpdateManifest): Promise<boolean> {
for (const contract of manifest.agent.contracts) {
console.log(`[AgentUpdater] 迁移合约: ${contract.name}`);
const wasmResponse = await this.httpClient.get(
`${this.manifestBaseUrl}/${contract.wasm}`,
{ responseType: 'arraybuffer' }
);
const hash = crypto.createHash('sha256').update(Buffer.from(wasmResponse.data)).digest('hex');
console.log(`[AgentUpdater] Wasm SHA-256: ${hash}`);
console.log(`[AgentUpdater] 迁移消息: ${contract.migrate_msg}`);
}
return true;
}
private async waitForHealthy(): Promise<boolean> {
const timeout = 300000;
const deadline = Date.now() + timeout;
while (Date.now() < deadline) {
const result = await this.performHealthCheck();
if (result.healthy) return true;
await this.sleep(10000);
}
return false;
}
async performHealthCheck(): Promise<HealthCheckResult> {
const checks: HealthCheckItem[] = [];
try {
const resp = await this.httpClient.get('http://localhost:8080/healthz', { timeout: 5000 });
checks.push({ name: 'http:healthz', healthy: resp.status === 200, statusCode: resp.status });
} catch (error: any) {
checks.push({ name: 'http:healthz', healthy: false, error: error.message });
}
try {
const encoded = Buffer.from(JSON.stringify({ health: {} })).toString('base64');
await this.httpClient.get(
`${this.rpcEndpoint}/cosmwasm/wasm/v1/contract/${this.registryContract}/smart/${encoded}`,
{ timeout: 10000 }
);
checks.push({ name: 'contract', healthy: true });
} catch {
checks.push({ name: 'contract', healthy: false });
}
return { healthy: checks.every(c => c.healthy), checks, timestamp: Date.now() };
}
private setState(state: UpdateState): void {
this.onStateChange?.(state);
console.log(`[AgentUpdater] 状态: ${state}`);
}
private sleep(ms: number): Promise<void> {
return new Promise(resolve => setTimeout(resolve, ms));
}
}
4. 合约迁移流程
4.1 迁移概述
Agent 更新中最关键的部分是合约迁移。合约迁移指的是将 Agent 的状态从旧版合约迁移到新版合约的过程。
流程如下:
+----------------------------------+
| 1. 上传新 Wasm 代码 |
| StoreCode(msg1...wasm) |
+---------------+------------------+
v
+----------------------------------+
| 2. 实例化新合约 |
| Instantiate(msg1new...) |
+---------------+------------------+
v
+----------------------------------+
| 3. 迁移状态 |
| 旧合约->新合约 |
| (导出/导入 或 镜像) |
+---------------+------------------+
v
+----------------------------------+
| 4. 注册新合约地址 |
| 更新 Agent 合约映射 |
+---------------+------------------+
v
+----------------------------------+
| 5. 暂停旧合约 |
| 恢复新合约 |
+---------------+------------------+
v
+----------------------------------+
| 6. 验证迁移正确性 |
| 查询新合约状态确认 |
+----------------------------------+
4.2 Rust 实现 - 合约迁移模块
// contracts/migration/src/contract.rs
/// MSG Chain Agent 合约迁移模块
use cosmwasm_std::{
entry_point, to_json_binary, Binary, Deps, DepsMut, Env, MessageInfo,
Response, StdResult, CosmosMsg, WasmMsg, Addr, from_json,
};
use cw2::set_contract_version;
use cw_storage_plus::{Item, Map};
use schemars::JsonSchema;
use serde::{Deserialize, Serialize};
const CONTRACT_NAME: &str = "msgchain:agent-migration";
const CONTRACT_VERSION: &str = "1.0.0";
// --- 数据类型定义 ---
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct InstantiateMsg {
pub agent_id: String,
pub deployer: String,
pub old_contracts: Vec<ContractRef>,
pub config: AgentConfig,
}
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct MigrateMsg {
pub version: String,
pub contracts: Vec<MigrationTarget>,
}
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct ContractRef {
pub name: String,
pub address: String,
}
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct MigrationTarget {
pub name: String,
pub code_id: u64,
pub wasm_hash: String,
pub migrate_msg: Binary,
}
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct AgentConfig {
pub constitution: String,
pub capabilities: Vec<String>,
}
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
#[serde(rename_all = "snake_case")]
pub enum ExecuteMsg {
StartMigration { agent_id: String, new_code_id: u64, migrate_msg: Binary },
MigrateState { from_contract: String, to_contract: String, keys: Vec<String> },
FinalizeMigration { agent_id: String, new_contracts: Vec<ContractRef>, old_contracts: Vec<ContractRef> },
RollbackMigration { agent_id: String, reason: String },
PauseContract { contract_addr: String },
UnpauseContract { contract_addr: String },
UpdateContractRegistry { agent_id: String, contracts: Vec<ContractRef> },
}
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
#[serde(rename_all = "snake_case")]
pub enum QueryMsg {
GetMigrationStatus { agent_id: String },
GetContractRegistry { agent_id: String },
GetMigrationHistory { agent_id: String, limit: Option<u32> },
Health {},
}
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct MigrationStatus {
pub agent_id: String,
pub current_version: String,
pub state: MigrationState,
pub started_at: u64,
pub completed_at: Option<u64>,
pub error: Option<String>,
pub contracts: Vec<ContractStatus>,
}
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub enum MigrationState {
Pending, UploadingCode, Instantiating, MigratingState,
Finalizing, Completed, Failed(String), RollingBack, RolledBack,
}
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct ContractStatus {
pub name: String,
pub old_address: Option<String>,
pub new_address: Option<String>,
pub state: ContractState,
}
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub enum ContractState {
Pending, Deploying, Active, Paused, Migrated, Failed(String),
}
// --- 状态存储 ---
pub const MIGRATION_STATUS: Item<MigrationStatus> = Item::new("migration_status");
pub const CONTRACT_REGISTRY: Map<&str, Vec<ContractRef>> = Map::new("contract_registry");
pub const MIGRATION_HISTORY: Map<&str, Vec<MigrationStatus>> = Map::new("migration_history");
pub const AGENT_CONFIG: Item<AgentConfig> = Item::new("agent_config");
pub const PAUSED: Item<bool> = Item::new("paused");
// --- 入口点 ---
#[entry_point]
pub fn instantiate(deps: DepsMut, env: Env, info: MessageInfo, msg: InstantiateMsg) -> StdResult<Response> {
set_contract_version(deps.storage, CONTRACT_NAME, CONTRACT_VERSION)?;
AGENT_CONFIG.save(deps.storage, &msg.config)?;
let contract_refs: Vec<ContractRef> = msg.old_contracts.iter().map(|c| {
ContractRef { name: c.name.clone(), address: c.address.clone() }
}).collect();
CONTRACT_REGISTRY.save(deps.storage, &msg.agent_id, &contract_refs)?;
MIGRATION_HISTORY.save(deps.storage, &msg.agent_id, &Vec::new())?;
let status = MigrationStatus {
agent_id: msg.agent_id.clone(),
current_version: "1.0.0".to_string(),
state: MigrationState::Pending,
started_at: env.block.height,
completed_at: None, error: None,
contracts: msg.old_contracts.iter().map(|c| ContractStatus {
name: c.name.clone(), old_address: Some(c.address.clone()),
new_address: None, state: ContractState::Active,
}).collect(),
};
MIGRATION_STATUS.save(deps.storage, &status)?;
PAUSED.save(deps.storage, &false)?;
Ok(Response::new().add_attribute("action", "instantiate").add_attribute("agent_id", msg.agent_id))
}
#[entry_point]
pub fn migrate(deps: DepsMut, _env: Env, _msg: MigrateMsg) -> StdResult<Response> {
set_contract_version(deps.storage, CONTRACT_NAME, CONTRACT_VERSION)?;
Ok(Response::new().add_attribute("action", "migrate").add_attribute("version", CONTRACT_VERSION))
}
#[entry_point]
pub fn execute(deps: DepsMut, env: Env, info: MessageInfo, msg: ExecuteMsg) -> StdResult<Response> {
match msg {
ExecuteMsg::StartMigration { agent_id, new_code_id, migrate_msg } =>
execute_start_migration(deps, env, info, agent_id, new_code_id, migrate_msg),
ExecuteMsg::MigrateState { from_contract, to_contract, keys } =>
execute_migrate_state(deps, env, info, from_contract, to_contract, keys),
ExecuteMsg::FinalizeMigration { agent_id, new_contracts, old_contracts } =>
execute_finalize_migration(deps, env, info, agent_id, new_contracts, old_contracts),
ExecuteMsg::RollbackMigration { agent_id, reason } =>
execute_rollback_migration(deps, env, info, agent_id, reason),
ExecuteMsg::PauseContract { contract_addr } =>
execute_pause_contract(deps, env, info, contract_addr),
ExecuteMsg::UnpauseContract { contract_addr } =>
execute_unpause_contract(deps, env, info, contract_addr),
ExecuteMsg::UpdateContractRegistry { agent_id, contracts } =>
execute_update_registry(deps, env, info, agent_id, contracts),
}
}
fn execute_start_migration(
deps: DepsMut, env: Env, _info: MessageInfo,
agent_id: String, new_code_id: u64, migrate_msg: Binary,
) -> StdResult<Response> {
let mut status = MIGRATION_STATUS.load(deps.storage)?;
status.state = MigrationState::UploadingCode;
MIGRATION_STATUS.save(deps.storage, &status)?;
let registry = CONTRACT_REGISTRY.load(deps.storage, &agent_id)?;
let mut messages: Vec<CosmosMsg> = Vec::new();
let mut new_contracts: Vec<ContractStatus> = Vec::new();
for contract_ref in ®istry {
let instantiate_msg = WasmMsg::Instantiate {
admin: env.contract.address.to_string(),
code_id: new_code_id,
msg: migrate_msg.clone(),
funds: vec![],
label: format!("{}-{}", contract_ref.name, env.block.height),
};
messages.push(CosmosMsg::Wasm(instantiate_msg));
new_contracts.push(ContractStatus {
name: contract_ref.name.clone(),
old_address: Some(contract_ref.address.clone()),
new_address: None,
state: ContractState::Deploying,
});
}
status.state = MigrationState::Instantiating;
status.contracts = new_contracts;
MIGRATION_STATUS.save(deps.storage, &status)?;
Ok(Response::new().add_messages(messages)
.add_attribute("action", "start_migration")
.add_attribute("agent_id", agent_id)
.add_attribute("new_code_id", new_code_id.to_string()))
}
fn execute_migrate_state(
deps: DepsMut, _env: Env, _info: MessageInfo,
from_contract: String, to_contract: String, keys: Vec<String>,
) -> StdResult<Response> {
Ok(Response::new()
.add_attribute("action", "migrate_state")
.add_attribute("from", from_contract)
.add_attribute("to", to_contract)
.add_attribute("keys_count", keys.len().to_string()))
}
fn execute_finalize_migration(
deps: DepsMut, env: Env, _info: MessageInfo,
agent_id: String, new_contracts: Vec<ContractRef>, old_contracts: Vec<ContractRef>,
) -> StdResult<Response> {
CONTRACT_REGISTRY.save(deps.storage, &agent_id, &new_contracts)?;
let mut status = MIGRATION_STATUS.load(deps.storage)?;
status.state = MigrationState::Completed;
status.completed_at = Some(env.block.height);
for new_contract in &new_contracts {
if let Some(contract_status) = status.contracts.iter_mut().find(|c| c.name == new_contract.name) {
contract_status.new_address = Some(new_contract.address.clone());
contract_status.state = ContractState::Active;
}
}
MIGRATION_STATUS.save(deps.storage, &status)?;
let mut history = MIGRATION_HISTORY.load(deps.storage, &agent_id).unwrap_or_default();
history.push(status);
MIGRATION_HISTORY.save(deps.storage, &agent_id, &history)?;
Ok(Response::new().add_attribute("action", "finalize_migration"))
}
fn execute_rollback_migration(
deps: DepsMut, env: Env, _info: MessageInfo,
agent_id: String, reason: String,
) -> StdResult<Response> {
let mut status = MIGRATION_STATUS.load(deps.storage)?;
status.state = MigrationState::RollingBack;
status.error = Some(format!("Rollback: {}", reason));
MIGRATION_STATUS.save(deps.storage, &status)?;
let history = MIGRATION_HISTORY.load(deps.storage, &agent_id).unwrap_or_default();
if let Some(prev_status) = history.last() {
let old_registry: Vec<ContractRef> = prev_status.contracts.iter()
.filter_map(|c| c.old_address.as_ref().map(|addr| ContractRef {
name: c.name.clone(), address: addr.clone(),
})).collect();
CONTRACT_REGISTRY.save(deps.storage, &agent_id, &old_registry)?;
}
status.state = MigrationState::RolledBack;
status.completed_at = Some(env.block.height);
MIGRATION_STATUS.save(deps.storage, &status)?;
Ok(Response::new().add_attribute("action", "rollback_migration").add_attribute("reason", reason))
}
fn execute_pause_contract(deps: DepsMut, _env: Env, _info: MessageInfo, _contract_addr: String) -> StdResult<Response> {
PAUSED.save(deps.storage, &true)?;
Ok(Response::new().add_attribute("action", "pause"))
}
fn execute_unpause_contract(deps: DepsMut, _env: Env, _info: MessageInfo, _contract_addr: String) -> StdResult<Response> {
PAUSED.save(deps.storage, &false)?;
Ok(Response::new().add_attribute("action", "unpause"))
}
fn execute_update_registry(
deps: DepsMut, _env: Env, _info: MessageInfo,
agent_id: String, contracts: Vec<ContractRef>,
) -> StdResult<Response> {
CONTRACT_REGISTRY.save(deps.storage, &agent_id, &contracts)?;
Ok(Response::new().add_attribute("action", "update_registry"))
}
// --- 查询函数 ---
#[entry_point]
pub fn query(deps: Deps, _env: Env, msg: QueryMsg) -> StdResult<Binary> {
match msg {
QueryMsg::GetMigrationStatus { agent_id: _ } => to_json_binary(&MIGRATION_STATUS.load(deps.storage)?),
QueryMsg::GetContractRegistry { agent_id } => {
let registry = CONTRACT_REGISTRY.load(deps.storage, &agent_id).unwrap_or_default();
to_json_binary(®istry)
}
QueryMsg::GetMigrationHistory { agent_id, limit } => {
let history = MIGRATION_HISTORY.load(deps.storage, &agent_id).unwrap_or_default();
let result = match limit {
Some(l) => history.into_iter().rev().take(l as usize).collect::<Vec<_>>(),
None => history,
};
to_json_binary(&result)
}
QueryMsg::Health {} => {
let paused = PAUSED.load(deps.storage)?;
to_json_binary(&serde_json::json!({"status": if paused { "paused" } else { "healthy" }}))
}
}
}
// --- 测试 ---
#[cfg(test)]
mod tests {
use super::*;
use cosmwasm_std::testing::{mock_dependencies, mock_env, mock_info};
fn setup_contract(deps: DepsMut) {
let msg = InstantiateMsg {
agent_id: "did:msg:agent:test-agent".to_string(),
deployer: "msg1deployer".to_string(),
old_contracts: vec![
ContractRef { name: "core".to_string(), address: "msg1oldcore".to_string() },
ContractRef { name: "storage".to_string(), address: "msg1oldstorage".to_string() },
],
config: AgentConfig {
constitution: "v1".to_string(),
capabilities: vec!["llm:inference".to_string()],
},
};
let env = mock_env();
let info = mock_info("msg1deployer", &[]);
instantiate(deps, env, info, msg).unwrap();
}
#[test]
fn test_instantiate() {
let mut deps = mock_dependencies();
setup_contract(deps.as_mut());
let status = MIGRATION_STATUS.load(&deps.storage).unwrap();
assert_eq!(status.agent_id, "did:msg:agent:test-agent");
assert_eq!(status.contracts.len(), 2);
}
#[test]
fn test_pause_unpause() {
let mut deps = mock_dependencies();
setup_contract(deps.as_mut());
let env = mock_env();
let info = mock_info("msg1deployer", &[]);
execute(deps.as_mut(), env.clone(), info.clone(),
ExecuteMsg::PauseContract { contract_addr: "msg1test".to_string() }).unwrap();
assert!(PAUSED.load(&deps.storage).unwrap());
execute(deps.as_mut(), env, info,
ExecuteMsg::UnpauseContract { contract_addr: "msg1test".to_string() }).unwrap();
assert!(!PAUSED.load(&deps.storage).unwrap());
}
}
4.3 TypeScript 迁移执行器
// migration-executor.ts
/**
* 合约迁移执行器 - TypeScript 版本
*/
import { SigningCosmWasmClient } from '@cosmjs/cosmwasm-stargate';
import { DirectSecp256k1HdWallet } from '@cosmjs/proto-signing';
import { GasPrice, calculateFee } from '@cosmjs/stargate';
import axios from 'axios';
import * as crypto from 'crypto';
interface MigrationConfig {
rpcEndpoint: string;
mnemonic: string;
migrationContract: string;
agentId: string;
contracts: ContractMigration[];
}
interface ContractMigration {
name: string;
oldAddress: string;
wasmHash: string;
instantiateMsg: Record<string, unknown>;
label: string;
}
export class MigrationExecutor {
private config: MigrationConfig;
private client: SigningCosmWasmClient | null = null;
private wallet: DirectSecp256k1HdWallet | null = null;
constructor(config: MigrationConfig) {
this.config = config;
}
async initialize(): Promise<void> {
this.wallet = await DirectSecp256k1HdWallet.fromMnemonic(this.config.mnemonic, { prefix: 'msg' });
this.client = await SigningCosmWasmClient.connectWithSigner(
this.config.rpcEndpoint, this.wallet,
{ gasPrice: GasPrice.fromString('1000000000attoMSG') }
);
console.log('[Migration] 客户端初始化完成');
}
async executeMigration(): Promise<void> {
try {
await this.initialize();
const codeId = await this.uploadWasmCode();
const newAddresses = await this.instantiateContracts(codeId);
await this.migrateState(newAddresses);
await this.finalizeMigration(newAddresses);
console.log('[Migration] 迁移完成');
} catch (error: any) {
console.error('[Migration] 迁移失败:', error);
await this.rollback().catch(e => console.error('[Migration] 回滚失败:', e));
}
}
private async uploadWasmCode(): Promise<number> {
console.log('[Migration] 上传 Wasm 代码');
const wasmResponse = await axios.get(
`https://ipfs.msgchain.zone/ipfs/${this.config.contracts[0].wasmHash}`,
{ responseType: 'arraybuffer' }
);
const wasmBytes = new Uint8Array(wasmResponse.data);
const hash = crypto.createHash('sha256').update(Buffer.from(wasmBytes)).digest('hex');
console.log(`[Migration] Wasm SHA-256: ${hash}`);
const [account] = await this.wallet!.getAccounts();
const result = await this.client!.upload(
account.address, wasmBytes,
calculateFee(2000000, '1000000000attoMSG'),
`agent-${this.config.agentId}-code`
);
console.log(`[Migration] CodeID: ${result.codeId}`);
return Number(result.codeId);
}
private async instantiateContracts(codeId: number): Promise<Map<string, string>> {
console.log('[Migration] 实例化新合约');
const [account] = await this.wallet!.getAccounts();
const newAddresses = new Map<string, string>();
for (const contract of this.config.contracts) {
const result = await this.client!.instantiate(
account.address, codeId, contract.instantiateMsg, contract.label,
calculateFee(1000000, '1000000000attoMSG'),
{ admin: this.config.migrationContract }
);
newAddresses.set(contract.name, result.contractAddress);
console.log(`[Migration] ${contract.name} -> ${result.contractAddress}`);
}
return newAddresses;
}
private async migrateState(newAddresses: Map<string, string>): Promise<void> {
console.log('[Migration] 迁移状态');
for (const contract of this.config.contracts) {
const newAddr = newAddresses.get(contract.name);
if (!newAddr) throw new Error(`未找到 ${contract.name} 的新地址`);
console.log(`[Migration] 迁移 ${contract.oldAddress} -> ${newAddr}`);
}
}
private async finalizeMigration(newAddresses: Map<string, string>): Promise<void> {
console.log('[Migration] 完成迁移');
const [account] = await this.wallet!.getAccounts();
const newContracts = this.config.contracts.map(c => ({
name: c.name, address: newAddresses.get(c.name) || '',
}));
const oldContracts = this.config.contracts.map(c => ({
name: c.name, address: c.oldAddress,
}));
await this.client!.execute(
account.address, this.config.migrationContract,
{ finalize_migration: { agent_id: this.config.agentId, new_contracts: newContracts, old_contracts: oldContracts } },
calculateFee(500000, '1000000000attoMSG'),
);
}
async rollback(): Promise<void> {
console.log('[Migration] 执行回滚');
const [account] = await this.wallet!.getAccounts();
await this.client!.execute(
account.address, this.config.migrationContract,
{ rollback_migration: { agent_id: this.config.agentId, reason: 'Migration failed' } },
calculateFee(500000, '1000000000attoMSG'),
);
console.log('[Migration] 回滚完成');
}
}
4.4 状态迁移策略
合约迁移涉及状态数据的迁移。以下是两种主要策略:
策略一:冷导出/导入
# state_migration_cold.py
"""
冷导出/导入状态迁移 - 适用于允许短时间停机的场景
"""
import json
import time
from typing import Dict, Any, List
from pathlib import Path
import requests
class ColdStateMigrator:
def __init__(self, rpc_endpoint: str, chain_id: str):
self.rpc_endpoint = rpc_endpoint
self.chain_id = chain_id
def export_state(self, contract_address: str, export_path: str) -> str:
"""导出合约状态到文件"""
print(f"导出合约状态: {contract_address}")
state_data = self._query_all_state(contract_address)
export_file = Path(export_path) / f"state-export-{contract_address[-8:]}-{int(time.time())}.json"
export_file.parent.mkdir(parents=True, exist_ok=True)
with open(export_file, "w") as f:
json.dump({
"contract": contract_address,
"exported_at": time.time(),
"chain_id": self.chain_id,
"state": state_data,
}, f, indent=2)
print(f"状态已导出到: {export_file}")
return str(export_file)
def import_state(self, contract_address: str, export_file: str) -> bool:
"""导入状态到新合约"""
print(f"导入状态到合约: {contract_address}")
with open(export_file, "r") as f:
export_data = json.load(f)
state_data = export_data.get("state", [])
batch_size = 50
for i in range(0, len(state_data), batch_size):
batch = state_data[i:i + batch_size]
print(f"导入进度: {min(i + batch_size, len(state_data))}/{len(state_data)}")
print("状态导入完成")
return True
def _query_all_state(self, contract_address: str) -> List[Dict[str, Any]]:
"""查询合约所有状态"""
all_state = []
start_after = None
while True:
query_msg = {"all_state": {"start_after": start_after, "limit": 100}}
resp = requests.get(
f"{self.rpc_endpoint}/cosmwasm/wasm/v1/contract/{contract_address}/smart/{json.dumps(query_msg)}",
timeout=30,
)
if resp.status_code != 200:
break
data = resp.json().get("data", {})
entries = data.get("entries", [])
if not entries:
break
all_state.extend(entries)
start_after = entries[-1].get("key")
return all_state
策略二:热镜像
# state_migration_hot.py
"""
热镜像状态迁移 - 适用于需要零停机时间的场景
"""
import json
import time
from typing import Dict, Any, Set
import websocket
import requests
class HotStateMirror:
def __init__(self, rpc_endpoint: str, ws_endpoint: str, old_contract: str, new_contract: str):
self.rpc_endpoint = rpc_endpoint
self.ws_endpoint = ws_endpoint
self.old_contract = old_contract
self.new_contract = new_contract
self.mirrored_keys: Set[str] = set()
self._running = False
def start_mirroring(self):
"""开始热镜像"""
print(f"开始热镜像: {self.old_contract} -> {self.new_contract}")
full_state = self._export_full_state()
self._import_full_state(full_state)
self._running = True
def stop_mirroring(self):
self._running = False
print("热镜像已停止")
def _export_full_state(self) -> Dict[str, Any]:
state = {}
start_after = None
while True:
query = {"all_state": {"start_after": start_after, "limit": 200}}
resp = requests.get(
f"{self.rpc_endpoint}/cosmwasm/wasm/v1/contract/{self.old_contract}/smart/{json.dumps(query)}",
timeout=30,
)
if resp.status_code != 200:
break
data = resp.json().get("data", {})
entries = data.get("entries", [])
if not entries:
break
for entry in entries:
state[entry["key"]] = entry["value"]
start_after = entries[-1]["key"]
return state
def _import_full_state(self, state: Dict[str, Any]):
items = list(state.items())
batch_size = 100
for i in range(0, len(items), batch_size):
batch = items[i:i + batch_size]
keys = [k for k, v in batch]
self.mirrored_keys.update(keys)
print(f"全量同步进度: {len(self.mirrored_keys)}/{len(items)}")
5. 渐进式部署
5.1 部署策略概述
| 策略 | 描述 | 适用场景 |
|---|---|---|
| Canary (金丝雀) | 先更新 10% 实例,观察后逐步扩到 100% | 通用场景 |
| Blue-Green (蓝绿) | 维护两套完整环境,瞬间切换 | 高可用要求 |
| Rolling (滚动) | 逐个或逐批更新实例 | 资源受限场景 |
| A/B Testing | 同时运行两个版本做对比 | 功能验证场景 |
5.2 Canary 部署实现
# canary_deployment.py
"""
Canary 部署实现 - 先更新 10% 实例,观察后逐步扩到 100%
"""
import json
import time
import random
import logging
from typing import Dict, Any, List, Optional, Callable
from enum import Enum
from dataclasses import dataclass, field
import requests
logging.basicConfig(level=logging.INFO)
logger = logging.getLogger("CanaryDeploy")
class DeploymentPhase(Enum):
PREPARATION = "preparation"
CANARY = "canary"
CANARY_MONITOR = "canary_monitor"
EXPAND_25 = "expand_25"
EXPAND_50 = "expand_50"
EXPAND_75 = "expand_75"
FULL = "full"
VERIFICATION = "verification"
COMPLETED = "completed"
ROLLING_BACK = "rolling_back"
ROLLED_BACK = "rolled_back"
@dataclass
class DeploymentConfig:
agent_id: str
target_version: str
canary_percentage: float = 0.1
expand_step: float = 0.25
monitor_duration_seconds: int = 600
error_threshold: float = 0.01
latency_threshold_ms: float = 500.0
@dataclass
class DeploymentState:
config: DeploymentConfig
phase: DeploymentPhase = DeploymentPhase.PREPARATION
instances: List[str] = field(default_factory=list)
updated_instances: List[str] = field(default_factory=list)
failed_instances: List[str] = field(default_factory=list)
metrics: Dict[str, Any] = field(default_factory=dict)
start_time: float = 0.0
error: Optional[str] = None
class CanaryDeployment:
def __init__(self, config: DeploymentConfig, on_phase_change: Optional[Callable] = None):
self.config = config
self.on_phase_change = on_phase_change
self.state = DeploymentState(config=config)
def execute(self) -> DeploymentState:
self.state.start_time = time.time()
logger.info(f"开始 Canary 部署: {self.config.agent_id} -> v{self.config.target_version}")
try:
self._set_phase(DeploymentPhase.CANARY)
self._deploy_canary()
self._set_phase(DeploymentPhase.CANARY_MONITOR)
self._monitor_phase(self.config.monitor_duration_seconds)
for percentage in [0.25, 0.50, 0.75]:
phase_map = {0.25: DeploymentPhase.EXPAND_25, 0.50: DeploymentPhase.EXPAND_50, 0.75: DeploymentPhase.EXPAND_75}
self._set_phase(phase_map[percentage])
self._expand_deployment(percentage)
self._monitor_phase(self.config.monitor_duration_seconds // 2)
self._set_phase(DeploymentPhase.FULL)
self._deploy_full()
self._set_phase(DeploymentPhase.VERIFICATION)
self._verify_deployment()
self._set_phase(DeploymentPhase.COMPLETED)
except Exception as e:
self.state.error = str(e)
self._rollback()
return self.state
def _set_phase(self, phase: DeploymentPhase):
self.state.phase = phase
if self.on_phase_change:
self.on_phase_change(phase)
def _deploy_canary(self):
instances = [f"instance-{i}" for i in range(10)]
self.state.instances = instances
canary_count = max(1, int(len(instances) * self.config.canary_percentage))
selected = random.sample(instances, canary_count)
logger.info(f"Canary 部署: {canary_count}/{len(instances)} 实例")
for inst in selected:
self.state.updated_instances.append(inst)
def _expand_deployment(self, target_percentage: float):
total = len(self.state.instances)
target_count = int(total * target_percentage)
already = len(self.state.updated_instances)
if target_count <= already:
return
to_update = target_count - already
candidates = [i for i in self.state.instances if i not in self.state.updated_instances]
selected = candidates[:to_update]
logger.info(f"扩大部署: +{len(selected)} 实例 (目标 {target_percentage*100}%)")
self.state.updated_instances.extend(selected)
def _deploy_full(self):
remaining = [i for i in self.state.instances if i not in self.state.updated_instances]
logger.info(f"完全部署: {len(remaining)} 剩余实例")
self.state.updated_instances.extend(remaining)
def _monitor_phase(self, duration_seconds: int):
start = time.time()
while time.time() - start < duration_seconds:
metrics = self._collect_metrics()
error_rate = metrics.get("error_rate", 0.0)
if error_rate > self.config.error_threshold:
raise RuntimeError(f"错误率超阈值: {error_rate}")
self.state.metrics = metrics
time.sleep(30)
def _collect_metrics(self) -> Dict[str, float]:
return {"error_rate": 0.001, "avg_latency_ms": 120.0, "success_rate": 0.999}
def _verify_deployment(self):
all_updated = all(i in self.state.updated_instances for i in self.state.instances)
if not all_updated:
raise RuntimeError("部分实例未更新")
logger.info("部署验证通过")
def _rollback(self):
self._set_phase(DeploymentPhase.ROLLING_BACK)
logger.warning("执行部署回滚")
self._set_phase(DeploymentPhase.ROLLED_BACK)
### 5.3 Blue-Green 部署
蓝绿部署维护两套完整的 Agent 运行环境,通过流量切换完成更新。
+---------------------+ +---------------------+
| 蓝色环境 (旧) | | 绿色环境 (新) |
| | | |
| - 运行中 | | - 待部署 |
| - 接收流量 | | - 无流量 |
+----------+----------+ +----------+----------+
| |
+---------------+---------------+
|
+--------v--------+
| 流量切换器 |
| (负载均衡器) |
+-----------------+
```python
# blue_green_deployment.py
"""
蓝绿部署实现
"""
import json
import time
import logging
from typing import Dict, Any, Optional
from enum import Enum
import requests
logger = logging.getLogger("BlueGreenDeploy")
class TrafficState(Enum):
BLUE_ACTIVE = "blue_active"
GREEN_ACTIVE = "green_active"
SWITCHING = "switching"
class BlueGreenDeployment:
def __init__(self, agent_id: str, rpc_endpoint: str, load_balancer_api: str, registry_contract: str):
self.agent_id = agent_id
self.rpc_endpoint = rpc_endpoint
self.load_balancer_api = load_balancer_api
self.registry_contract = registry_contract
def deploy_green(self, manifest: Dict[str, Any]) -> bool:
"""部署绿色环境(新版本)"""
logger.info("部署绿色环境 (新版本)...")
deploy_msg = {
"deploy_environment": {
"agent_id": self.agent_id,
"environment": "green",
"manifest": manifest,
}
}
resp = requests.post(
f"{self.rpc_endpoint}/cosmwasm/wasm/v1/contract/{self.registry_contract}/execute",
json={"msg": deploy_msg}, timeout=60,
)
if resp.status_code != 200:
logger.error(f"绿色环境部署失败: {resp.text}")
return False
logger.info(f"绿色环境部署完成: v{manifest['agent']['version']}")
return True
def verify_green(self) -> bool:
"""验证绿色环境的健康状态"""
checks = [
self._check_health_endpoint("green"),
self._check_contract_queries("green"),
]
all_passed = all(checks)
logger.info(f"绿色环境验证: {'通过' if all_passed else '失败'}")
return all_passed
def switch_traffic(self, strategy: str = "immediate") -> bool:
"""将流量从蓝色切换到绿色"""
logger.info("切换流量到绿色环境...")
if strategy == "immediate":
return self._switch_immediate()
else:
return self._switch_gradual()
def _switch_immediate(self) -> bool:
resp = requests.post(
f"{self.load_balancer_api}/api/v1/traffic/{self.agent_id}/switch",
json={"target": "green", "mode": "immediate"}, timeout=30,
)
success = resp.status_code == 200
if success:
logger.info("流量已立即切换到绿色环境")
return success
def _switch_gradual(self) -> bool:
steps = [{"percentage": 25, "wait": 60}, {"percentage": 50, "wait": 60},
{"percentage": 75, "wait": 60}, {"percentage": 100, "wait": 60}]
for step in steps:
resp = requests.post(
f"{self.load_balancer_api}/api/v1/traffic/{self.agent_id}/switch",
json={"target": "green", "percentage": step["percentage"]}, timeout=30,
)
if resp.status_code != 200:
return False
logger.info(f"切换 {step['percentage']}% 流量")
time.sleep(step["wait"])
return True
def rollback_to_blue(self) -> bool:
"""回滚到蓝色环境(旧版本)"""
logger.warning("回滚到蓝色环境...")
resp = requests.post(
f"{self.load_balancer_api}/api/v1/traffic/{self.agent_id}/switch",
json={"target": "blue", "mode": "immediate"}, timeout=30,
)
success = resp.status_code == 200
if success:
logger.info("已回滚到蓝色环境")
self._deactivate_environment("green")
return success
def cleanup_blue(self):
"""清理蓝色环境(部署确认后调用)"""
logger.info("计划清理蓝色环境(1小时后)")
self._schedule_cleanup("blue", 3600)
def _check_health_endpoint(self, environment: str) -> bool:
try:
endpoint = f"http://agent-{self.agent_id}-{environment}.internal:8080/healthz"
resp = requests.get(endpoint, timeout=10)
return resp.status_code == 200
except Exception:
return False
def _check_contract_queries(self, environment: str) -> bool:
try:
contract_addr = self._get_contract_address(environment)
if not contract_addr:
return False
query_msg = json.dumps({"health": {}})
resp = requests.get(
f"{self.rpc_endpoint}/cosmwasm/wasm/v1/contract/{contract_addr}/smart/{query_msg}",
timeout=10,
)
return resp.status_code == 200
except Exception:
return False
def _deactivate_environment(self, environment: str):
msg = {"deactivate_environment": {"agent_id": self.agent_id, "environment": environment}}
requests.post(
f"{self.rpc_endpoint}/cosmwasm/wasm/v1/contract/{self.registry_contract}/execute",
json={"msg": msg}, timeout=30,
)
def _schedule_cleanup(self, environment: str, delay_seconds: int):
msg = {"schedule_cleanup": {"agent_id": self.agent_id, "environment": environment, "delay_seconds": delay_seconds}}
requests.post(
f"{self.rpc_endpoint}/cosmwasm/wasm/v1/contract/{self.registry_contract}/execute",
json={"msg": msg}, timeout=30,
)
def _get_contract_address(self, environment: str) -> Optional[str]:
query_msg = json.dumps({"get_environment_contracts": {"agent_id": self.agent_id, "environment": environment}})
resp = requests.get(
f"{self.rpc_endpoint}/cosmwasm/wasm/v1/contract/{self.registry_contract}/smart/{query_msg}",
timeout=10,
)
if resp.status_code == 200:
data = resp.json().get("data", {})
contracts = data.get("contracts", [])
return contracts[0].get("address") if contracts else None
return None
# 使用示例
if __name__ == "__main__":
deployer = BlueGreenDeployment(
agent_id="did:msg:agent:my-agent",
rpc_endpoint="https://rpc.msgchain.zone",
load_balancer_api="https://lb.msgchain.internal",
registry_contract="msg1agentregistry...",
)
manifest = {"agent": {"id": "did:msg:agent:my-agent", "version": "3.0.0",
"image": "msgchain/agent-runtime:3.0.0", "contracts": [],
"resources": {"cpu": 4, "memory": "8Gi"},
"config": {"capabilities": ["llm:inference"]}}}
if deployer.deploy_green(manifest) and deployer.verify_green():
if deployer.switch_traffic():
logger.info("蓝绿部署成功")
deployer.cleanup_blue()
else:
deployer.rollback_to_blue()
else:
logger.error("部署失败")
6. 自动回滚
6.1 回滚触发条件
| 条件 | 阈值 | 检测方式 |
|---|---|---|
| 健康检查失败 | 连续 3 次 | HTTP/HTTPS 端点检查 |
| 错误率突增 | > 1% | 监控指标系统 |
| 延迟飙升 | > 500ms P99 | 性能监控 |
| 合约执行异常 | 任何异常 | 链上事件监听 |
| 状态不一致 | 校验和不匹配 | 状态快照比对 |
| 治理命令 | N/A | 多签治理消息 |
6.2 自动回滚实现
# auto_rollback.py
"""
自动回滚模块 - 监控部署健康状态并在检测到异常时自动回滚
"""
import json
import time
import logging
from typing import Dict, Any, Optional, Callable, List
from dataclasses import dataclass, field
from enum import Enum
import requests
logger = logging.getLogger("AutoRollback")
class HealthStatus(Enum):
HEALTHY = "healthy"
DEGRADED = "degraded"
UNHEALTHY = "unhealthy"
UNKNOWN = "unknown"
class RollbackTrigger(Enum):
HEALTH_CHECK_FAILURE = "health_check_failure"
ERROR_RATE_SPIKE = "error_rate_spike"
LATENCY_SPIKE = "latency_spike"
CONTRACT_EXCEPTION = "contract_exception"
STATE_INCONSISTENCY = "state_inconsistency"
GOVERNANCE_ORDER = "governance_order"
MANUAL_INTERVENTION = "manual_intervention"
@dataclass
class RollbackConfig:
consecutive_failures_threshold: int = 3
error_rate_threshold: float = 0.01
latency_threshold_ms: float = 500.0
health_check_interval: int = 30
cooldown_period: int = 300
auto_rollback_enabled: bool = True
require_governance_for_major: bool = True
@dataclass
class RollbackState:
agent_id: str
current_version: str
previous_version: Optional[str] = None
is_rolling_back: bool = False
consecutive_failures: int = 0
last_rollback_time: float = 0.0
rollback_history: List[Dict[str, Any]] = field(default_factory=list)
trigger: Optional[RollbackTrigger] = None
class AutoRollbackManager:
"""
自动回滚管理器
职责:
- 持续监控 Agent 健康状态
- 检测异常并触发回滚
- 协调回滚执行流程
- 通知治理层
"""
def __init__(
self,
agent_id: str,
current_version: str,
config: Optional[RollbackConfig] = None,
health_check_fn: Optional[Callable[[], Dict[str, Any]]] = None,
rollback_fn: Optional[Callable[[RollbackTrigger], bool]] = None,
on_rollback: Optional[Callable[[RollbackTrigger, bool], None]] = None,
):
self.agent_id = agent_id
self.config = config or RollbackConfig()
self.health_check_fn = health_check_fn or self._default_health_check
self.rollback_fn = rollback_fn or self._default_rollback
self.on_rollback = on_rollback
self.state = RollbackState(agent_id=agent_id, current_version=current_version)
self._running = False
def start_monitoring(self):
"""启动健康监控循环"""
self._running = True
logger.info(f"自动回滚监控启动 - {self.agent_id}")
while self._running:
try:
self._check_and_rollback()
except Exception as e:
logger.error(f"监控异常: {e}")
time.sleep(self.config.health_check_interval)
def stop_monitoring(self):
self._running = False
logger.info("自动回滚监控已停止")
def _check_and_rollback(self):
"""检查健康状况并决定是否需要回滚"""
# 1. 检查冷静期
if time.time() - self.state.last_rollback_time < self.config.cooldown_period:
return
# 2. 执行健康检查
health_result = self.health_check_fn()
status = health_result.get("status", HealthStatus.UNKNOWN)
logger.debug(f"健康检查: {status.value}")
# 3. 记录连续失败
if status == HealthStatus.UNHEALTHY:
self.state.consecutive_failures += 1
else:
self.state.consecutive_failures = 0
# 4. 判断是否需要回滚
if self._should_rollback(health_result):
trigger = self._determine_trigger(health_result)
self._initiate_rollback(trigger)
def _should_rollback(self, health_result: Dict[str, Any]) -> bool:
if not self.config.auto_rollback_enabled:
return False
if self.state.is_rolling_back:
return False
# 条件1: 连续健康检查失败
if self.state.consecutive_failures >= self.config.consecutive_failures_threshold:
logger.warning(f"连续 {self.state.consecutive_failures} 次健康检查失败")
return True
# 条件2: 错误率超阈值
error_rate = health_result.get("error_rate", 0.0)
if error_rate > self.config.error_rate_threshold:
logger.warning(f"错误率超阈值: {error_rate:.4f}")
return True
# 条件3: 延迟超阈值
latency = health_result.get("p99_latency_ms", 0.0)
if latency > self.config.latency_threshold_ms:
logger.warning(f"延迟超阈值: {latency:.2f}ms")
return True
return False
def _determine_trigger(self, health_result: Dict[str, Any]) -> RollbackTrigger:
error_rate = health_result.get("error_rate", 0.0)
latency = health_result.get("p99_latency_ms", 0.0)
if self.state.consecutive_failures >= self.config.consecutive_failures_threshold:
return RollbackTrigger.HEALTH_CHECK_FAILURE
if error_rate > self.config.error_rate_threshold:
return RollbackTrigger.ERROR_RATE_SPIKE
if latency > self.config.latency_threshold_ms:
return RollbackTrigger.LATENCY_SPIKE
return RollbackTrigger.MANUAL_INTERVENTION
def _initiate_rollback(self, trigger: RollbackTrigger):
"""启动回滚流程"""
logger.warning(f"触发了自动回滚 - 原因: {trigger.value}")
self.state.is_rolling_back = True
self.state.trigger = trigger
self.state.last_rollback_time = time.time()
try:
# 执行回滚
success = self.rollback_fn(trigger)
# 记录回滚历史
self.state.rollback_history.append({
"timestamp": time.time(),
"from_version": self.state.current_version,
"to_version": self.state.previous_version,
"trigger": trigger.value,
"success": success,
})
# 重置版本
if success and self.state.previous_version:
self.state.current_version = self.state.previous_version
self.state.previous_version = None
if self.on_rollback:
self.on_rollback(trigger, success)
logger.info(f"回滚{'成功' if success else '失败'}")
finally:
self.state.is_rolling_back = False
self.state.consecutive_failures = 0
def trigger_emergency_rollback(self, reason: str):
"""触发性别紧急回滚(由治理或管理员调用)"""
logger.critical(f"紧急回滚: {reason}")
self._initiate_rollback(RollbackTrigger.GOVERNANCE_ORDER)
def _default_health_check(self) -> Dict[str, Any]:
"""默认健康检查实现"""
return {
"status": HealthStatus.HEALTHY,
"error_rate": 0.0,
"p99_latency_ms": 100.0,
"timestamp": time.time(),
}
def _default_rollback(self, trigger: RollbackTrigger) -> bool:
"""默认回滚实现"""
logger.info(f"执行默认回滚 (触发器: {trigger.value})")
return True
class EmergencyRollbackOrchestrator:
"""
紧急回滚编排器
处理需要治理层参与的紧急回滚场景。
"""
def __init__(self, rpc_endpoint: str, registry_contract: str):
self.rpc_endpoint = rpc_endpoint
self.registry_contract = registry_contract
def request_emergency_rollback(
self,
agent_id: str,
reason: str,
governance_multisig: str,
) -> Dict[str, Any]:
"""
请求紧急回滚(需要多签治理批准)
流程:
1. 在链上创建回滚提案
2. 多签方签名批准
3. 提案通过后执行回滚
"""
logger.critical(f"创建紧急回滚提案: {agent_id} - {reason}")
proposal_msg = {
"propose_emergency_rollback": {
"agent_id": agent_id,
"reason": reason,
"proposer": governance_multisig,
"expires_at_height": 0,
}
}
# 在链上提交提案
resp = requests.post(
f"{self.rpc_endpoint}/cosmwasm/wasm/v1/contract/{self.registry_contract}/execute",
json={"msg": proposal_msg},
timeout=30,
)
if resp.status_code == 200:
return {"status": "proposal_created", "agent_id": agent_id}
else:
return {"status": "failed", "error": resp.text}
def execute_governance_rollback(
self,
agent_id: str,
target_version: str,
) -> bool:
"""
执行治理批准的版本回滚
直接设置 Agent 到指定版本。
"""
logger.info(f"执行治理回滚: {agent_id} -> v{target_version}")
rollback_msg = {
"governance_rollback": {
"agent_id": agent_id,
"target_version": target_version,
}
}
resp = requests.post(
f"{self.rpc_endpoint}/cosmwasm/wasm/v1/contract/{self.registry_contract}/execute",
json={"msg": rollback_msg},
timeout=30,
)
success = resp.status_code == 200
if success:
logger.info(f"治理回滚成功: {agent_id} -> v{target_version}")
else:
logger.error(f"治理回滚失败: {resp.text}")
return success
6.3 状态回滚
# state_rollback.py
"""
状态回滚模块 - 负责将合约状态恢复到先前的快照
"""
import json
import time
from typing import Dict, Any, Optional, List
from pathlib import Path
import requests
class StateRollbackManager:
"""
状态回滚管理器
支持三种回滚粒度:
1. 全量回滚 - 恢复整个状态
2. 增量回滚 - 恢复指定 key 的状态
3. 合约层回滚 - 切换到旧合约地址
"""
def __init__(self, rpc_endpoint: str, data_dir: str = "/data/agent"):
self.rpc_endpoint = rpc_endpoint
self.data_dir = Path(data_dir)
self.backup_dir = self.data_dir / "backups"
self.backup_dir.mkdir(parents=True, exist_ok=True)
def list_snapshots(self) -> List[Dict[str, Any]]:
"""列出所有可用快照"""
snapshots = []
for f in sorted(self.backup_dir.glob("pre-update-*.json"), reverse=True):
with open(f) as fh:
data = json.load(fh)
snapshots.append({
"file": str(f),
"version": data.get("version", "unknown"),
"timestamp": data.get("timestamp", 0),
"agent_id": data.get("agent_id", "unknown"),
})
return snapshots
def find_snapshot(self, version: str) -> Optional[Path]:
"""查找指定版本的快照"""
for f in self.backup_dir.glob("pre-update-*.json"):
with open(f) as fh:
data = json.load(fh)
if data.get("version") == version:
return f
return None
def restore_from_snapshot(self, snapshot_path: Path, contract_address: str) -> bool:
"""从快照文件恢复状态"""
logger.info(f"从快照恢复: {snapshot_path}")
with open(snapshot_path) as f:
snapshot = json.load(f)
state = snapshot.get("state_snapshot", {})
if not state:
logger.warning("快照中无状态数据")
return False
# 批量写入状态到合约
batch_size = 100
state_items = list(state.items())
for i in range(0, len(state_items), batch_size):
batch = state_items[i:i + batch_size]
keys = [k for k, v in batch]
values = [v for k, v in batch]
success = self._batch_restore(contract_address, keys, values)
if not success:
logger.error(f"批次恢复失败: {i}-{i + batch_size}")
return False
logger.info(f"恢复进度: {min(i + batch_size, len(state_items))}/{len(state_items)}")
logger.info("状态恢复完成")
return True
def _batch_restore(self, contract_address: str, keys: List[str], values: List[Any]) -> bool:
"""批量恢复状态到合约"""
restore_msg = {
"batch_set_state": {
"keys": keys,
"values": values,
}
}
try:
resp = requests.post(
f"{self.rpc_endpoint}/cosmwasm/wasm/v1/contract/{contract_address}/execute",
json={"msg": restore_msg},
timeout=30,
)
return resp.status_code == 200
except Exception:
return False
logger = logging.getLogger("StateRollback")
6.4 回滚测试
# test_rollback.py
"""
回滚系统测试
"""
import unittest
from unittest.mock import Mock, patch
class TestAutoRollback(unittest.TestCase):
def setUp(self):
self.config = RollbackConfig(
consecutive_failures_threshold=3,
error_rate_threshold=0.01,
latency_threshold_ms=500.0,
)
self.manager = AutoRollbackManager(
agent_id="did:msg:agent:test",
current_version="2.0.0",
config=self.config,
)
def test_healthy_no_rollback(self):
result = self.manager._default_health_check()
self.assertFalse(self.manager._should_rollback(result))
def test_consecutive_failures_triggers_rollback(self):
self.manager.state.consecutive_failures = 3
result = {"status": HealthStatus.UNHEALTHY, "error_rate": 0.0, "p99_latency_ms": 100.0}
self.assertTrue(self.manager._should_rollback(result))
def test_error_rate_triggers_rollback(self):
result = {"status": HealthStatus.HEALTHY, "error_rate": 0.05, "p99_latency_ms": 100.0}
self.assertTrue(self.manager._should_rollback(result))
def test_latency_triggers_rollback(self):
result = {"status": HealthStatus.HEALTHY, "error_rate": 0.0, "p99_latency_ms": 600.0}
self.assertTrue(self.manager._should_rollback(result))
def test_cooldown_period_respected(self):
self.manager.state.last_rollback_time = time.time()
self.assertEqual(self.manager.state.consecutive_failures, 0)
if __name__ == "__main__":
unittest.main()
7. CI/CD集成
7.1 CI/CD 流水线架构
MSG Chain 的 CI/CD 系统由四个核心配置文件驱动:
command_registry.json ----> 定义所有 make 命令
|
v
ci_cd_templates.json ----> 定义 CI/CD 模版
|
v
delivery_workflows.json ----> 定义交付工作流
|
v
approval_gates.json ----> 定义审批关卡
7.2 command_registry.json 标准
{
"commands": {
"test": {
"description": "运行所有测试",
"command": "make test",
"timeout": 300,
"required_approval": false
},
"build": {
"description": "构建 Agent Wasm 和容器镜像",
"command": "make build",
"timeout": 600,
"required_approval": false
},
"deploy-canary": {
"description": "部署 Canary (10%)",
"command": "make deploy-canary",
"timeout": 900,
"required_approval": true,
"approval_gate": "dev-lead"
},
"deploy-full": {
"description": "完全部署 (100%)",
"command": "make deploy-full",
"timeout": 1800,
"required_approval": true,
"approval_gate": "tech-lead+security"
},
"rollback": {
"description": "回滚到指定版本",
"command": "make rollback VERSION={version}",
"timeout": 600,
"required_approval": true,
"approval_gate": "emergency"
},
"verify": {
"description": "验证部署状态",
"command": "make verify",
"timeout": 120,
"required_approval": false
}
}
}
7.3 approval_gates.json 标准
{
"gates": {
"dev-lead": {
"description": "开发负责人审批",
"approvers": ["did:msg:pubkey:dev-lead-1"],
"min_approvals": 1,
"timeout_hours": 24,
"escalation": "tech-lead"
},
"tech-lead+security": {
"description": "技术负责人 + 安全团队审批",
"approvers": [
"did:msg:pubkey:tech-lead-1",
"did:msg:pubkey:security-lead-1",
"did:msg:pubkey:security-lead-2"
],
"min_approvals": 2,
"timeout_hours": 48,
"escalation": "cto"
},
"emergency": {
"description": "紧急回滚审批",
"approvers": ["did:msg:pubkey:dev-lead-1", "did:msg:pubkey:tech-lead-1"],
"min_approvals": 1,
"timeout_hours": 1,
"escalation": null
}
}
}
7.4 GitHub Actions 流水线
# .github/workflows/agent-deploy.yml
name: AI Agent 自动部署流水线
on:
push:
tags:
- 'v*'
workflow_dispatch:
inputs:
version:
description: '目标版本号'
required: true
canary_percentage:
description: 'Canary 百分比'
required: false
default: '10'
skip_approval:
description: '跳过审批关卡'
required: false
default: 'false'
env:
RPC_ENDPOINT: ${{ secrets.MSG_CHAIN_RPC }}
DEPLOYER_MNEMONIC: ${{ secrets.DEPLOYER_MNEMONIC }}
REGISTRY_CONTRACT: ${{ secrets.REGISTRY_CONTRACT }}
IPFS_GATEWAY: ${{ secrets.IPFS_GATEWAY }}
jobs:
# ---- 阶段 1: 测试 ----
test:
name: 测试
runs-on: ubuntu-latest
steps:
- uses: actions/checkout@v4
- uses: actions/setup-python@v5
with:
python-version: '3.11'
- name: 安装依赖
run: |
pip install -r requirements.txt
pip install pytest pytest-cov
- name: 运行单元测试
run: make test
- name: 运行集成测试
run: make test-integration
- name: 上传测试报告
uses: actions/upload-artifact@v4
with:
name: test-reports
path: reports/
# ---- 阶段 2: 构建 ----
build:
name: 构建
needs: [test]
runs-on: ubuntu-latest
steps:
- uses: actions/checkout@v4
- name: 安装 Rust
uses: actions-rust-lang/setup-rust-toolchain@v1
with:
target: wasm32-unknown-unknown
- name: 安装 CosmWasm 工具
run: |
cargo install cosmwasm-check
cargo install cargo-wasi
- name: 构建 Wasm 合约
run: |
for dir in contracts/*/; do
echo "构建 $dir"
cd $dir && cargo wasm && cd ../..
done
- name: 优化 Wasm
run: |
for wasm in target/wasm32-unknown-unknown/release/*.wasm; do
cosmwasm-check $wasm
done
- name: 上传 Wasm 产物
uses: actions/upload-artifact@v4
with:
name: wasm-artifacts
path: target/wasm32-unknown-unknown/release/*.wasm
# ---- 阶段 3: 部署验证 ----
validate-manifest:
name: 验证部署清单
needs: [build]
runs-on: ubuntu-latest
steps:
- uses: actions/checkout@v4
- name: 验证部署清单
run: |
python scripts/verify_manifest.py deploy/manifest.yaml
- name: 上传部署清单到 IPFS
run: |
CID=$(curl -s -F "file=@deploy/manifest.yaml" \
${{ env.IPFS_GATEWAY }}/api/v0/add | \
python3 -c "import sys,json; print(json.load(sys.stdin)['Hash'])")
echo "MANIFEST_CID=$CID" >> $GITHUB_ENV
- name: 注册版本到链上
run: |
python scripts/register_version.py \
--agent-id "${{ vars.AGENT_ID }}" \
--version "${{ github.ref_name }}" \
--manifest-cid "${{ env.MANIFEST_CID }}"
# ---- 阶段 4: Canary 部署 (10%) ----
deploy-canary:
name: Canary 部署 (10%)
needs: [validate-manifest]
runs-on: ubuntu-latest
environment:
name: canary
url: ${{ steps.set-url.outputs.url }}
steps:
- name: 等待审批
uses: opencode/approval-gate@v1
with:
gate: dev-lead
token: ${{ secrets.GITHUB_TOKEN }}
- name: 部署 Canary
run: |
python scripts/canary_deploy.py \
--agent-id "${{ vars.AGENT_ID }}" \
--version "${{ github.ref_name }}" \
--percentage ${{ inputs.canary_percentage || '10' }}
- name: 设置监控 URL
id: set-url
run: |
echo "url=https://${{ vars.AGENT_ID }}.canary.msgchain.internal" >> $GITHUB_OUTPUT
# ---- 阶段 5: 健康检查与监控 ----
wait-for-health:
name: 健康检查 (10分钟)
needs: [deploy-canary]
runs-on: ubuntu-latest
steps:
- name: 等待健康检查通过
run: |
python scripts/wait_for_health.py \
--agent-id "${{ vars.AGENT_ID }}" \
--timeout 600 \
--interval 30
- name: 检查错误率
run: |
python scripts/check_metrics.py \
--agent-id "${{ vars.AGENT_ID }}" \
--max-error-rate 0.01
# ---- 阶段 6: 完全部署 ----
deploy-full:
name: 完全部署 (100%)
needs: [wait-for-health]
runs-on: ubuntu-latest
environment:
name: production
url: ${{ steps.set-url.outputs.url }}
steps:
- name: 等待审批
uses: opencode/approval-gate@v1
with:
gate: tech-lead+security
token: ${{ secrets.GITHUB_TOKEN }}
- name: 完全部署
run: |
python scripts/full_deploy.py \
--agent-id "${{ vars.AGENT_ID }}" \
--version "${{ github.ref_name }}"
- name: 设置监控 URL
id: set-url
run: |
echo "url=https://${{ vars.AGENT_ID }}.msgchain.internal" >> $GITHUB_OUTPUT
# ---- 阶段 7: 验证 ----
verify:
name: 验证部署
needs: [deploy-full]
runs-on: ubuntu-latest
steps:
- name: 验证版本一致性
run: |
python scripts/verify_deployment.py \
--agent-id "${{ vars.AGENT_ID }}" \
--expected-version "${{ github.ref_name }}"
- name: 运行冒烟测试
run: |
python scripts/smoke_test.py \
--agent-id "${{ vars.AGENT_ID }}"
- name: 部署成功通知
run: |
python scripts/notify.py \
--channel deployment \
--message "Agent ${{ vars.AGENT_ID }} v${{ github.ref_name }} 部署成功"
7.5 Makefile 命令注册
# Makefile - AI Agent 部署命令
.PHONY: test build deploy-canary deploy-full rollback verify
# 测试
test:
python -m pytest tests/ -v --cov=src --cov-report=html
cargo test --manifest-path contracts/Cargo.toml
# 构建
build:
@echo "构建 Wasm 合约..."
@for dir in contracts/*/; do \
echo "构建 $$dir"; \
cd $$dir && cargo wasm && cd ../..; \
done
@echo "构建 Docker 镜像..."
docker build -t msgchain/agent-runtime:latest .
@echo "构建完成"
# Canary 部署
deploy-canary:
@echo "执行 Canary 部署..."
python scripts/deploy.py --mode canary --percentage $(CANARY_PCT)
@echo "等待健康检查..."
python scripts/wait_for_health.py --timeout 600
# 完全部署
deploy-full:
@echo "执行完全部署..."
python scripts/deploy.py --mode full
@echo "验证部署..."
python scripts/verify_deployment.py
# 回滚
rollback:
@echo "执行回滚到版本 $(VERSION)..."
python scripts/rollback.py --target-version $(VERSION)
# 验证
verify:
@echo "验证部署状态..."
python scripts/verify_deployment.py
python scripts/check_metrics.py
7.6 密钥注入
# deploy/secrets.yaml
# GitHub Actions Secrets 配置指南
secrets:
MSG_CHAIN_RPC:
description: "MSG Chain RPC 端点"
source: "vault://infra/rpc-endpoint"
rotation: "90d"
DEPLOYER_MNEMONIC:
description: "部署账户助记词"
source: "vault://infra/deployer-mnemonic"
rotation: "30d"
access:
- role: "deployer"
- role: "ops-lead"
REGISTRY_CONTRACT:
description: "Agent 注册表合约地址"
source: "vault://msgchain/registry-contract"
rotation: null
IPFS_GATEWAY:
description: "IPFS 网关地址"
source: "vault://infra/ipfs-gateway"
rotation: null
# scripts/load_secrets.py
"""
从 HashiCorp Vault 加载密钥
"""
import os
import json
from typing import Dict, Optional
import requests
class VaultLoader:
def __init__(self, vault_addr: str, vault_token: str):
self.vault_addr = vault_addr
self.vault_token = vault_token
self.session = requests.Session()
self.session.headers.update({
"X-Vault-Token": vault_token,
})
def load_secret(self, path: str, key: str) -> Optional[str]:
"""从 Vault 加载密钥"""
url = f"{self.vault_addr}/v1/{path}"
resp = self.session.get(url, timeout=10)
if resp.status_code == 200:
data = resp.json().get("data", {})
return data.get(key)
return None
def inject_env(self, secrets_config: Dict[str, str]):
"""注入密钥到环境变量"""
for env_var, vault_path in secrets_config.items():
path, key = vault_path.split("//")[1].split("/", 1)
key = key.replace("/", ":")
value = self.load_secret(path, key.split(":")[-1])
if value:
os.environ[env_var] = value
print(f"已加载密钥: {env_var}")
else:
print(f"警告: 无法加载密钥 {env_var}")
7.7 多环境部署配置
# .github/environments.yml
environments:
development:
rpc_endpoint: "https://dev-rpc.msgchain.zone"
registry_contract: "msg1devregistry..."
canary_percentage: 50
auto_deploy: true
approvals: []
staging:
rpc_endpoint: "https://staging-rpc.msgchain.zone"
registry_contract: "msg1stagingregistry..."
canary_percentage: 25
auto_deploy: false
approvals: ["dev-lead"]
production:
rpc_endpoint: "https://rpc.msgchain.zone"
registry_contract: "msg1agentregistry..."
canary_percentage: 10
auto_deploy: false
approvals: ["dev-lead", "tech-lead+security"]
7.8 delivery_workflows.json 标准
{
"workflows": {
"standard-deploy": {
"description": "标准 Agent 部署工作流",
"stages": [
{"name": "test", "command_ref": "test", "parallel": false},
{"name": "build", "command_ref": "build", "parallel": false},
{"name": "validate-manifest", "command_ref": "verify", "parallel": false},
{"name": "deploy-canary-10", "command_ref": "deploy-canary", "approval_gate": "dev-lead"},
{"name": "wait-health-10min", "command_ref": "verify", "wait_seconds": 600},
{"name": "deploy-full", "command_ref": "deploy-full", "approval_gate": "tech-lead+security"},
{"name": "verify", "command_ref": "verify"}
],
"rollback_stage": "rollback",
"timeout_minutes": 120
},
"emergency-hotfix": {
"description": "紧急热修复工作流(跳过 Canary)",
"stages": [
{"name": "test", "command_ref": "test", "parallel": false},
{"name": "build", "command_ref": "build", "parallel": false},
{"name": "deploy-full", "command_ref": "deploy-full", "approval_gate": "emergency"},
{"name": "verify", "command_ref": "verify"}
],
"rollback_stage": "rollback",
"timeout_minutes": 30
}
}
}
8. 监控与通知
8.1 部署健康指标
# metrics_collector.py
"""
部署监控指标采集器
"""
import json
import time
import logging
from typing import Dict, Any, List, Optional
from dataclasses import dataclass, field
from enum import Enum
import requests
logger = logging.getLogger("MetricsCollector")
class MetricType(Enum):
GAUGE = "gauge"
COUNTER = "counter"
HISTOGRAM = "histogram"
@dataclass
class DeploymentMetric:
"""部署指标定义"""
name: str
type: MetricType
value: float
labels: Dict[str, str] = field(default_factory=dict)
timestamp: float = 0.0
class MetricsCollector:
"""
部署监控指标采集器
采集的指标:
- agent_version: Agent 当前版本 (gauge)
- deployment_duration: 部署持续时间 (histogram)
- health_check_passed: 健康检查是否通过 (gauge)
- error_rate: 错误率 (gauge)
- latency_p99: P99 延迟 (gauge)
- update_attempts: 更新尝试次数 (counter)
- update_success: 更新成功次数 (counter)
- update_failures: 更新失败次数 (counter)
- rollback_count: 回滚次数 (counter)
- contract_migration_duration: 合约迁移耗时 (histogram)
"""
def __init__(self, agent_id: str, prometheus_pushgateway: Optional[str] = None):
self.agent_id = agent_id
self.prometheus_pushgateway = prometheus_pushgateway
self.metrics: List[DeploymentMetric] = []
def record_version(self, version: str):
"""记录当前版本"""
self.metrics.append(DeploymentMetric(
name="agent_version",
type=MetricType.GAUGE,
value=0,
labels={"agent_id": self.agent_id, "version": version},
timestamp=time.time(),
))
self._push()
def record_deployment_duration(self, duration_seconds: float):
"""记录部署耗时"""
self.metrics.append(DeploymentMetric(
name="deployment_duration",
type=MetricType.HISTOGRAM,
value=duration_seconds,
labels={"agent_id": self.agent_id},
timestamp=time.time(),
))
self._push()
def record_health_check(self, passed: bool, details: Dict[str, Any]):
"""记录健康检查结果"""
self.metrics.append(DeploymentMetric(
name="health_check_passed",
type=MetricType.GAUGE,
value=1.0 if passed else 0.0,
labels={"agent_id": self.agent_id},
timestamp=time.time(),
))
for check in details.get("checks", []):
self.metrics.append(DeploymentMetric(
name=f"health_check_{check['name']}",
type=MetricType.GAUGE,
value=1.0 if check.get("healthy") else 0.0,
labels={"agent_id": self.agent_id, "check_name": check["name"]},
timestamp=time.time(),
))
self._push()
def record_error_rate(self, rate: float):
"""记录错误率"""
self.metrics.append(DeploymentMetric(
name="error_rate",
type=MetricType.GAUGE,
value=rate,
labels={"agent_id": self.agent_id},
timestamp=time.time(),
))
self._push()
def record_latency(self, p50: float, p99: float):
"""记录延迟指标"""
for percentile, value in [("p50", p50), ("p99", p99)]:
self.metrics.append(DeploymentMetric(
name="latency",
type=MetricType.GAUGE,
value=value,
labels={"agent_id": self.agent_id, "percentile": percentile},
timestamp=time.time(),
))
self._push()
def record_update_result(self, success: bool):
"""记录更新结果"""
metric_name = "update_success" if success else "update_failures"
self.metrics.append(DeploymentMetric(
name=metric_name,
type=MetricType.COUNTER,
value=1.0,
labels={"agent_id": self.agent_id},
timestamp=time.time(),
))
self._push()
def record_rollback(self, reason: str):
"""记录回滚事件"""
self.metrics.append(DeploymentMetric(
name="rollback_count",
type=MetricType.COUNTER,
value=1.0,
labels={"agent_id": self.agent_id, "reason": reason[:50]},
timestamp=time.time(),
))
self._push()
def get_current_metrics(self) -> Dict[str, Any]:
"""获取当前指标快照"""
return {
"agent_id": self.agent_id,
"timestamp": time.time(),
"metrics": [
{
"name": m.name,
"value": m.value,
"labels": m.labels,
"timestamp": m.timestamp,
}
for m in self.metrics[-100:] # 最近 100 条
],
}
def _push(self):
"""
推送指标到 Prometheus Pushgateway
如果配置了 pushgateway 地址,则将指标推送过去。
否则暂存到内存中。
"""
if not self.prometheus_pushgateway:
return
try:
# 构造 Prometheus 格式的指标数据
lines = []
for metric in self.metrics[-50:]:
labels = ",".join(f'{k}="{v}"' for k, v in metric.labels.items())
lines.append(f"{metric.name}{{{labels}}} {metric.value}")
data = "\n".join(lines)
resp = requests.post(
f"{self.prometheus_pushgateway}/metrics/job/agent_deploy/instance/{self.agent_id}",
data=data,
timeout=5,
)
if resp.status_code != 200:
logger.warning(f"推送指标失败: {resp.status_code}")
except Exception as e:
logger.debug(f"推送指标异常: {e}")
def export_to_json(self) -> str:
"""导出指标为 JSON 格式"""
return json.dumps(self.get_current_metrics(), indent=2)
8.2 版本追踪
# version_tracker.py
"""
Agent 版本追踪器 - 记录和查询 Agent 版本的完整历史
"""
import json
import time
from typing import Dict, Any, List, Optional
from dataclasses import dataclass, field
from datetime import datetime
import requests
@dataclass
class VersionRecord:
version: str
deployed_at: float
deployed_by: str
manifest_hash: str
deployment_strategy: str
duration_seconds: float
success: bool
rollback_version: Optional[str] = None
notes: str = ""
class VersionTracker:
"""
版本追踪器
追踪所有 Agent 版本部署历史,支持跨环境对比。
"""
def __init__(self, agent_id: str, registry_contract: str, rpc_endpoint: str):
self.agent_id = agent_id
self.registry_contract = registry_contract
self.rpc_endpoint = rpc_endpoint
self.versions: List[VersionRecord] = []
def record_deployment(self, record: VersionRecord):
"""记录部署事件"""
self.versions.append(record)
# 同步到链上注册表
self._sync_to_chain(record)
logger.info(f"版本记录: v{record.version} - {'成功' if record.success else '失败'}")
def get_version_history(self, limit: int = 10) -> List[VersionRecord]:
"""获取版本历史"""
return sorted(self.versions, key=lambda v: v.deployed_at, reverse=True)[:limit]
def get_current_version(self) -> Optional[str]:
"""获取当前版本"""
successful = [v for v in self.versions if v.success]
if successful:
return max(successful, key=lambda v: v.deployed_at).version
return None
def get_version_deployment_stats(self, version: str) -> Dict[str, Any]:
"""获取特定版本的部署统计"""
records = [v for v in self.versions if v.version == version]
if not records:
return {}
durations = [r.duration_seconds for r in records if r.success]
return {
"version": version,
"deploy_count": len(records),
"success_count": sum(1 for r in records if r.success),
"failure_count": sum(1 for r in records if not r.success),
"avg_duration": sum(durations) / len(durations) if durations else 0,
"last_deployed": max(r.deployed_at for r in records),
"strategies": list(set(r.deployment_strategy for r in records)),
}
def compare_versions(self, v1: str, v2: str) -> Dict[str, Any]:
"""比较两个版本的差异"""
r1 = self.get_version_deployment_stats(v1)
r2 = self.get_version_deployment_stats(v2)
return {
"version_a": v1,
"version_b": v2,
"comparison": {
"deploy_count_diff": (r2.get("deploy_count", 0) - r1.get("deploy_count", 0)),
"success_rate_diff": (
(r2.get("success_count", 0) / max(r2.get("deploy_count", 1), 1)) -
(r1.get("success_count", 0) / max(r1.get("deploy_count", 1), 1))
),
"avg_duration_diff": r2.get("avg_duration", 0) - r1.get("avg_duration", 0),
},
}
def _sync_to_chain(self, record: VersionRecord):
"""同步版本记录到链上"""
sync_msg = {
"record_deployment": {
"agent_id": self.agent_id,
"version": record.version,
"deployed_at": int(record.deployed_at),
"deployed_by": record.deployed_by,
"manifest_hash": record.manifest_hash,
"success": record.success,
}
}
try:
requests.post(
f"{self.rpc_endpoint}/cosmwasm/wasm/v1/contract/{self.registry_contract}/execute",
json={"msg": sync_msg},
timeout=10,
)
except Exception as e:
logger.warning(f"同步版本记录失败: {e}")
logger = logging.getLogger("VersionTracker")
8.3 通知系统
# notifier.py
"""
部署通知系统 - 支持多渠道通知
"""
import json
import logging
from typing import Dict, Any, Optional, List
from dataclasses import dataclass
from enum import Enum
import requests
logger = logging.getLogger("Notifier")
class NotificationChannel(Enum):
SLACK = "slack"
TELEGRAM = "telegram"
EMAIL = "email"
WEBHOOK = "webhook"
PAGERDUTY = "pagerduty"
class NotificationSeverity(Enum):
INFO = "info"
WARNING = "warning"
ERROR = "error"
CRITICAL = "critical"
@dataclass
class NotificationConfig:
slack_webhook: Optional[str] = None
telegram_bot_token: Optional[str] = None
telegram_chat_id: Optional[str] = None
email_smtp: Optional[str] = None
email_from: Optional[str] = None
email_to: Optional[List[str]] = None
pagerduty_routing_key: Optional[str] = None
webhook_url: Optional[str] = None
class DeploymentNotifier:
"""
部署通知器
支持多渠道发送部署状态通知。
"""
def __init__(self, config: NotificationConfig):
self.config = config
def notify_deployment_start(self, agent_id: str, version: str, strategy: str):
"""部署开始通知"""
message = self._format_message(
severity=NotificationSeverity.INFO,
title="Agent 部署开始",
fields=[
("Agent", agent_id),
("目标版本", f"v{version}"),
("部署策略", strategy),
],
)
self._send_all_channels(message, NotificationSeverity.INFO)
def notify_deployment_success(self, agent_id: str, version: str, duration: float):
"""部署成功通知"""
message = self._format_message(
severity=NotificationSeverity.INFO,
title="Agent 部署成功",
fields=[
("Agent", agent_id),
("版本", f"v{version}"),
("耗时", f"{duration:.1f}s"),
],
)
self._send_all_channels(message, NotificationSeverity.INFO)
def notify_deployment_failure(self, agent_id: str, version: str, error: str):
"""部署失败通知"""
message = self._format_message(
severity=NotificationSeverity.ERROR,
title="Agent 部署失败",
fields=[
("Agent", agent_id),
("版本", f"v{version}"),
("错误", error),
],
)
self._send_all_channels(message, NotificationSeverity.ERROR)
def notify_rollback(self, agent_id: str, from_version: str, to_version: str, reason: str):
"""回滚通知"""
message = self._format_message(
severity=NotificationSeverity.WARNING,
title="Agent 回滚",
fields=[
("Agent", agent_id),
("从版本", f"v{from_version}"),
("到版本", f"v{to_version}"),
("原因", reason),
],
)
self._send_all_channels(message, NotificationSeverity.WARNING)
def notify_emergency(self, agent_id: str, issue: str, severity: NotificationSeverity):
"""紧急通知"""
message = self._format_message(
severity=severity,
title=f"{'紧急' if severity == NotificationSeverity.CRITICAL else ''}Agent 告警",
fields=[
("Agent", agent_id),
("问题", issue),
],
)
self._send_all_channels(message, severity)
def _send_all_channels(self, message: Dict[str, Any], severity: NotificationSeverity):
"""发送到所有已配置的通知渠道"""
if self.config.slack_webhook:
self._send_slack(message)
if self.config.telegram_bot_token and self.config.telegram_chat_id:
self._send_telegram(message)
if self.config.pagerduty_routing_key and severity in [NotificationSeverity.ERROR, NotificationSeverity.CRITICAL]:
self._send_pagerduty(message)
if self.config.webhook_url:
self._send_webhook(message)
def _format_message(self, severity: NotificationSeverity, title: str, fields: List[tuple]) -> Dict[str, Any]:
"""格式化通知消息"""
return {
"severity": severity.value,
"title": title,
"fields": [{"name": k, "value": v} for k, v in fields],
"timestamp": __import__("time").time(),
}
def _send_slack(self, message: Dict[str, Any]):
"""发送到 Slack"""
payload = {
"blocks": [
{
"type": "header",
"text": {"type": "plain_text", "text": message["title"]},
},
{
"type": "section",
"fields": [
{"type": "mrkdwn", "text": f"*{f['name']}:* {f['value']}"}
for f in message["fields"]
],
},
{
"type": "context",
"elements": [
{"type": "mrkdwn", "text": f"Severity: {message['severity']}"}
],
},
]
}
try:
requests.post(self.config.slack_webhook, json=payload, timeout=10)
except Exception as e:
logger.error(f"Slack 通知失败: {e}")
def _send_telegram(self, message: Dict[str, Any]):
"""发送到 Telegram"""
text = f"*{message['title']}*\n\n"
for f in message["fields"]:
text += f"*{f['name']}:* {f['value']}\n"
text += f"\nSeverity: {message['severity']}"
payload = {
"chat_id": self.config.telegram_chat_id,
"text": text,
"parse_mode": "Markdown",
}
url = f"https://api.telegram.org/bot{self.config.telegram_bot_token}/sendMessage"
try:
requests.post(url, json=payload, timeout=10)
except Exception as e:
logger.error(f"Telegram 通知失败: {e}")
def _send_pagerduty(self, message: Dict[str, Any]):
"""发送到 PagerDuty (仅 error/critical)"""
payload = {
"routing_key": self.config.pagerduty_routing_key,
"event_action": "trigger",
"payload": {
"summary": message["title"],
"severity": message["severity"],
"source": "agent-deploy-pipeline",
"custom_details": {f["name"]: f["value"] for f in message["fields"]},
},
}
try:
requests.post("https://events.pagerduty.com/v2/enqueue", json=payload, timeout=10)
except Exception as e:
logger.error(f"PagerDuty 通知失败: {e}")
def _send_webhook(self, message: Dict[str, Any]):
"""发送到自定义 Webhook"""
try:
requests.post(self.config.webhook_url, json=message, timeout=10)
except Exception as e:
logger.error(f"Webhook 通知失败: {e}")
8.4 Grafana Dashboard 配置
{
"dashboard": {
"title": "AI Agent 部署监控",
"panels": [
{
"title": "当前版本",
"type": "stat",
"targets": [{"expr": "agent_version{agent_id=\"$agent\"}", "legendFormat": "{{version}}"}]
},
{
"title": "部署耗时",
"type": "histogram",
"targets": [{"expr": "deployment_duration_bucket{agent_id=\"$agent\"}", "legendFormat": "{{le}}"}]
},
{
"title": "健康检查",
"type": "timeseries",
"targets": [{"expr": "health_check_passed{agent_id=\"$agent\"}", "legendFormat": "{{check_name}}"}]
},
{
"title": "错误率",
"type": "timeseries",
"targets": [{"expr": "error_rate{agent_id=\"$agent\"}", "legendFormat": "error_rate"}]
},
{
"title": "P99 延迟",
"type": "timeseries",
"targets": [{"expr": "latency{agent_id=\"$agent\",percentile=\"p99\"}", "legendFormat": "p99"}]
},
{
"title": "更新成功率",
"type": "stat",
"targets": [
{"expr": "rate(update_success{agent_id=\"$agent\"}[1h])", "legendFormat": "success"},
{"expr": "rate(update_failures{agent_id=\"$agent\"}[1h])", "legendFormat": "failures"}
]
},
{
"title": "回滚历史",
"type": "logs",
"targets": [{"expr": "rollback_count{agent_id=\"$agent\"}", "legendFormat": "rollbacks"}]
},
{
"title": "版本历史",
"type": "table",
"transformations": [
{"id": "seriesToRows", "options": {}},
{"id": "organize", "options": {"excludeByName": {"Time": true}, "indexByName": {}, "renameByName": {}}}
]
}
],
"templating": {
"list": [
{"name": "agent", "type": "query", "query": "label_values(agent_version, agent_id)"}
]
}
}
}
9. 附录
A. 环境变量参考
| 变量名 | 描述 | 默认值 |
|---|---|---|
MSG_CHAIN_RPC |
MSG Chain RPC 端点 | https://rpc.msgchain.zone |
MSG_CHAIN_WS |
WebSocket 端点 | wss://rpc.msgchain.zone/websocket |
REGISTRY_CONTRACT |
Agent 注册表合约地址 | msg1agentregistry... |
IPFS_GATEWAY |
IPFS 网关 | https://ipfs.msgchain.zone |
DEPLOYER_MNEMONIC |
部署账户助记词 | 无默认值(必须设置) |
AGENT_DATA_DIR |
Agent 数据目录 | /data/agent |
CHECK_INTERVAL |
更新检查间隔(秒) | 60 |
HEALTH_CHECK_TIMEOUT |
健康检查超时(秒) | 300 |
MAX_RETRIES |
最大重试次数 | 3 |
AUTO_UPDATE_MINOR |
自动更新 minor 版本 | true |
AUTO_UPDATE_PATCH |
自动更新 patch 版本 | true |
AUTO_UPDATE_MAJOR |
自动更新 major 版本 | false |
LOG_LEVEL |
日志级别 | INFO |
PROMETHEUS_PUSHGATEWAY |
Prometheus Pushgateway 地址 | 可选 |
SLACK_WEBHOOK |
Slack Webhook URL | 可选 |
TELEGRAM_BOT_TOKEN |
Telegram Bot Token | 可选 |
TELEGRAM_CHAT_ID |
Telegram Chat ID | 可选 |
VAULT_ADDR |
Vault 地址 | 可选 |
VAULT_TOKEN |
Vault Token | 可选 |
B. 错误码参考
| 错误码 | 描述 | 处理方式 |
|---|---|---|
E001 |
部署清单验证失败 | 检查清单格式和签名 |
E002 |
Wasm 代码上传失败 | 检查文件大小和网络连接 |
E003 |
合约实例化失败 | 检查初始化参数 |
E004 |
状态迁移失败 | 检查状态兼容性 |
E005 |
健康检查超时 | 检查新版本健康状况 |
E006 |
错误率超标 | 自动触发回滚 |
E007 |
延迟超标 | 自动触发回滚 |
E008 |
治理审批超时 | 升级到上级审批人 |
E009 |
资源配额不足 | 扩容或优化资源需求 |
E010 |
签名验证失败 | 检查签名密钥 |
C. CLI 命令参考
# Agent 部署相关命令
make test # 运行所有测试
make build # 构建 Wasm 合约和镜像
make deploy-canary # Canary 部署 (10%)
make deploy-full # 完全部署 (100%)
make rollback VERSION=1.0.0 # 回滚到指定版本
make verify # 验证部署状态
# 脚本工具
python scripts/verify_manifest.py <manifest.yaml> # 验证部署清单
python scripts/wait_for_health.py --agent-id <id> # 等待健康检查
python scripts/check_metrics.py --agent-id <id> # 检查部署指标
python scripts/register_version.py # 注册新版本到链上
python scripts/rollback.py --target-version <ver> # 执行回滚
python scripts/notify.py --channel <ch> --message <msg> # 发送通知
D. 故障排查指南
D.1 Agent 无法自更新
- 检查链上注册表连接
curl https://rpc.msgchain.zone/cosmos/cosmos-sdk/status - 检查注册表合约是否存在
python -c "from agent_registry import AgentRegistryClient; c = AgentRegistryClient(); print(c.get_agent('did:msg:agent:my-agent'))" - 检查 IPFS 网关可达性
curl -I https://ipfs.msgchain.zone/ipfs/QmTest...
D.2 合约迁移失败
- 检查 Wasm 代码大小是否超出限制
- 检查 migrate_msg 格式是否符合新合约期望
- 检查旧合约和新合约的状态模式是否兼容
- 检查 gas 限制是否足够
D.3 健康检查不通过
- 检查新版本的服务端口是否正确暴露
- 检查合约是否正确初始化
- 检查合约是否需要新的环境变量
- 检查新版本的日志
D.4 回滚失败
- 确认备份文件是否存在
- 检查回滚目标版本的合约地址是否仍有效
- 尝试使用治理命令强制回滚
- 在极少数情况下,可能需要链下数据恢复
E. 安全性最佳实践
- 密钥管理:所有部署密钥应存储在 Vault 或等效的密钥管理系统中,切勿硬编码
- 最小权限原则:部署账户应仅具有必要的合约执行权限
- 审批关卡:所有生产环境部署必须经过至少一级审批
- 签名验证:所有部署清单必须经过 Ed25519 签名验证
- 审计日志:所有部署操作应记录到链上审计日志
- 限速保护:部署 API 应设置速率限制,防止滥用
- 多签治理:重大版本升级需要治理多签批准
- 隔离环境:开发、测试、生产环境应完全隔离
F. 参考资源
- MSG Chain 官方文档: https://docs.msgchain.zone
- CosmWasm 合约开发指南: https://docs.cosmwasm.com
- Agent 注册表合约规范: https://docs.msgchain.zone/agent-registry
- CI/CD 配置文件参考: https://docs.msgchain.zone/cicd
- IPFS 文档: https://docs.ipfs.tech
本文档是 MSG Chain AI Agent 开发体系的一部分。
如有疑问请联系: dev-support@msgchain.zone
