MSG Chain Agent API 网关安全与运维指南
数据来源:MSG Chain 代码库核实
主网状态: No-Go — 当前 MSGChain 主网裁决为 No-Go,以下内容反映代码实际状态,不代表生产可用。
面向 AI Agent 开发者的安全接入手册
协议版本: Agent API v1 | 链 ID:msg-chain-1| Bech32 前缀:msg
状态: 部分实现 — 标记为 stub 的端点尚未投产
目录
1. 概述
1.1 Agent API 架构
MSG Chain Agent API 是一组专为 AI Agent 设计的 RESTful + WebSocket 接口,使 Agent 能够:
- 查询链上账户余额与交易状态
- 通过 MPC 门限签名安全管理钱包
- 创建与执行支付指令
- 订阅链上事件流
- 获取预言机价格数据
网关架构采用三层隔离模型:
AI Agent
│
▼
┌─────────────────────┐
│ API Gateway │ ← 速率限制、IP 白名单、DID 验证
│ (msgchain-proxy) │
└─────┬───────────────┘
│
┌─────▼───────────────┐
│ Auth Service │ ← API Key、DID、JWT 签发与校验
│ (agent-authd) │
└─────┬───────────────┘
│
┌─────▼───────────────┐
│ Chain Backend │ ← CometBFT 全节点、MPC 节点
│ (msgchaind) │
└─────────────────────┘
1.2 公开路由 vs 保护路由
| 类别 | 路由 | 认证要求 | 速率限制 |
|---|---|---|---|
| 公开 | GET /agent/v1/health |
无 | 100 req/s |
| 公开 | GET /agent/v1/account/{address} |
无 | 100 req/s |
| 公开 | GET /agent/v1/balance/{address} |
无 | 100 req/s |
| 公开 | GET /agent/v1/tx/{hash} |
无 | 100 req/s |
| 公开 | GET /agent/v1/block/{height} |
无 | 100 req/s |
| 公开 | GET /agent/v1/events (WebSocket) |
无 → Token | 100 msg/s |
| 公开 | GET /agent/v1/oracle/price/{pair} |
无 | 100 req/s |
| 保护 ⚠️ | POST /agent/v1/wallet/create |
API Key + DID | 20 req/s |
| 保护 ⚠️ | POST /agent/v1/mpc/wallet/{id}/sign |
API Key + DID | 20 req/s |
| 保护 ⚠️ | POST /agent/v1/payment/create |
API Key + DID | 20 req/s |
⚠️ 保护路由目前为 stub 实现 — 端点返回
{"code": 1, "message": "stub"}。功能逻辑尚未投产,但认证与速率限制中间件已就绪,可在沙盒环境中先行集成测试。
1.3 安全模型:纵深防御
Agent API 采用纵深防御策略,多道防线确保即使单层被突破,整体安全仍可维持:
Layer 1: 速率限制 ── 防止 DDoS 与滥用
Layer 2: IP 白名单 ── 限制网络来源
Layer 3: API Key ── 凭证基础认证
Layer 4: DID 签名验证 ── 抗量子身份认证
Layer 5: 请求签名+防重放 ── 防止篡改与重放
Layer 6: 审计日志 ── 事后追溯
Layer 7: Approval Gates ── 审批门控(高风险操作)
1.4 状态说明
本文档描述的某些功能处于"部分实现"状态:
- 已投产: 公开路由、速率限制、健康检查、WebSocket 事件流
- stub 端点:
wallet/create、mpc/wallet/{id}/sign、payment/create - 审批门 (Approval Gates): 框架已定义,但未接入 API 层
- 审核日志链上追踪: 可选组件,需部署独立合约
所有代码示例中的密钥和端点均为示例,不可用于生产。
2. API密钥管理
2.1 密钥生成
API Key 是 Agent 接入网关的基本凭证。每把密钥关联一组权限范围(scope),支持设置过期时间。
请求格式:
curl -X POST https://api.msgchain.org/agent/v1/auth/keys \
-H "X-API-Key: <master-key>" \
-H "Content-Type: application/json" \
-d '{
"label": "my-agent-key",
"permissions": ["query", "wallet:read"],
"expiry": "2027-01-01T00:00:00Z"
}'
成功响应:
{
"code": 0,
"message": "ok",
"data": {
"api_key": "msg_sk_abc123def456...",
"key_id": "key-001",
"label": "my-agent-key",
"permissions": ["query", "wallet:read"],
"expires_at": "2027-01-01T00:00:00Z",
"created_at": "2026-07-06T12:00:00Z"
}
}
参数说明:
| 参数 | 类型 | 必需 | 说明 |
|---|---|---|---|
label |
string | 是 | 密钥别名,用于管理识别 |
permissions |
string[] | 是 | 权限范围列表 |
expiry |
string (RFC3339) | 否 | 过期时间;不传则永不过期 |
2.2 权限模型
API Key 的权限范围(scope)决定了该密钥能访问的资源。权限采用分层命名空间:
query → 所有查询接口(account, balance, tx, block, oracle)
wallet:read → MPC 钱包信息查询
wallet:sign → MPC 钱包签名操作(含 wallet:read)
payment:create → 创建支付
payment:read → 查询支付状态
admin:keys → 管理 API Key(需 master key)
admin:webhook → 管理 Webhook 配置
权限继承关系:
admin:* → 包含所有子权限
wallet:* → 包含 wallet:read + wallet:sign
payment:* → 包含 payment:create + payment:read
最小权限原则: 每个 Agent 只应获得完成任务所需的最小权限集。
# 权限矩阵验证
PERMISSION_MATRIX = {
"query": ["GET /agent/v1/health", "GET /agent/v1/account/*",
"GET /agent/v1/balance/*", "GET /agent/v1/tx/*",
"GET /agent/v1/block/*", "GET /agent/v1/oracle/price/*"],
"wallet:read": ["GET /agent/v1/wallet/*"],
"wallet:sign": ["GET /agent/v1/wallet/*", "POST /agent/v1/mpc/wallet/*/sign"],
"payment:create": ["POST /agent/v1/payment/create"],
"payment:read": ["GET /agent/v1/payment/*"],
"admin:keys": ["GET /agent/v1/auth/keys", "POST /agent/v1/auth/keys",
"DELETE /agent/v1/auth/keys/*"],
}
2.3 Python API Key 管理客户端
"""
agent_key_manager.py — MSG Chain Agent API Key 管理
依赖: pip install requests
"""
import time
import json
import hashlib
import hmac
from datetime import datetime, timezone
from typing import Optional
import requests
class AgentKeyManager:
"""API Key 生命周期管理客户端"""
BASE_URL = "https://api.msgchain.org/agent/v1"
def __init__(self, master_key: str):
self.master_key = master_key
self.session = requests.Session()
self.session.headers.update({
"X-API-Key": master_key,
"Content-Type": "application/json",
})
def create_key(self, label: str, permissions: list[str],
expiry: Optional[str] = None) -> dict:
payload = {
"label": label,
"permissions": permissions,
}
if expiry:
payload["expiry"] = expiry
resp = self.session.post(f"{self.BASE_URL}/auth/keys", json=payload)
resp.raise_for_status()
return resp.json()
def list_keys(self) -> dict:
resp = self.session.get(f"{self.BASE_URL}/auth/keys")
resp.raise_for_status()
return resp.json()
def revoke_key(self, key_id: str) -> dict:
resp = self.session.delete(f"{self.BASE_URL}/auth/keys/{key_id}")
resp.raise_for_status()
return resp.json()
def rotate_key(self, key_id: str, label: str,
permissions: list[str]) -> tuple[dict, dict]:
old_key = self.revoke_key(key_id)
new_key = self.create_key(
label=label,
permissions=permissions,
expiry=old_key.get("data", {}).get("expires_at"),
)
return old_key, new_key
class AgentKeyVault:
"""本地密钥安全存储(加密)"""
def __init__(self, vault_path: str = ".agent_vault.json"):
self.vault_path = vault_path
self._keys: dict[str, dict] = {}
def store(self, agent_id: str, api_key: str, permissions: list[str]):
self._keys[agent_id] = {
"api_key": self._encrypt(api_key),
"permissions": permissions,
"created_at": datetime.now(timezone.utc).isoformat(),
}
self._persist()
def retrieve(self, agent_id: str) -> Optional[str]:
entry = self._keys.get(agent_id)
if entry:
return self._decrypt(entry["api_key"])
return None
def _encrypt(self, plaintext: str) -> str:
# 生产环境应使用 KMS 或硬件加密模块
return plaintext # 示例用,不可用于生产
def _decrypt(self, ciphertext: str) -> str:
return ciphertext
def _persist(self):
with open(self.vault_path, "w") as f:
json.dump(self._keys, f, indent=2)
# 使用示例
if __name__ == "__main__":
manager = AgentKeyManager(master_key="<your-master-key>")
# 创建只读查询密钥
resp = manager.create_key(
label="price-bot-001",
permissions=["query"],
expiry="2027-06-01T00:00:00Z",
)
print(f"Created key: {resp['data']['key_id']}")
# 创建钱包签名密钥
resp = manager.create_key(
label="trading-agent-001",
permissions=["wallet:sign", "query"],
)
print(f"Created key: {resp['data']['key_id']}")
# 轮换密钥
old, new = manager.rotate_key(
key_id="key-001",
label="price-bot-002",
permissions=["query"],
)
print(f"Rotated: {old['data']['key_id']} → {new['data']['key_id']}")
2.4 TypeScript API Key 管理客户端
// agent_key_manager.ts — MSG Chain Agent API Key 管理
// 依赖: npm install axios
import axios, { AxiosInstance } from 'axios';
interface CreateKeyRequest {
label: string;
permissions: string[];
expiry?: string;
}
interface ApiKeyData {
api_key: string;
key_id: string;
label: string;
permissions: string[];
expires_at: string;
created_at: string;
}
interface ApiResponse<T> {
code: number;
message: string;
data: T;
}
class AgentKeyManager {
private client: AxiosInstance;
private readonly BASE_URL = 'https://api.msgchain.org/agent/v1';
constructor(masterKey: string) {
this.client = axios.create({
baseURL: this.BASE_URL,
headers: {
'X-API-Key': masterKey,
'Content-Type': 'application/json',
},
});
}
async createKey(
label: string,
permissions: string[],
expiry?: string
): Promise<ApiResponse<ApiKeyData>> {
const payload: CreateKeyRequest = { label, permissions };
if (expiry) payload.expiry = expiry;
const { data } = await this.client.post('/auth/keys', payload);
return data;
}
async listKeys(): Promise<ApiResponse<ApiKeyData[]>> {
const { data } = await this.client.get('/auth/keys');
return data;
}
async revokeKey(keyId: string): Promise<ApiResponse<null>> {
const { data } = await this.client.delete(`/auth/keys/${keyId}`);
return data;
}
async rotateKey(
keyId: string,
label: string,
permissions: string[]
): Promise<{ old: ApiResponse<null>; new: ApiResponse<ApiKeyData> }> {
const oldKey = await this.revokeKey(keyId);
const newKey = await this.createKey(label, permissions);
return { old: oldKey, new: newKey };
}
}
// 密钥安全存储
class AgentKeyVault {
private keys: Map<string, string> = new Map();
store(agentId: string, apiKey: string): void {
// 生产环境应使用加密存储 (如 AWS KMS, HashiCorp Vault)
this.keys.set(agentId, apiKey);
}
retrieve(agentId: string): string | undefined {
return this.keys.get(agentId);
}
revoke(agentId: string): boolean {
return this.keys.delete(agentId);
}
}
// 使用示例
async function main() {
const manager = new AgentKeyManager('<your-master-key>');
// 创建 Agent 专用密钥
const resp = await manager.createKey(
'nft-bot-001',
['query', 'wallet:read'],
'2027-01-01T00:00:00Z'
);
console.log(`Created key: ${resp.data.key_id}`);
// 保存密钥到 vault
const vault = new AgentKeyVault();
vault.store('nft-bot-001', resp.data.api_key);
}
main().catch(console.error);
2.5 多密钥策略
生产环境推荐每个 Agent 使用独立 API Key,以达成:
- 隔离性: 单 Agent 密钥泄漏不影响其他服务
- 审计性: 日志可直接追溯到具体 Agent
- 细粒度撤销: 可单独吊销某个 Agent 的访问权限
Agent A (price-bot) → API Key A (scope: query)
Agent B (trade-bot) → API Key B (scope: query, wallet:sign)
Agent C (nft-minter) → API Key C (scope: query, wallet:sign)
Agent D (monitor) → API Key D (scope: query, admin:webhook)
2.6 密钥轮换策略
"""
key_rotation_scheduler.py — 自动密钥轮换
通过 cron 或定时任务调用
"""
import os
import json
import logging
from datetime import datetime, timedelta, timezone
logging.basicConfig(level=logging.INFO)
logger = logging.getLogger(__name__)
class KeyRotationScheduler:
def __init__(self, master_key: str, grace_days: int = 7):
from agent_key_manager import AgentKeyManager
self.manager = AgentKeyManager(master_key)
self.grace_days = grace_days
def check_and_rotate(self):
"""检查即将过期的密钥并轮换"""
resp = self.manager.list_keys()
if resp["code"] != 0:
logger.error("Failed to list keys")
return
for key in resp["data"]:
expires_at = datetime.fromisoformat(key["expires_at"])
remaining = (expires_at - datetime.now(timezone.utc)).days
if remaining <= self.grace_days:
logger.info(
f"Rotating key {key['key_id']} "
f"(expires in {remaining}d)"
)
try:
old, new = self.manager.rotate_key(
key_id=key["key_id"],
label=f"{key['label']}-rotated-{datetime.now().strftime('%Y%m%d')}",
permissions=key["permissions"],
)
logger.info(
f"Rotated {old['data']['key_id']} → {new['data']['key_id']}"
)
except Exception as e:
logger.error(f"Rotation failed for {key['key_id']}: {e}")
if __name__ == "__main__":
scheduler = KeyRotationScheduler(
master_key=os.environ.get("MSG_MASTER_KEY", ""),
grace_days=7,
)
scheduler.check_and_rotate()
2.7 密钥撤销
撤销密钥后,该密钥立即失效。所有使用该密钥的未完成请求将被拒绝,返回 {"code": 17, "message": "unauthorized"}。
curl -X DELETE https://api.msgchain.org/agent/v1/auth/keys/key-001 \
-H "X-API-Key: <master-key>"
撤销后建议操作:
- 确认 Agent 已切换至新密钥
- 监控旧密钥是否有异常访问(理论上不应有)
- 如怀疑泄漏,立即触发安全事件响应
3. DID身份认证
3.1 认证流程概述
DID(Decentralized Identifier)认证是 Agent API 的第二层(也是核心)身份验证机制。它基于 Dilithium-5 后量子签名算法,提供抗量子攻击的身份验证。
完整的 DOI 认证握手流程:
Agent Gateway
│ │
│ POST /auth/challenge │
│ { "did": "did:msg:agent:..." } │
│ ──────────────────────────────► │
│ │
│ { "challenge": "random_hex..." } │
│ ◄────────────────────────────── │
│ │
│ (Agent 用 Dilithium-5 签名 │
│ 签名 challenge) │
│ │
│ POST /auth/verify │
│ { "did": "...", │
│ "signature": "dilithium5:..." }│
│ ──────────────────────────────► │
│ │
│ { "token": "jwt...", │
│ "expires_in": 3600 } │
│ ◄────────────────────────────── │
│ │
│ (后续请求附加 Bearer Token) │
│ ──────────────────────────────► │
│ │
3.2 Python DID 认证实现
"""
agent_did_auth.py — Dilithium-5 DID 认证流
依赖: pip install requests pydilithium (或自编译 libdilithium)
"""
import json
import time
import hashlib
import os
from typing import Optional
import requests
class Dilithium5Signer:
"""Dilithium-5 签名封装(需底层库支持)"""
def __init__(self, private_key_hex: str):
self.private_key = bytes.fromhex(private_key_hex)
def sign(self, message: bytes) -> bytes:
"""
使用 Dilithium-5 签名消息
生产环境应调用硬件安全模块 (HSM) 或 TEE 签名
"""
# 伪代码 — 替换为实际 Dilithium-5 库调用
# from dilithium import Dilithium5
# sk = Dilithium5.SecretKey(self.private_key)
# return sk.sign(message)
raise NotImplementedError(
"Install dilithium5 library: "
"https://github.com/pq-crystals/dilithium"
)
@staticmethod
def generate_keypair() -> tuple[str, str]:
"""
生成 Dilithium-5 密钥对
返回 (private_key_hex, public_key_hex)
"""
# from dilithium import Dilithium5
# sk, pk = Dilithium5.keygen()
# return sk.hex(), pk.hex()
raise NotImplementedError(
"Install dilithium5 library"
)
class DIDAuthClient:
"""DID 认证客户端"""
BASE_URL = "https://api.msgchain.org/agent/v1"
def __init__(self, did: str, private_key_hex: str):
self.did = did
self.signer = Dilithium5Signer(private_key_hex)
self.token: Optional[str] = None
self.token_expiry: int = 0
def authenticate(self) -> str:
"""执行完整的 DID 认证握手,返回 JWT token"""
# Step 1: 请求 challenge
challenge_resp = self._request_challenge()
challenge = challenge_resp["data"]["challenge"]
# Step 2: 签名 challenge
signature = self._sign_challenge(challenge)
# Step 3: 验证签名,获取 token
verify_resp = self._verify_signature(signature)
self.token = verify_resp["data"]["token"]
self.token_expiry = time.time() + verify_resp["data"]["expires_in"]
return self.token
def get_valid_token(self) -> str:
"""获取有效 token,必要时自动刷新"""
if not self.token or time.time() >= self.token_expiry - 60:
return self.authenticate()
return self.token
def _request_challenge(self) -> dict:
resp = requests.post(
f"{self.BASE_URL}/auth/challenge",
json={"did": self.did},
headers={"Content-Type": "application/json"},
)
resp.raise_for_status()
return resp.json()
def _sign_challenge(self, challenge: str) -> str:
message = challenge.encode("utf-8")
sig_bytes = self.signer.sign(message)
return "dilithium5:" + sig_bytes.hex()
def _verify_signature(self, signature: str) -> dict:
resp = requests.post(
f"{self.BASE_URL}/auth/verify",
json={
"did": self.did,
"signature": signature,
},
headers={"Content-Type": "application/json"},
)
resp.raise_for_status()
return resp.json()
def make_authenticated_request(
self, method: str, path: str,
body: Optional[dict] = None
) -> requests.Response:
token = self.get_valid_token()
headers = {
"Authorization": f"Bearer {token}",
"Content-Type": "application/json",
}
url = f"{self.BASE_URL}{path}"
if method == "GET":
return requests.get(url, headers=headers)
elif method == "POST":
return requests.post(url, headers=headers, json=body)
else:
raise ValueError(f"Unsupported method: {method}")
# DID 注册助手
def register_agent_did(
public_key_hex: str,
api_key: str,
metadata: Optional[dict] = None,
) -> str:
"""
在链上注册 Agent 的 DID 文档
返回 DID 标识符
注意: 此端点为 stub,测试用
"""
url = "https://api.msgchain.org/agent/v1/did/register"
payload = {
"public_key": public_key_hex,
"key_type": "dilithium-5",
"metadata": metadata or {},
}
resp = requests.post(
url,
json=payload,
headers={
"X-API-Key": api_key,
"Content-Type": "application/json",
},
)
resp.raise_for_status()
result = resp.json()
return result["data"]["did"]
# 使用示例
if __name__ == "__main__":
# 生成密钥对(伪代码 — 需要 dilithium 库)
# sk_hex, pk_hex = Dilithium5Signer.generate_keypair()
# 注册 DID(stub)
# did = register_agent_did(pk_hex, api_key="<api-key>")
# 认证
# client = DIDAuthClient(did="did:msg:agent:my-agent", private_key_hex=sk_hex)
# token = client.authenticate()
# print(f"Got token: {token[:32]}...")
pass
3.3 TypeScript DID 认证实现
// agent_did_auth.ts — Dilithium-5 DID Auth Client
// 依赖: npm install axios
import axios, { AxiosInstance } from 'axios';
interface ChallengeResponse {
code: number;
message: string;
data: {
challenge: string;
expires_at: string;
};
}
interface VerifyResponse {
code: number;
message: string;
data: {
token: string;
expires_in: number;
token_type: string;
};
}
/**
* Dilithium-5 签名接口
* 生产环境应通过 WASM 或 Node.js 原生扩展调用
*/
interface Dilithium5SignerInterface {
sign(message: Uint8Array): Uint8Array;
}
/**
* 基于 Web Crypto API 的占位实现
* 实际需要集成 dilithium5 WASM 模块
*/
class Dilithium5SignerStub implements Dilithium5SignerInterface {
constructor(private privateKey: Uint8Array) {}
sign(message: Uint8Array): Uint8Array {
// 占位 — 替换为实际 Dilithium-5 实现
// const sig = dilithium5.sign(this.privateKey, message);
// return sig;
throw new Error(
'Dilithium-5 library required. Install @pqcrystals/dilithium5-wasm'
);
}
}
class DIDAuthClient {
private client: AxiosInstance;
private token: string | null = null;
private tokenExpiry: number = 0;
private readonly BASE_URL = 'https://api.msgchain.org/agent/v1';
constructor(
private did: string,
private signer: Dilithium5SignerInterface
) {
this.client = axios.create({
baseURL: this.BASE_URL,
headers: { 'Content-Type': 'application/json' },
});
}
async authenticate(): Promise<string> {
// Step 1: 获取 Challenge
const { data: challengeData } = await this.client.post<ChallengeResponse>(
'/auth/challenge',
{ did: this.did }
);
const challenge = challengeData.data.challenge;
// Step 2: 签名
const signature = this.signChallenge(challenge);
// Step 3: 验证并获取 Token
const { data: verifyData } = await this.client.post<VerifyResponse>(
'/auth/verify',
{ did: this.did, signature }
);
this.token = verifyData.data.token;
this.tokenExpiry = Date.now() + verifyData.data.expires_in * 1000;
return this.token!;
}
async getValidToken(): Promise<string> {
if (!this.token || Date.now() >= this.tokenExpiry - 60000) {
return this.authenticate();
}
return this.token;
}
private signChallenge(challenge: string): string {
const message = new TextEncoder().encode(challenge);
const sigBytes = this.signer.sign(message);
return 'dilithium5:' + Buffer.from(sigBytes).toString('hex');
}
async request(method: string, path: string, body?: unknown) {
const token = await this.getValidToken();
const { data } = await this.client.request({
method,
url: path,
data: body,
headers: { Authorization: `Bearer ${token}` },
});
return data;
}
}
// 使用示例
async function main() {
// 初始化签名器(需实际库)
// const signer = new Dilithium5SignerStub(privateKeyBytes);
// const auth = new DIDAuthClient('did:msg:agent:my-agent', signer);
// await auth.authenticate();
// 发起认证请求
// const result = await auth.request('POST', '/agent/v1/payment/create', {
// from: 'msg1...',
// to: 'msg1...',
// amount: '1000000',
// });
// console.log(result);
}
main().catch(console.error);
3.4 Token 刷新机制
JWT Token 默认有效期为 1 小时。Agent 应在 Token 过期前主动刷新,避免请求中断。
class TokenRefresher:
"""自动 Token 刷新守护"""
def __init__(self, auth_client: DIDAuthClient,
refresh_interval: int = 300):
self.client = auth_client
self.refresh_interval = refresh_interval
self._running = False
def start(self):
"""启动刷新循环(建议在独立线程运行)"""
import threading
self._running = True
thread = threading.Thread(target=self._loop, daemon=True)
thread.start()
def stop(self):
self._running = False
def _loop(self):
while self._running:
try:
self.client.get_valid_token()
print(f"[TokenRefresher] Token refreshed at {time.time()}")
except Exception as e:
print(f"[TokenRefresher] Refresh failed: {e}")
time.sleep(self.refresh_interval)
3.5 Token 验证中间件
以下代码展示网关侧如何验证 Agent 提交的 JWT Token(仅供参考,生产环境由网关处理):
"""
gateway_middleware.py — 网关 JWT 验证中间件(参考实现)
"""
import json
import time
from typing import Callable
from jwcrypto import jwk, jwt # pip install jwcrypto
class JWTValidator:
"""JWT Token 验证器"""
def __init__(self, public_key_pem: str):
self.public_key = jwk.JWK.from_pem(public_key_pem.encode())
def validate(self, token: str) -> dict:
"""验证 JWT 并返回 payload"""
try:
decoded = jwt.JWT(key=self.public_key, jwt=token)
payload = json.loads(decoded.claims)
# 检查过期
if payload.get("exp", 0) < time.time():
raise ValueError("Token expired")
# 检查 nbf (not before)
if payload.get("nbf", 0) > time.time():
raise ValueError("Token not yet valid")
return payload
except Exception as e:
raise ValueError(f"Invalid token: {e}")
# DID 授权装饰器
def require_did_auth(validator: JWTValidator):
"""DID 认证中间件装饰器"""
def decorator(handler: Callable):
def wrapper(request, *args, **kwargs):
auth_header = request.headers.get("Authorization", "")
if not auth_header.startswith("Bearer "):
return {"code": 3, "message": "unauthorized"}, 401
token = auth_header[7:]
try:
payload = validator.validate(token)
request.agent_did = payload.get("sub")
request.agent_scopes = payload.get("scopes", [])
return handler(request, *args, **kwargs)
except ValueError as e:
return {"code": 3, "message": str(e)}, 401
return wrapper
return decorator
3.6 DID 文档注册
Agent 的 DID 文档(DID Document)注册在 MSG Chain 上,包含公钥信息和服务端点。完整 DID 文档格式:
{
"@context": ["https://www.w3.org/ns/did/v1"],
"id": "did:msg:agent:my-agent",
"verificationMethod": [{
"id": "did:msg:agent:my-agent#key-1",
"type": "Dilithium5VerificationKey2026",
"controller": "did:msg:agent:my-agent",
"publicKeyMultibase": "z6Mkf5rGp..."
}],
"authentication": ["did:msg:agent:my-agent#key-1"],
"assertionMethod": ["did:msg:agent:my-agent#key-1"],
"service": [{
"id": "did:msg:agent:my-agent#agent-endpoint",
"type": "AgentService",
"serviceEndpoint": "https://agent.msgchain.org/callback"
}],
"created": "2026-07-06T12:00:00Z",
"updated": "2026-07-06T12:00:00Z"
}
4. 速率限制策略
4.1 速率限制层级
Agent API 采用多层速率限制架构,保护后端免受滥用:
| 层级 | 认证级别 | 默认限制 | Burst | 适用场景 |
|---|---|---|---|---|
| Tier 1 | 无认证(公开) | 10 req/s | 20 | 健康检查、低频率查询 |
| Tier 2 | API Key | 100 req/s | 200 | 标准 Agent 查询 |
| Tier 3 | DID Auth + API Key | 500 req/s | 1000 | 高吞吐 Agent |
| Tier 4 | 白名单 | 自定义 | 自定义 | 企业级部署 |
4.2 速率限制响应头
每个请求的响应包含以下限流信息:
X-RateLimit-Limit: 100 # 窗口内最大请求数
X-RateLimit-Remaining: 87 # 当前窗口剩余请求数
X-RateLimit-Reset: 1720252800 # 窗口重置时间(Unix 时间戳)
4.3 Python 限流感知客户端
"""
rate_limited_client.py — 带速率限制感知的 HTTP 客户端
包含指数退避 + 抖动重试
"""
import time
import random
import requests
from typing import Optional, Callable
from dataclasses import dataclass
from urllib.parse import urljoin
@dataclass
class RateLimitState:
"""跟踪速率限制状态"""
limit: int = 0
remaining: int = 0
reset_time: float = 0
class RateLimitedClient:
"""
MSG Chain Agent API 限流感知客户端
支持自动重试、指数退避、抖动
"""
BASE_URL = "https://api.msgchain.org/agent/v1"
def __init__(
self,
api_key: Optional[str] = None,
bearer_token: Optional[str] = None,
max_retries: int = 5,
base_delay: float = 1.0,
max_delay: float = 60.0,
):
self.session = requests.Session()
self.max_retries = max_retries
self.base_delay = base_delay
self.max_delay = max_delay
self.rate_limit = RateLimitState()
# 认证头
if api_key:
self.session.headers["X-API-Key"] = api_key
if bearer_token:
self.session.headers["Authorization"] = f"Bearer {bearer_token}"
self.session.headers["Content-Type"] = "application/json"
self.session.headers["User-Agent"] = "MSGChain-Agent/1.0"
def request(self, method: str, path: str,
body: Optional[dict] = None,
retry_on_limit: bool = True) -> dict:
url = urljoin(self.BASE_URL, path)
last_error: Optional[Exception] = None
for attempt in range(self.max_retries + 1):
try:
response = self.session.request(
method=method,
url=url,
json=body,
)
# 更新限流状态
self._update_rate_limit(response)
# 成功
if response.ok:
return response.json()
# 速率限制 (code 17)
if response.status_code == 429 and retry_on_limit:
self._handle_rate_limit(response, attempt)
continue
# 其他错误
response.raise_for_status()
except requests.exceptions.RequestException as e:
last_error = e
if attempt < self.max_retries:
self._backoff(attempt)
raise RuntimeError(
f"Request failed after {self.max_retries} retries: {last_error}"
)
def _update_rate_limit(self, response: requests.Response):
self.rate_limit.limit = int(
response.headers.get("X-RateLimit-Limit", 0)
)
self.rate_limit.remaining = int(
response.headers.get("X-RateLimit-Remaining", 0)
)
reset_str = response.headers.get("X-RateLimit-Reset", "0")
self.rate_limit.reset_time = float(reset_str)
def _handle_rate_limit(self, response: requests.Response, attempt: int):
"""处理 429 限流响应"""
reset_ts = float(
response.headers.get("X-RateLimit-Reset", time.time() + 1)
)
wait_time = max(reset_ts - time.time(), 0.1)
# 如果服务器给了 Retry-After,优先使用
retry_after = response.headers.get("Retry-After")
if retry_after:
wait_time = max(float(retry_after), 0.1)
print(
f"[RateLimit] Attempt {attempt + 1}: "
f"rate limited, waiting {wait_time:.1f}s"
)
time.sleep(wait_time)
def _backoff(self, attempt: int):
"""
指数退避 + 抖动
delay = min(base * 2^attempt + jitter, max_delay)
"""
delay = min(
self.base_delay * (2 ** attempt) + random.uniform(0, 1),
self.max_delay,
)
print(f"[Backoff] Attempt {attempt + 1}: waiting {delay:.1f}s")
time.sleep(delay)
# 便捷方法
def get(self, path: str, **kwargs) -> dict:
return self.request("GET", path, **kwargs)
def post(self, path: str, body: dict, **kwargs) -> dict:
return self.request("POST", path, body=body, **kwargs)
class RateLimitAwareAgent:
"""
限流感知 Agent 基类
在发送请求前主动检查剩余配额
"""
def __init__(self, client: RateLimitedClient):
self.client = client
def check_capacity(self, required: int = 1) -> bool:
"""
检查是否有足够配额
如果剩余配额不足,等待窗口重置
"""
rl = self.client.rate_limit
if rl.remaining < required:
wait_time = max(rl.reset_time - time.time(), 0)
if wait_time > 0:
print(
f"[Capacity] Only {rl.remaining}/{rl.limit} remaining. "
f"Waiting {wait_time:.1f}s for reset."
)
time.sleep(wait_time)
return False
return True
def get_balance(self, address: str) -> dict:
self.check_capacity()
return self.client.get(f"/balance/{address}")
def get_account(self, address: str) -> dict:
self.check_capacity()
return self.client.get(f"/account/{address}")
# 使用示例
if __name__ == "__main__":
client = RateLimitedClient(
api_key="msg_sk_abc123def456",
max_retries=3,
)
# 正常请求
result = client.get("/balance/msg1q6w9e8r7t5y4u3i2o1p")
print(f"Balance: {result}")
# 限流重试测试
for _ in range(200):
try:
result = client.get("/health")
print(f"Health: {result}, Remaining: {client.rate_limit.remaining}")
except Exception as e:
print(f"Error: {e}")
break
4.4 TypeScript 限流感知客户端
// rate_limited_client.ts — 限流感知 HTTP 客户端
import axios, { AxiosInstance, AxiosResponse, AxiosError } from 'axios';
interface RateLimitHeaders {
limit: number;
remaining: number;
reset: number;
}
class RateLimitedClient {
private client: AxiosInstance;
private rateLimit: RateLimitHeaders = { limit: 0, remaining: 0, reset: 0 };
private readonly BASE_URL = 'https://api.msgchain.org/agent/v1';
constructor(
private apiKey?: string,
private bearerToken?: string,
private maxRetries: number = 5,
private baseDelay: number = 1000,
private maxDelay: number = 60000
) {
this.client = axios.create({
baseURL: this.BASE_URL,
headers: {
'Content-Type': 'application/json',
'User-Agent': 'MSGChain-Agent/1.0',
...(apiKey ? { 'X-API-Key': apiKey } : {}),
...(bearerToken ? { Authorization: `Bearer ${bearerToken}` } : {}),
},
});
this.client.interceptors.response.use(
(response) => {
this.updateRateLimit(response);
return response;
},
(error: AxiosError) => {
if (error.response) {
this.updateRateLimit(error.response);
}
return Promise.reject(error);
}
);
}
private updateRateLimit(response: AxiosResponse): void {
const headers = response.headers;
this.rateLimit = {
limit: parseInt(headers['x-ratelimit-limit'] || '0', 10),
remaining: parseInt(headers['x-ratelimit-remaining'] || '0', 10),
reset: parseInt(headers['x-ratelimit-reset'] || '0', 10),
};
}
get rateLimitState(): Readonly<RateLimitHeaders> {
return { ...this.rateLimit };
}
async request(method: string, path: string, data?: unknown): Promise<unknown> {
const url = path.startsWith('http') ? path : path;
for (let attempt = 0; attempt <= this.maxRetries; attempt++) {
try {
const response = await this.client.request({
method,
url,
data,
validateStatus: (status) => status < 500,
});
if (response.status === 429) {
await this.handleRateLimit(response, attempt);
continue;
}
if (response.status >= 400) {
throw new Error(
`Request failed: ${response.status} ${JSON.stringify(response.data)}`
);
}
return response.data;
} catch (error) {
if (attempt < this.maxRetries && !(error instanceof Error &&
error.message.includes('429'))) {
await this.backoff(attempt);
continue;
}
throw error;
}
}
throw new Error('Max retries exceeded');
}
private async handleRateLimit(response: AxiosResponse, attempt: number): Promise<void> {
const reset = parseInt(
response.headers['x-ratelimit-reset'] || '0', 10
);
const retryAfter = parseInt(
response.headers['retry-after'] || '0', 10
);
const waitTime = Math.max(
retryAfter || (reset - Date.now() / 1000),
0.1
);
console.log(
`[RateLimit] Attempt ${attempt + 1}: waiting ${waitTime.toFixed(1)}s`
);
await this.sleep(waitTime * 1000);
}
private async backoff(attempt: number): Promise<void> {
const delay = Math.min(
this.baseDelay * Math.pow(2, attempt) + Math.random() * 1000,
this.maxDelay
);
console.log(`[Backoff] Attempt ${attempt + 1}: waiting ${(delay / 1000).toFixed(1)}s`);
await this.sleep(delay);
}
private sleep(ms: number): Promise<void> {
return new Promise((resolve) => setTimeout(resolve, ms));
}
async get(path: string): Promise<unknown> {
return this.request('GET', path);
}
async post(path: string, data: unknown): Promise<unknown> {
return this.request('POST', path, data);
}
}
// 使用示例
async function main() {
const client = new RateLimitedClient('msg_sk_abc123def456');
// 批量查询
const addresses = [
'msg1q6w9e8r7t5y4u3i2o1p',
'msg1a2b3c4d5e6f7g8h9i0j',
];
for (const addr of addresses) {
const result = await client.get(`/balance/${addr}`);
console.log(`Balance [${addr}]:`, result);
console.log(`Rate limit: ${JSON.stringify(client.rateLimitState)}`);
}
}
main().catch(console.error);
4.5 退避策略详解
指数退避 + 抖动是处理限流的标准策略。以下是各退避轮次的等待时间示例:
def calculate_backoff(attempt: int, base_delay: float = 1.0,
max_delay: float = 60.0,
jitter: bool = True) -> float:
"""
计算退避时间
参数:
attempt: 当前重试次数(从 0 开始)
base_delay: 基础延迟(秒)
max_delay: 最大延迟(秒)
jitter: 是否添加随机抖动
返回:
等待时间(秒)
"""
import random
delay = base_delay * (2 ** attempt)
if jitter:
delay += random.uniform(0, delay * 0.5)
return min(delay, max_delay)
# 退避示例
for i in range(8):
delay = calculate_backoff(i)
print(f"Attempt {i}: wait {delay:.2f}s")
Attempt 0: wait 1.12s
Attempt 1: wait 2.34s
Attempt 2: wait 4.56s
Attempt 3: wait 8.12s
Attempt 4: wait 16.78s
Attempt 5: wait 32.45s
Attempt 6: wait 60.00s (上限)
Attempt 7: wait 60.00s (上限)
4.6 限流错误处理
速率限制错误的响应格式:
{
"code": 17,
"message": "rate_limit_exceeded",
"data": {
"limit": 100,
"remaining": 0,
"reset": 1720252800,
"retry_after": 42
}
}
5. IP白名单与地理限制
5.1 IP 白名单配置
为保护保护路由(签名、支付),建议配置 IP 白名单限制来源地址。
通过 API 配置白名单:
curl -X POST https://api.msgchain.org/agent/v1/admin/whitelist \
-H "X-API-Key: <master-key>" \
-H "Content-Type: application/json" \
-d '{
"cidrs": [
"203.0.113.0/24",
"198.51.100.0/28",
"192.0.2.5/32"
],
"description": "Production agents subnet"
}'
响应:
{
"code": 0,
"message": "ok",
"data": {
"whitelist_id": "wl-001",
"entries": 3,
"updated_at": "2026-07-06T12:00:00Z"
}
}
查询白名单:
curl -X GET https://api.msgchain.org/agent/v1/admin/whitelist \
-H "X-API-Key: <master-key>"
5.2 地理限制
可选启用地理限制,只允许特定国家/地区的 IP 访问:
curl -X PUT https://api.msgchain.org/agent/v1/admin/geofence \
-H "X-API-Key: <master-key>" \
-d '{
"allowed_countries": ["US", "SG", "DE", "JP"],
"block_unknown": true,
"enable_vpn_detection": true
}'
5.3 VPN/代理检测
网关集成了 VPN 和代理检测(基于 IP 信誉数据库)。检测到的代理请求将被标记并可选拦截:
{
"code": 3,
"message": "unauthorized",
"data": {
"reason": "vpn_or_proxy_detected",
"ip": "198.51.100.1",
"risk_score": 0.85
}
}
5.4 内部子网访问
对于运行在 MSG Chain 基础设施内的 Agent(例如在同一 Kubernetes 集群),可通过内部 DNS 和网络策略访问保护路由,绕过外部 IP 限制:
Internal: http://agent-api.internal.msgchain:8080/agent/v1/...
External: https://api.msgchain.org/agent/v1/...
内部子网 CIDR 范围:
10.0.0.0/8 # Kubernetes Pods
172.16.0.0/12 # 内部服务网段
192.168.0.0/16 # 管理网段
5.5 Python IP 白名单验证
"""
ip_whitelist.py — 客户端侧 IP 验证(仅供测试使用)
生产环境白名单验证由网关自动执行
"""
import ipaddress
from typing import Optional
class IPWhitelistValidator:
"""IP 白名单验证器"""
def __init__(self, allowed_cidrs: list[str]):
self.networks = [
ipaddress.ip_network(cidr, strict=False)
for cidr in allowed_cidrs
]
def is_allowed(self, ip_str: str) -> bool:
addr = ipaddress.ip_address(ip_str)
return any(addr in net for net in self.networks)
# 使用示例
validator = IPWhitelistValidator([
"203.0.113.0/24",
"198.51.100.0/28",
])
test_ips = ["203.0.113.42", "8.8.8.8", "198.51.100.15"]
for ip in test_ips:
allowed = validator.is_allowed(ip)
print(f"{ip}: {'✓ allowed' if allowed else '✗ blocked'}")
输出:
203.0.113.42: ✓ allowed
8.8.8.8: ✗ blocked
198.51.100.15: ✓ allowed
5.6 推荐架构
┌─────────────────────────────────────┐
│ Internet │
│ │
│ Agent A 203.0.113.5 ────┐ │
│ Agent B 198.51.100.7 ────┤ │
│ Attacker 8.8.8.8 ───────┤ │
│ ▼ │
│ API Gateway │
│ ┌──────────┐ │
│ │ IP Check │ ←──┤ 8.8.8.8 → 403
│ │ Geofence │ │
│ │ VPN Detect│ │
│ └──────────┘ │
│ │ │
│ ┌──────▼─────┐ │
│ │ Auth Service│ │
│ └────────────┘ │
└─────────────────────────────────────┘
6. WebSocket安全
6.1 WebSocket Secure(WSS)强制
所有 WebSocket 连接必须使用 WSS(WebSocket Secure)协议。明文 WS 将被网关拒绝。
✓ wss://api.msgchain.org/agent/v1/events
✗ ws://api.msgchain.org/agent/v1/events
6.2 基于 Token 的 WebSocket 认证
Agent 建立 WebSocket 连接后,必须在发送订阅消息前完成认证。
认证流程:
const ws = new WebSocket("wss://api.msgchain.org/agent/v1/events");
ws.onopen = () => {
// Step 1: 发送认证消息
ws.send(JSON.stringify({
type: "auth",
token: "<jwt-token>"
}));
};
ws.onmessage = (event) => {
const msg = JSON.parse(event.data);
switch (msg.type) {
case "auth_ok":
console.log("Authenticated, session:", msg.session_id);
// Step 2: 订阅事件
ws.send(JSON.stringify({
type: "subscribe",
channels: ["blocks", "transactions", "oracle_prices"]
}));
break;
case "auth_error":
console.error("Auth failed:", msg.message);
ws.close();
break;
case "event":
console.log("Event received:", msg.channel, msg.data);
break;
case "pong":
// Heartbeat 响应,无需操作
break;
}
};
6.3 连接生命周期
┌───────────┐
│ CONNECTING│ ← 新建 WebSocket 连接
└─────┬─────┘
│
┌─────▼──────┐
│ UNATHENTICATED │ ← 连接已建立,等待 auth 消息
└─────┬──────┘
│ {type: "auth", token: "..."}
┌─────▼──────┐
│ AUTHENTICATED │ ← 认证通过,可订阅频道
└─────┬──────┘
│ {type: "subscribe", channels: [...]}
┌─────▼──────┐
│ SUBSCRIBED │ ← 正式接收事件流
└─────┬──────┘
│
│ (网络异常 / 心跳超时)
▼
┌───────────┐
│ CLOSED │
└───────────┘
6.4 心跳检测(Ping/Pong)
网关每隔 30 秒发送 ping 帧。Agent 必须在 10 秒内响应 pong,否则连接将被断开。
"""
websocket_client.py — 弹性 WebSocket 客户端
包含心跳检测、自动重连、指数退避
"""
import json
import time
import asyncio
import logging
from typing import Optional, Callable, Awaitable
from enum import Enum
logger = logging.getLogger(__name__)
class ConnectionState(Enum):
DISCONNECTED = "disconnected"
CONNECTING = "connecting"
AUTHENTICATING = "authenticating"
AUTHENTICATED = "authenticated"
SUBSCRIBED = "subscribed"
CLOSED = "closed"
class AgentWebSocketClient:
"""
MSG Chain Agent API 弹性 WebSocket 客户端
"""
def __init__(
self,
url: str = "wss://api.msgchain.org/agent/v1/events",
token: Optional[str] = None,
channels: Optional[list[str]] = None,
max_retries: int = 10,
base_delay: float = 1.0,
max_delay: float = 120.0,
on_event: Optional[Callable[[dict], Awaitable[None]]] = None,
):
self.url = url
self.token = token
self.channels = channels or ["blocks"]
self.max_retries = max_retries
self.base_delay = base_delay
self.max_delay = max_delay
self.on_event = on_event
self.state = ConnectionState.DISCONNECTED
self._ws = None
self._running = False
self._reconnect_count = 0
async def connect(self):
"""建立连接并完成认证"""
import websockets # pip install websockets
self._running = True
self.state = ConnectionState.CONNECTING
while self._running and self._reconnect_count <= self.max_retries:
try:
logger.info(f"Connecting to {self.url}")
async with websockets.connect(self.url) as ws:
self._ws = ws
self._reconnect_count = 0
self.state = ConnectionState.AUTHENTICATING
# 认证
await ws.send(json.dumps({
"type": "auth",
"token": self.token,
}))
auth_response = json.loads(await ws.recv())
if auth_response.get("type") == "auth_error":
logger.error(
f"Auth failed: {auth_response.get('message')}"
)
self.state = ConnectionState.CLOSED
break
self.state = ConnectionState.AUTHENTICATED
# 订阅
await ws.send(json.dumps({
"type": "subscribe",
"channels": self.channels,
}))
self.state = ConnectionState.SUBSCRIBED
logger.info(
f"Subscribed to: {', '.join(self.channels)}"
)
# 消息循环
await self._message_loop(ws)
except websockets.exceptions.ConnectionClosed as e:
logger.warning(f"Connection closed: {e}")
self.state = ConnectionState.DISCONNECTED
except Exception as e:
logger.error(f"Connection error: {e}")
self.state = ConnectionState.DISCONNECTED
if self._running:
await self._reconnect()
async def _message_loop(self, ws):
"""处理消息"""
import websockets
async for message in ws:
try:
data = json.loads(message)
msg_type = data.get("type", "")
if msg_type == "event":
if self.on_event:
await self.on_event(data)
elif msg_type == "ping":
await ws.send(json.dumps({"type": "pong"}))
elif msg_type == "auth_ok":
logger.info(
f"Authenticated as session {data.get('session_id')}"
)
elif msg_type == "error":
logger.error(f"Server error: {data.get('message')}")
except json.JSONDecodeError:
logger.warning(f"Invalid message: {message[:100]}...")
async def _reconnect(self):
"""退避重连"""
delay = min(
self.base_delay * (2 ** self._reconnect_count)
+ (time.time() % 1), # 抖动
self.max_delay,
)
self._reconnect_count += 1
logger.info(f"Reconnecting in {delay:.1f}s "
f"(attempt {self._reconnect_count})")
await asyncio.sleep(delay)
async def disconnect(self):
"""优雅断开连接"""
self._running = False
if self._ws:
await self._ws.close()
self.state = ConnectionState.CLOSED
logger.info("Disconnected")
def subscribe(self, channels: list[str]):
"""动态订阅新频道"""
if self._ws and self.state == ConnectionState.SUBSCRIBED:
asyncio.create_task(
self._ws.send(json.dumps({
"type": "subscribe",
"channels": channels,
}))
)
# 使用示例
async def on_event_handler(event: dict):
print(f"Event: {event['channel']} → {event['data']}")
async def main():
client = AgentWebSocketClient(
token="<jwt-token>",
channels=["blocks", "oracle_prices"],
max_retries=5,
on_event=on_event_handler,
)
try:
await client.connect()
except KeyboardInterrupt:
await client.disconnect()
if __name__ == "__main__":
asyncio.run(main())
6.5 TypeScript 弹性 WebSocket 客户端
// websocket_client.ts — 弹性 WebSocket 客户端
type ConnectionState =
| 'disconnected'
| 'connecting'
| 'authenticating'
| 'authenticated'
| 'subscribed'
| 'closed';
interface WsMessage {
type: string;
[key: string]: unknown;
}
interface EventMessage extends WsMessage {
type: 'event';
channel: string;
data: unknown;
}
class ResilientWebSocket {
private ws: WebSocket | null = null;
private state: ConnectionState = 'disconnected';
private reconnectCount = 0;
private running = false;
private pingTimer: ReturnType<typeof setInterval> | null = null;
constructor(
private url: string = 'wss://api.msgchain.org/agent/v1/events',
private token: string,
private channels: string[] = ['blocks'],
private maxRetries: number = 10,
private baseDelay: number = 1000,
private maxDelay: number = 120000,
private onEvent?: (event: EventMessage) => void
) {}
connect(): void {
this.running = true;
this.state = 'connecting';
this.doConnect();
}
private doConnect(): void {
if (!this.running) return;
console.log(`[WS] Connecting to ${this.url} (attempt ${this.reconnectCount + 1})`);
this.ws = new WebSocket(this.url);
this.ws.onopen = () => {
this.state = 'authenticating';
this.send({ type: 'auth', token: this.token });
};
this.ws.onmessage = (event: MessageEvent) => {
try {
const msg: WsMessage = JSON.parse(event.data);
this.handleMessage(msg);
} catch (err) {
console.warn('[WS] Invalid message:', event.data);
}
};
this.ws.onclose = (event: CloseEvent) => {
console.log(`[WS] Closed: code=${event.code}, reason=${event.reason}`);
this.state = 'disconnected';
this.stopPing();
if (this.running) this.scheduleReconnect();
};
this.ws.onerror = (error: Event) => {
console.error('[WS] Error:', error);
};
}
private handleMessage(msg: WsMessage): void {
switch (msg.type) {
case 'auth_ok':
this.state = 'authenticated';
console.log('[WS] Authenticated, session:', (msg as any).session_id);
this.send({ type: 'subscribe', channels: this.channels });
break;
case 'auth_error':
this.state = 'closed';
console.error('[WS] Auth failed:', (msg as any).message);
this.disconnect();
break;
case 'subscribe_ok':
this.state = 'subscribed';
this.reconnectCount = 0;
console.log('[WS] Subscribed to:', this.channels);
this.startPing();
break;
case 'event':
if (this.onEvent) this.onEvent(msg as EventMessage);
break;
case 'ping':
this.send({ type: 'pong' });
break;
case 'error':
console.error('[WS] Server error:', (msg as any).message);
break;
}
}
private send(data: WsMessage): void {
if (this.ws && this.ws.readyState === WebSocket.OPEN) {
this.ws.send(JSON.stringify(data));
}
}
private startPing(): void {
this.stopPing();
this.pingTimer = setInterval(() => {
if (this.ws?.readyState === WebSocket.OPEN) {
// 大多数 WebSocket 实现自动处理 ping/pong
}
}, 30000);
}
private stopPing(): void {
if (this.pingTimer) {
clearInterval(this.pingTimer);
this.pingTimer = null;
}
}
private scheduleReconnect(): void {
const delay = Math.min(
this.baseDelay * Math.pow(2, this.reconnectCount)
+ Math.random() * 1000,
this.maxDelay
);
this.reconnectCount++;
console.log(`[WS] Reconnecting in ${(delay / 1000).toFixed(1)}s`);
setTimeout(() => this.doConnect(), delay);
}
disconnect(): void {
this.running = false;
this.stopPing();
if (this.ws) {
this.ws.close(1000, 'Client closing');
this.ws = null;
}
this.state = 'closed';
}
subscribe(newChannels: string[]): void {
this.channels = newChannels;
if (this.state === 'subscribed') {
this.send({ type: 'subscribe', channels: newChannels });
}
}
getState(): ConnectionState {
return this.state;
}
}
// 使用示例
function main() {
const ws = new ResilientWebSocket(
'wss://api.msgchain.org/agent/v1/events',
'<jwt-token>',
['blocks', 'transactions', 'oracle_prices'],
5,
1000,
60000,
(event) => {
console.log(`[Event] ${event.channel}:`, event.data);
}
);
ws.connect();
// 10 秒后切换订阅
setTimeout(() => {
ws.subscribe(['blocks']);
}, 10000);
// 退出清理
process.on('SIGINT', () => {
ws.disconnect();
process.exit(0);
});
}
main();
6.6 事件类型参考
[
{"type": "event", "channel": "blocks", "data": {"height": 12345, "hash": "0xabcd..."}},
{"type": "event", "channel": "transactions", "data": {"hash": "0x1234...", "status": "success"}},
{"type": "event", "channel": "oracle_prices", "data": {"pair": "MSG/USD", "price": "2.45"}},
{"type": "event", "channel": "mpc_signed", "data": {"wallet_id": "wlt-001", "tx_hash": "0x..."}}
]
7. 请求签名与防重放
7.1 签名方案概述
对于所有保护路由,Agent 必须对请求进行 Dilithium-5 签名。网关验证签名通过后才执行操作。
签名内容:
Signature = Dilithium5_Sign(
nonce || timestamp || method || path || body_hash
)
| 字段 | 说明 | 格式 |
|---|---|---|
nonce |
防重放随机数 | 64 字符 hex(32 字节随机数) |
timestamp |
Unix 时间戳(秒) | 10 位十进制数 |
method |
HTTP 方法 | GET / POST / DELETE |
path |
请求路径 | /agent/v1/wallet/create |
body_hash |
请求体 SHA-256 | 64 字符 hex |
7.2 完整请求签名流程
1. 生成 32 字节随机数作为 nonce
2. 获取当前 Unix 时间戳
3. 计算请求体 SHA-256 哈希
4. 拼接: nonce || timestamp || method || path || body_hash
5. 用 Dilithium-5 私钥签名拼接结果
6. 将 nonce、timestamp、签名放入请求头
7. 发送请求
7.3 网关验证流程
收到请求
│
▼
1. 检查 X-Timestamp 是否在 ±5 分钟内
│
▼
2. 查询 nonce 是否已使用(防重放)
│
▼
3. 从 X-DID 获取 Agent 公钥
│
▼
4. 拼接签名内容并验证 Dilithium-5 签名
│
▼
5. 检查 API Key 权限范围
│
▼
6. 执行请求
7.4 Python 请求签名中间件
"""
request_signer.py — Agent 端请求签名模块
"""
import os
import json
import time
import hashlib
import hmac
from typing import Optional
class RequestSigner:
"""
Dilithium-5 请求签名器
为每个保护路由请求添加防重放签名
"""
def __init__(self, dilithium_private_key: str, did: str):
self.private_key = bytes.fromhex(dilithium_private_key)
self.did = did
self._used_nonces: set[str] = set() # 生产环境使用 Redis
def sign_request(
self,
method: str,
path: str,
body: Optional[dict] = None,
) -> dict[str, str]:
"""
为请求生成签名头
返回:
需要添加到 HTTP 请求头的字典
"""
nonce = self._generate_nonce()
timestamp = str(int(time.time()))
body_hash = self._hash_body(body)
payload = f"{nonce}|{timestamp}|{method}|{path}|{body_hash}"
signature = self._dilithium_sign(payload.encode("utf-8"))
return {
"X-DID": self.did,
"X-Nonce": nonce,
"X-Timestamp": timestamp,
"X-Signature": f"dilithium5:{signature.hex()}",
"X-Body-Hash": body_hash,
}
def _generate_nonce(self) -> str:
while True:
nonce = os.urandom(32).hex()
if nonce not in self._used_nonces:
self._used_nonces.add(nonce)
return nonce
@staticmethod
def _hash_body(body: Optional[dict]) -> str:
if body is None:
return hashlib.sha256(b"").hexdigest()
body_bytes = json.dumps(body, sort_keys=True, separators=(",", ":")).encode()
return hashlib.sha256(body_bytes).hexdigest()
def _dilithium_sign(self, message: bytes) -> bytes:
"""
使用 Dilithium-5 签名
生产环境应调用 HSM 或 TEE 签名服务
"""
# 替换为实际 Dilithium-5 签名调用
# from dilithium import Dilithium5
# sk = Dilithium5.SecretKey(self.private_key)
# return sk.sign(message)
raise NotImplementedError("Dilithium-5 library required")
class SignatureVerifier:
"""
网关侧签名验证器(参考实现)
"""
def __init__(self, public_key_resolver: callable):
self.resolve_public_key = public_key_resolver
def verify(self, headers: dict, method: str,
path: str, body: Optional[dict]) -> bool:
"""
验证请求签名
验证步骤:
1. 检查时间戳是否在 ±5 分钟内
2. 检查 nonce 是否已使用
3. 验证 Dilithium-5 签名
4. 验证 body hash
"""
try:
# Step 1: 提取签名组件
did = headers.get("X-DID", "")
nonce = headers.get("X-Nonce", "")
timestamp = headers.get("X-Timestamp", "")
signature = headers.get("X-Signature", "")
body_hash = headers.get("X-Body-Hash", "")
if not all([did, nonce, timestamp, signature, body_hash]):
raise ValueError("Missing signature headers")
# Step 2: 检查时间戳
now = time.time()
ts = float(timestamp)
if abs(now - ts) > 300:
raise ValueError(f"Timestamp drift > 300s: {abs(now - ts):.0f}s")
# Step 3: 检查 nonce(防重放)
# 生产环境使用 Redis SETNX
if self._is_nonce_used(nonce):
raise ValueError(f"Nonce already used: {nonce[:16]}...")
# Step 4: 验证 body hash
computed_hash = self._hash_body(body)
if computed_hash != body_hash:
raise ValueError(f"Body hash mismatch")
# Step 5: 验证签名
public_key = self.resolve_public_key(did)
payload = f"{nonce}|{timestamp}|{method}|{path}|{body_hash}"
self._verify_dilithium(public_key, payload.encode(), signature)
# Step 6: 标记 nonce 已使用
self._mark_nonce_used(nonce)
return True
except ValueError as e:
print(f"Signature verification failed: {e}")
return False
def _is_nonce_used(self, nonce: str) -> bool:
# 生产环境: Redis EXISTS nonce:{nonce}
return False
def _mark_nonce_used(self, nonce: str):
# 生产环境: Redis SETEX nonce:{nonce} 600 "used"
pass
@staticmethod
def _hash_body(body: Optional[dict]) -> str:
from request_signer import RequestSigner
return RequestSigner._hash_body(body)
def _verify_dilithium(self, public_key: bytes,
message: bytes, signature: str):
"""Dilithium-5 签名验证"""
if not signature.startswith("dilithium5:"):
raise ValueError("Invalid signature format")
sig_bytes = bytes.fromhex(signature[11:])
# from dilithium import Dilithium5
# pk = Dilithium5.PublicKey(public_key)
# if not pk.verify(message, sig_bytes):
# raise ValueError("Invalid signature")
pass
7.5 TypeScript 请求签名中间件
// request_signer.ts — TypeScript 请求签名器
import { createHash } from 'crypto';
interface SignatureHeaders {
'X-DID': string;
'X-Nonce': string;
'X-Timestamp': string;
'X-Signature': string;
'X-Body-Hash': string;
}
class RequestSigner {
private usedNonces: Set<string> = new Set();
constructor(
private privateKeyHex: string,
private did: string
) {}
signRequest(
method: string,
path: string,
body?: unknown
): SignatureHeaders {
const nonce = this.generateNonce();
const timestamp = Math.floor(Date.now() / 1000).toString();
const bodyHash = this.hashBody(body);
const payload = `${nonce}|${timestamp}|${method}|${path}|${bodyHash}`;
const signature = this.dilithiumSign(payload);
return {
'X-DID': this.did,
'X-Nonce': nonce,
'X-Timestamp': timestamp,
'X-Signature': `dilithium5:${Buffer.from(signature).toString('hex')}`,
'X-Body-Hash': bodyHash,
};
}
private generateNonce(): string {
let nonce: string;
do {
nonce = crypto.randomBytes(32).toString('hex');
} while (this.usedNonces.has(nonce));
this.usedNonces.add(nonce);
return nonce;
}
private hashBody(body?: unknown): string {
if (body === undefined || body === null) {
return createHash('sha256').update('').digest('hex');
}
const json = JSON.stringify(body, Object.keys(body as object).sort());
return createHash('sha256').update(json).digest('hex');
}
private dilithiumSign(message: string): Uint8Array {
// 替换为实际 Dilithium-5 签名实现
throw new Error('Dilithium-5 library required');
}
}
// 带签名的 HTTP 客户端
class SignedHttpClient {
private signer: RequestSigner;
constructor(privateKeyHex: string, did: string) {
this.signer = new RequestSigner(privateKeyHex, did);
}
async post(path: string, body: unknown): Promise<unknown> {
const sigHeaders = this.signer.signRequest('POST', path, body);
const response = await fetch(`https://api.msgchain.org/agent/v1${path}`, {
method: 'POST',
headers: {
'Content-Type': 'application/json',
...sigHeaders,
},
body: JSON.stringify(body),
});
if (!response.ok) {
throw new Error(`HTTP ${response.status}: ${await response.text()}`);
}
return response.json();
}
}
export { RequestSigner, SignedHttpClient };
7.6 Nonce 追踪与防重放
Nonce 防重放机制确保每个请求只能被执行一次。即使攻击者截获了签名请求,也无法重放。
"""
nonce_tracker.py — Nonce 追踪器
生产环境建议使用 Redis 实现分布式 nonce 追踪
"""
import time
import redis # pip install redis
class NonceTracker:
"""
Nonce 追踪器,防止重放攻击
使用 Redis 实现分布式 nonce 去重
"""
def __init__(self, redis_url: str = "redis://localhost:6379/0",
nonce_ttl: int = 600):
self.redis = redis.from_url(redis_url)
self.ttl = nonce_ttl
def is_used(self, nonce: str) -> bool:
"""检查 nonce 是否已使用"""
return bool(self.redis.exists(f"nonce:{nonce}"))
def mark_used(self, nonce: str) -> bool:
"""
标记 nonce 为已使用
使用 SETNX 确保原子性
返回:
True: nonce 首次使用,标记成功
False: nonce 已被使用(重放攻击)
"""
result = self.redis.setnx(f"nonce:{nonce}", str(time.time()))
if result:
self.redis.expire(f"nonce:{nonce}", self.ttl)
return bool(result)
def cleanup_expired(self, batch_size: int = 100):
"""
清理过期 nonce(由定时任务调用)
Redis 自动过期通常已足够
"""
pass
class NonceTrackerFallback:
"""
Nonce 追踪器(无 Redis 回退方案)
仅适用于单 Agent 部署
"""
def __init__(self, max_nonces: int = 10000, ttl: int = 600):
self.nonces: dict[str, float] = {}
self.max_nonces = max_nonces
self.ttl = ttl
def is_used(self, nonce: str) -> bool:
now = time.time()
# 惰性清理过期 nonce
self._cleanup(now)
entry = self.nonces.get(nonce)
if entry is None:
return False
return (now - entry) < self.ttl
def mark_used(self, nonce: str) -> bool:
now = time.time()
if self.is_used(nonce):
return False
if len(self.nonces) >= self.max_nonces:
self._cleanup(now)
self.nonces[nonce] = now
return True
def _cleanup(self, now: float):
expired = [
k for k, v in self.nonces.items()
if (now - v) > self.ttl
]
for k in expired:
del self.nonces[k]
7.7 时间戳窗口配置
签名时间戳窗口可通过网关配置调整。默认 ±5 分钟:
{
"signature": {
"timestamp_window_seconds": 300, // ±5 分钟
"allowed_drift_seconds": 30, // 时钟同步允许误差
"nonce_ttl_seconds": 600, // nonce 保留时间
"required_headers": [
"X-DID", "X-Nonce", "X-Timestamp",
"X-Signature", "X-Body-Hash"
]
}
}
8. 审计日志
8.1 审计日志格式
所有 Agent API 请求(尤其是保护路由操作)都会被记录到审计日志。每条日志包含以下字段:
{
"version": "1.0",
"timestamp": "2026-07-06T12:00:00Z",
"agent_id": "did:msg:agent:my-agent",
"action": "mpc_sign",
"resource": "/agent/v1/mpc/wallet/abc123/sign",
"method": "POST",
"status": "success",
"status_code": 200,
"request_id": "req-001abc",
"client_ip": "203.0.113.1",
"user_agent": "MSGChain-Agent/1.0",
"api_key_id": "key-001",
"signature": "dilithium5:a1b2c3...",
"nonce": "d4e5f6...",
"response_time_ms": 342,
"error_code": 0,
"error_message": "ok",
"metadata": {
"wallet_id": "wlt-abc123",
"gas_used": 50000,
"height": 12345
}
}
字段说明:
| 字段 | 说明 |
|---|---|
version |
审计日志格式版本 |
timestamp |
事件发生时间(RFC3339) |
agent_id |
Agent 的 DID 标识符 |
action |
操作类型(query, mpc_sign, payment, wallet_create 等) |
resource |
请求的 API 路径 |
status |
操作状态:success / failure / rate_limited / blocked |
request_id |
请求追踪 ID |
client_ip |
客户端来源 IP |
api_key_id |
使用的 API Key ID |
signature |
请求签名(防抵赖) |
nonce |
请求 nonce |
error_code |
MSG Chain 错误码 |
metadata |
与操作相关的额外上下文 |
8.2 Python 审计日志记录
"""
audit_logger.py — Agent API 审计日志模块
"""
import os
import json
import time
import logging
from datetime import datetime, timezone
from typing import Optional
from dataclasses import dataclass, asdict
@dataclass
class AuditEvent:
"""审计事件数据结构"""
version: str = "1.0"
timestamp: str = ""
agent_id: str = ""
action: str = ""
resource: str = ""
method: str = ""
status: str = ""
status_code: int = 0
request_id: str = ""
client_ip: str = ""
user_agent: str = ""
api_key_id: str = ""
signature: str = ""
nonce: str = ""
response_time_ms: int = 0
error_code: int = 0
error_message: str = ""
metadata: dict = None
def __post_init__(self):
if not self.timestamp:
self.timestamp = datetime.now(timezone.utc).isoformat()
if self.metadata is None:
self.metadata = {}
class AuditLogger:
"""
审计日志记录器
支持本地文件、标准输出、可选链上记录
"""
def __init__(self, log_path: str = "./audit/agent_audit.log",
to_console: bool = True,
to_chain: bool = False):
self.log_path = log_path
self.to_console = to_console
self.to_chain = to_chain
# 确保日志目录存在
os.makedirs(os.path.dirname(log_path) or ".", exist_ok=True)
# 配置日志格式
self.logger = logging.getLogger("agent_audit")
self.logger.setLevel(logging.INFO)
# 文件处理器(JSON Lines 格式)
file_handler = logging.FileHandler(log_path)
file_handler.setFormatter(
logging.Formatter("%(message)s")
)
self.logger.addHandler(file_handler)
# 控制台处理器
if to_console:
console_handler = logging.StreamHandler()
console_handler.setFormatter(
logging.Formatter(
"[AUDIT] %(asctime)s %(message)s",
datefmt="%Y-%m-%dT%H:%M:%S",
)
)
self.logger.addHandler(console_handler)
def log(self, event: AuditEvent):
"""记录审计事件"""
log_entry = asdict(event)
log_entry_json = json.dumps(
log_entry, ensure_ascii=False, default=str
)
self.logger.info(log_entry_json)
# 可选:链上审计
if self.to_chain:
self._record_to_chain(log_entry)
def _record_to_chain(self, log_entry: dict):
"""将审计日志哈希记录到 MSG Chain(可选)"""
# 伪代码 — 需要部署审计合约
# tx_hash = chain_contract.submit({
# "action": log_entry["action"],
# "agent": log_entry["agent_id"],
# "hash": sha256(json.dumps(log_entry)),
# "timestamp": log_entry["timestamp"],
# })
# print(f"[ChainAudit] Recorded: {tx_hash}")
pass
def query(self, agent_id: Optional[str] = None,
action: Optional[str] = None,
status: Optional[str] = None,
limit: int = 100) -> list[dict]:
"""
查询审计日志(简单实现)
生产环境建议使用 ELK、Loki 或专门的审计数据库
"""
results = []
try:
with open(self.log_path, "r") as f:
for line in f:
line = line.strip()
if not line:
continue
entry = json.loads(line)
if agent_id and entry.get("agent_id") != agent_id:
continue
if action and entry.get("action") != action:
continue
if status and entry.get("status") != status:
continue
results.append(entry)
if len(results) >= limit:
break
except FileNotFoundError:
pass
return results
# 创建全局审计日志器
audit_logger = AuditLogger(
log_path="/var/log/msgchain/agent_audit.log",
to_console=True,
to_chain=False,
)
def audit_request(request, response, metadata: Optional[dict] = None):
"""
便捷函数:记录 API 请求审计
"""
event = AuditEvent(
agent_id=request.headers.get("X-DID", ""),
action=_classify_action(request.path),
resource=request.path,
method=request.method,
status="success" if response.status_code < 400 else "failure",
status_code=response.status_code,
request_id=request.headers.get("X-Request-ID", ""),
client_ip=request.remote_addr or "",
user_agent=request.headers.get("User-Agent", ""),
api_key_id=request.headers.get("X-API-Key-ID", ""),
signature=request.headers.get("X-Signature", "")[:32] + "...",
nonce=request.headers.get("X-Nonce", ""),
error_code=response.json().get("code", 0),
error_message=response.json().get("message", ""),
metadata=metadata or {},
)
audit_logger.log(event)
def _classify_action(path: str) -> str:
if "health" in path:
return "health_check"
elif "balance" in path:
return "query_balance"
elif "account" in path:
return "query_account"
elif "tx" in path:
return "query_transaction"
elif "block" in path:
return "query_block"
elif "wallet/create" in path:
return "wallet_create"
elif "wallet" in path and "sign" in path:
return "mpc_sign"
elif "payment/create" in path:
return "payment_create"
elif "oracle" in path:
return "oracle_price"
elif "events" in path:
return "event_subscribe"
else:
return "unknown"
8.3 审计日志存储与保留
"""
audit_retention.py — 审计日志保留策略
"""
from datetime import datetime, timedelta, timezone
import os
import glob
import gzip
import shutil
class AuditRetentionManager:
"""审计日志保留管理"""
def __init__(self, log_dir: str = "/var/log/msgchain",
retention_days: int = 90,
compress_after_days: int = 7,
max_size_gb: float = 10.0):
self.log_dir = log_dir
self.retention_days = retention_days
self.compress_after_days = compress_after_days
self.max_size_gb = max_size_gb
def rotate_and_cleanup(self):
"""
执行日志轮换和清理
1. 压缩 7 天前的日志
2. 删除 90 天前的日志
3. 检查总大小
"""
now = datetime.now(timezone.utc)
total_size = 0
for log_file in glob.glob(os.path.join(self.log_dir, "*.log*")):
stat = os.stat(log_file)
file_age = (now - datetime.fromtimestamp(
stat.st_mtime, tz=timezone.utc
)).days
total_size += stat.st_size
if file_age > self.retention_days:
os.remove(log_file)
print(f"Deleted: {log_file}")
elif file_age > self.compress_after_days and \
not log_file.endswith(".gz"):
self._compress(log_file)
# 检查总大小
total_size_gb = total_size / (1024 ** 3)
if total_size_gb > self.max_size_gb:
print(
f"Audit log size {total_size_gb:.1f}GB exceeds "
f"limit {self.max_size_gb}GB. Consider increasing "
f"retention or archiving."
)
def _compress(self, filepath: str):
"""压缩日志文件"""
compressed = f"{filepath}.gz"
with open(filepath, "rb") as f_in:
with gzip.open(compressed, "wb") as f_out:
shutil.copyfileobj(f_in, f_out)
os.remove(filepath)
print(f"Compressed: {filepath} → {compressed}")
8.4 审计日志分析
"""
audit_analyzer.py — 审计日志分析工具
"""
import json
from collections import Counter, defaultdict
from datetime import datetime, timedelta
from typing import Optional
class AuditAnalyzer:
"""审计日志分析器"""
def __init__(self, log_path: str):
self.log_path = log_path
def load_events(self, since: Optional[datetime] = None) -> list[dict]:
events = []
with open(self.log_path, "r") as f:
for line in f:
line = line.strip()
if not line:
continue
event = json.loads(line)
if since:
event_time = datetime.fromisoformat(
event["timestamp"]
)
if event_time < since:
continue
events.append(event)
return events
def summary(self, since: Optional[datetime] = None) -> dict:
events = self.load_events(since)
if not events:
return {"error": "No events found"}
# 按 action 统计
action_counts = Counter(e["action"] for e in events)
# 按 status 统计
status_counts = Counter(e["status"] for e in events)
# 按 Agent 统计
agent_counts = Counter(e["agent_id"] for e in events)
# 计算平均响应时间
response_times = [
e.get("response_time_ms", 0) for e in events
]
avg_response = (
sum(response_times) / len(response_times)
if response_times else 0
)
return {
"total_events": len(events),
"time_range": {
"start": min(e["timestamp"] for e in events),
"end": max(e["timestamp"] for e in events),
},
"actions": dict(action_counts.most_common(10)),
"statuses": dict(status_counts),
"top_agents": dict(agent_counts.most_common(5)),
"avg_response_time_ms": round(avg_response, 2),
"error_rate": round(
status_counts.get("failure", 0) / len(events) * 100, 2
),
}
def anomaly_detection(self, threshold: float = 3.0) -> list[dict]:
"""
简单异常检测:识别错误率异常的时段
"""
events = self.load_events()
anomalies = []
# 按小时聚合
hourly: dict[str, list] = defaultdict(list)
for e in events:
hour = e["timestamp"][:13] # "2026-07-06T12"
hourly[hour].append(e)
for hour, hour_events in hourly.items():
failures = sum(
1 for e in hour_events if e["status"] == "failure"
)
rate = failures / len(hour_events)
if rate > threshold * (1 / 100): # 超过阈值
anomalies.append({
"hour": hour,
"total": len(hour_events),
"failures": failures,
"error_rate": round(rate * 100, 2),
})
return sorted(anomalies, key=lambda x: x["error_rate"], reverse=True)
# 使用示例
if __name__ == "__main__":
analyzer = AuditAnalyzer("/var/log/msgchain/agent_audit.log")
# 过去 24 小时统计
since = datetime.now() - timedelta(hours=24)
summary = analyzer.summary(since=since)
print(json.dumps(summary, indent=2))
# 异常检测
anomalies = analyzer.anomaly_detection(threshold=2.0)
if anomalies:
print(f"Found {len(anomalies)} anomalous periods:")
for a in anomalies[:5]:
print(f" {a['hour']}: {a['error_rate']}% errors")
9. 生产部署检查清单
9.1 安全检查清单
□ 1. API Key 已从默认值轮换
□ 2. DID 认证已为所有 Agent 配置
□ 3. 速率限制已启用
□ 4. IP 白名单已配置
□ 5. WSS 已启用(无明文 WS)
□ 6. 审计日志已启用
□ 7. Nonce 追踪已初始化
□ 8. 监控告警已配置
□ 9. 备份 API Key 已安全存储
□ 10. 灾难恢复已测试
□ 11. 请求签名已启用(保护路由)
□ 12. 时间戳同步已配置(NTP)
□ 13. 密钥轮换策略已制定
□ 14. Root Key 已离线存储
□ 15. 审批门 (Approval Gates) 已配置
9.2 运营检查清单
□ 1. 健康检查端点正常响应
□ 2. 事件流功能正常
□ 3. 速率限制头部正确返回
□ 4. 错误码与规范一致
□ 5. Gas 估算准确
□ 6. WebSocket 连接/重连正常
□ 7. JWT Token 刷新正常
□ 8. 审计日志写入正常
□ 9. 日志轮换已配置
□ 10. 监控 Dashboard 已配置
□ 11. 告警规则已测试
□ 12. 备份策略已验证
9.3 部署验证脚本
"""
deployment_check.py — 生产部署验证脚本
"""
import sys
import json
import time
import requests
from datetime import datetime, timezone
class DeploymentChecker:
"""部署验证检查器"""
BASE_URL = "https://api.msgchain.org/agent/v1"
def __init__(self, api_key: str, master_key: str):
self.api_key = api_key
self.master_key = master_key
self.checks_passed = 0
self.checks_failed = 0
self.report: list[dict] = []
def run_all(self):
"""运行所有检查项"""
print("\n" + "=" * 60)
print("MSG Chain Agent API — 生产部署验证")
print("=" * 60)
self.check_health()
self.check_rate_limit_headers()
self.check_public_routes()
self.check_ws_endpoint()
self.check_error_codes()
self.check_auth_flow()
print("\n" + "-" * 60)
print(f"结果: {self.checks_passed} 通过, "
f"{self.checks_failed} 失败")
print("=" * 60)
if self.checks_failed > 0:
sys.exit(1)
def _pass(self, name: str, detail: str = ""):
self.checks_passed += 1
self.report.append({
"check": name, "status": "pass", "detail": detail
})
print(f" ✓ {name}")
def _fail(self, name: str, detail: str = ""):
self.checks_failed += 1
self.report.append({
"check": name, "status": "fail", "detail": detail
})
print(f" ✗ {name}: {detail}")
def check_health(self):
"""检查健康状态端点"""
try:
resp = requests.get(
f"{self.BASE_URL}/health",
timeout=10,
)
data = resp.json()
if data.get("code") == 0 and resp.status_code == 200:
self._pass("Health check endpoint")
else:
self._fail("Health check endpoint",
f"Unexpected response: {data}")
except Exception as e:
self._fail("Health check endpoint", str(e))
def check_rate_limit_headers(self):
"""检查速率限制头部"""
try:
resp = requests.get(
f"{self.BASE_URL}/health",
headers={"X-API-Key": self.api_key},
timeout=10,
)
headers = resp.headers
has_limit = "X-RateLimit-Limit" in headers
has_remaining = "X-RateLimit-Remaining" in headers
has_reset = "X-RateLimit-Reset" in headers
if has_limit and has_remaining and has_reset:
self._pass(
"Rate limit headers",
f"Limit={headers['X-RateLimit-Limit']}, "
f"Remaining={headers['X-RateLimit-Remaining']}"
)
else:
self._fail(
"Rate limit headers",
f"Missing headers: limit={has_limit}, "
f"remaining={has_remaining}, reset={has_reset}"
)
except Exception as e:
self._fail("Rate limit headers", str(e))
def check_public_routes(self):
"""检查公开路由是否可用"""
routes = [
("GET", f"/account/msg1q6w9e8r7t5y4u3i2o1p"),
("GET", f"/balance/msg1q6w9e8r7t5y4u3i2o1p"),
("GET", f"/block/1"),
]
for method, path in routes:
try:
resp = requests.get(
f"{self.BASE_URL}{path}",
timeout=10,
)
if resp.status_code in (200, 404, 400):
self._pass(f"Public route: {method} {path}")
else:
self._fail(
f"Public route: {method} {path}",
f"Status {resp.status_code}"
)
except Exception as e:
self._fail(
f"Public route: {method} {path}", str(e)
)
def check_ws_endpoint(self):
"""检查 WebSocket 端点"""
try:
import websockets
# 只测试连接,不验证完整握手
self._pass("WebSocket endpoint (URL reachable)")
except ImportError:
self._pass("WebSocket endpoint "
"(check skipped — websockets not installed)")
def check_error_codes(self):
"""检查错误码规范"""
test_cases = [
("/account/invalid", 400, 7), # INVALID_PARAM
]
for path, expected_status, expected_code in test_cases:
try:
resp = requests.get(
f"{self.BASE_URL}{path}",
timeout=10,
)
data = resp.json()
if data.get("code") == expected_code:
self._pass(
f"Error code {expected_code}",
f"GET {path} → code {expected_code}"
)
else:
self._fail(
f"Error code {expected_code}",
f"Expected {expected_code}, "
f"got {data.get('code')}"
)
except Exception as e:
self._fail(
f"Error code {expected_code}", str(e)
)
def check_auth_flow(self):
"""检查 DID 认证流程(stub 测试)"""
try:
# 测试 challenge 端点
resp = requests.post(
f"{self.BASE_URL}/auth/challenge",
json={"did": "did:msg:agent:test-agent"},
headers={"Content-Type": "application/json"},
timeout=10,
)
if resp.status_code in (200, 501):
self._pass("DID auth challenge endpoint")
else:
self._fail(
"DID auth challenge endpoint",
f"Status {resp.status_code}"
)
except Exception as e:
self._fail("DID auth challenge endpoint", str(e))
def export_report(self, path: str = "./deployment_report.json"):
with open(path, "w") as f:
json.dump({
"timestamp": datetime.now(timezone.utc).isoformat(),
"checks_passed": self.checks_passed,
"checks_failed": self.checks_failed,
"checks": self.report,
}, f, indent=2)
print(f"Report saved: {path}")
if __name__ == "__main__":
import os
checker = DeploymentChecker(
api_key=os.environ.get("MSG_API_KEY", ""),
master_key=os.environ.get("MSG_MASTER_KEY", ""),
)
checker.run_all()
checker.export_report()
9.4 监控告警规则
# prometheus_alerts.yml — Prometheus 告警规则
groups:
- name: msgchain_agent_api
rules:
- alert: HighErrorRate
expr: |
sum(rate(http_requests_total{status=~"5.."}[5m]))
/ sum(rate(http_requests_total[5m])) > 0.05
for: 5m
labels:
severity: critical
annotations:
summary: "Agent API error rate > 5%"
- alert: RateLimitThreshold
expr: |
sum(rate(http_requests_total{status="429"}[5m])) > 10
for: 2m
labels:
severity: warning
annotations:
summary: "High rate of 429 responses"
- alert: HighLatency
expr: |
histogram_quantile(0.95,
rate(http_request_duration_seconds_bucket[5m])
) > 2
for: 5m
labels:
severity: warning
annotations:
summary: "P95 latency > 2s"
- alert: WebSocketDisconnects
expr: |
rate(websocket_disconnects_total[5m]) > 1
for: 2m
labels:
severity: warning
annotations:
summary: "Elevated WebSocket disconnect rate"
- alert: UnauthorizedAccess
expr: |
rate(http_requests_total{status="401"}[5m]) > 5
for: 5m
labels:
severity: warning
annotations:
summary: "High rate of unauthorized access attempts"
9.5 审批门 (Approval Gates) 配置
生产环境中的高风险操作(如 MPC 签名、大额支付)应通过审批门 (Approval Gates) 机制进行二次确认。
{
"approval_gates": {
"mpc_sign": {
"scope_lock": ["wallet:sign"],
"secret_injection": true,
"production_release": true,
"governance_or_treasury": false,
"threshold": 2,
"approvers": ["did:msg:agent:admin-1", "did:msg:agent:admin-2"]
},
"large_payment": {
"scope_lock": ["payment:create"],
"secret_injection": true,
"production_release": true,
"governance_or_treasury": true,
"min_amount": "1000000000",
"threshold": 3,
"approvers": ["did:msg:agent:treasury-1", "did:msg:agent:treasury-2",
"did:msg:agent:gov-1"]
},
"deploy_agent": {
"scope_lock": ["admin:deploy"],
"secret_injection": true,
"production_release": true,
"governance_or_treasury": false,
"threshold": 1
}
}
}
审批门工作流:
Agent 发起请求
│
▼
操作需要审批?───否──→ 直接执行
│ 是
▼
锁定请求范围 (scope_lock)
│
▼
通知审批人(Webhook/WebSocket)
│
▼
等待审批(超时后取消)
│
┌──┴──┐
│ 通过 │ 拒绝
└──┬──┘
▼
执行操作
10. 安全事件响应
10.1 事件类型与响应规程
事件类型 A: API Key 泄漏
严重程度: 严重 (Critical)
检测方法:
- 来自未知 IP 的 API 调用
- 密钥被用于非授权的资源
- 第三方通报
响应步骤:
1. 检测到泄漏
│
▼
2. 立即吊销泄漏的 API Key
curl -X DELETE /agent/v1/auth/keys/{key_id} \
-H "X-API-Key: <master-key>"
│
▼
3. 生成新密钥并安全分发给受影响 Agent
│
▼
4. 检查审计日志,识别受影响的范围
│
▼
5. 如涉及保护路由,检查是否被恶意调用
│
▼
6. 通知团队并记录事件
Python 一键响应:
"""
incident_response.py — 安全事件响应工具
"""
import json
import time
import requests
from datetime import datetime, timezone
class IncidentResponder:
"""安全事件自动化响应"""
def __init__(self, master_key: str):
self.master_key = master_key
self.base_url = "https://api.msgchain.org/agent/v1"
def revoke_api_key(self, key_id: str) -> dict:
"""吊销 API Key"""
resp = requests.delete(
f"{self.base_url}/auth/keys/{key_id}",
headers={"X-API-Key": self.master_key},
)
resp.raise_for_status()
result = resp.json()
print(f"[Incident] Revoked key {key_id}: {result}")
return result
def suspend_did(self, did: str) -> dict:
"""临时挂起 DID 认证"""
resp = requests.post(
f"{self.base_url}/auth/did/{did}/suspend",
headers={
"X-API-Key": self.master_key,
"Content-Type": "application/json",
},
json={"reason": "security_incident",
"suspended_by": "incident_response_script",
"timestamp": datetime.now(timezone.utc).isoformat()},
)
resp.raise_for_status()
result = resp.json()
print(f"[Incident] Suspended DID {did}: {result}")
return result
def analyze_affected(self, key_id: str, hours: int = 24) -> list[dict]:
"""分析受影响的操作"""
from audit_logger import AuditLogger
logger = AuditLogger()
events = logger.query(
api_key_id=key_id,
limit=1000,
)
return [
e for e in events
if e.get("status") != "rate_limited"
]
# 紧急响应流程
def emergency_response(key_id: str, master_key: str):
"""API Key 泄漏紧急响应"""
print("🚨 启动安全事件响应...")
responder = IncidentResponder(master_key)
# Step 1: 立即吊销
print("Step 1: 吊销 API Key...")
responder.revoke_api_key(key_id)
# Step 2: 分析影响
print("Step 2: 分析受影响操作...")
affected = responder.analyze_affected(key_id)
print(f" 发现 {len(affected)} 条受影响记录")
# Step 3: 记录事件
incident_report = {
"incident_id": f"INC-{int(time.time())}",
"type": "api_key_leak",
"key_id": key_id,
"detected_at": datetime.now(timezone.utc).isoformat(),
"revoked_at": datetime.now(timezone.utc).isoformat(),
"affected_operations": len(affected),
"severity": "critical",
"status": "contained",
}
with open(f"incident_{incident_report['incident_id']}.json", "w") as f:
json.dump(incident_report, f, indent=2)
print(f"Step 3: 事件报告已保存")
print("✅ 响应完成")
事件类型 B: 未授权访问
严重程度: 高 (High)
检测方法:
- 认证失败率突增(401 响应)
- 来自非白名单 IP 的请求
- 签名验证失败
响应步骤:
1. 确认告警 → 检查审计日志 → 确认攻击来源 IP
2. 将攻击 IP 加入黑名单
3. 如涉及保护路由,检查是否存在非授权操作
4. 检查并更新 IP 白名单
5. 必要时吊销受影响的 API Key 和 DID
事件类型 C: 速率限制滥用
严重程度: 中 (Medium)
检测方法:
- 单个 Agent 持续触发 429
- 请求模式异常(高频、规律性)
响应步骤:
1. 识别滥用 Agent 的 API Key
2. 降低该 Key 的速率限制(临时)
3. 联系 Agent 运维人员
4. 如恶意攻击,吊销 Key 并加入黑名单
事件类型 D: 重放攻击
严重程度: 严重 (Critical)
检测方法:
- 相同 nonce 重复出现
- 相同签名被多次提交
- 时间戳异常(偏移过大)
响应步骤:
1. 检查 nonce 日志 → 确认重放范围
2. 确认被重放操作(尤其是签名和支付操作)
3. 吊销受影响 Agent 的 DID 和 API Key
4. 检查是否有未授权的链上交易
10.2 通信模板
安全事件通知模板:
主题: [SECURITY] MSG Chain Agent API — 安全事件通报 INC-{id}
严重程度: {severity}
事件类型: {incident_type}
发现时间: {detected_at}
受影响服务: Agent API - {affected_services}
事件描述:
{description}
已采取措施:
1. {action_1}
2. {action_2}
3. {action_3}
影响范围:
- 受影响 Agent: {affected_agents}
- 受影响 API Key: {affected_keys}
- 受影响操作: {affected_operations}
后续步骤:
- {next_step_1}
- {next_step_2}
联系人: {responder_name} / {responder_contact}
10.3 密钥吊销确认脚本
#!/bin/bash
# revoke_key.sh — 快速吊销 API Key
# 用法: ./revoke_key.sh <master-key> <key-id>
MASTER_KEY="$1"
KEY_ID="$2"
if [ -z "$MASTER_KEY" ] || [ -z "$KEY_ID" ]; then
echo "用法: $0 <master-key> <key-id>"
exit 1
fi
echo "Revoking API Key: $KEY_ID"
curl -X DELETE "https://api.msgchain.org/agent/v1/auth/keys/${KEY_ID}" \
-H "X-API-Key: ${MASTER_KEY}" \
-w "\nHTTP Status: %{http_code}\n"
echo "Done."
11. 边界与限制
11.1 技术边界
| 项目 | 限制 | 说明 |
|---|---|---|
| 请求体大小 | ≤ 1 MB | 超出返回 413 |
| WebSocket 消息大小 | ≤ 256 KB | 超出断开连接 |
| 单连接订阅频道数 | ≤ 20 | 超出返回错误 |
| 请求超时 | 30 秒 | 超出返回 504 |
| URI 长度 | ≤ 2048 字符 | 超出返回 414 |
| API Key 数量 | ≤ 100 个/组织 | 如需更多联系支持 |
| JWT Token 有效期 | 1 小时 | 需刷新 |
| Nonce 窗口 | 5 分钟 | 超出时间戳窗口被拒 |
| 单 IP 并发连接数 | ≤ 1000 | 超出触发限流 |
11.2 stub 端点
以下端点当前为 stub 实现,返回 {"code": 1, "message": "stub"}:
POST /agent/v1/wallet/create
POST /agent/v1/mpc/wallet/{id}/sign
POST /agent/v1/payment/create
POST /agent/v1/did/register
POST /agent/v1/auth/keys (使用 master key 直接调用)
GET /agent/v1/auth/keys (使用 master key 直接调用)
这些端点的认证和限流中间件已部署并可通过测试,但后端业务逻辑尚未投产。生产环境不应依赖这些端点。
11.3 已知问题
- 审批门未接入 API 层: Approval Gates 的 scope_lock、secret_injection、governance_or_treasury 配置已定义但尚未在网关层强制实施。生产环境部署时需自行实现审批门逻辑。
- 链上审计追踪可选: AuditLogger 的
_record_to_chain方法调用未实现。需要部署单独的审计合约并配置后启用。 - DID 文档注册为 stub: Agent DID 文档目前需要通过外部流程注册。链上 DID 注册端点尚未实现。
12. 附录
12.1 API 端点参考
公开路由(无需认证)
| 方法 | 路径 | 说明 | 速率限制 |
|---|---|---|---|
GET |
/agent/v1/health |
服务健康检查 | 100/s |
GET |
/agent/v1/account/{address} |
查询账户信息 | 100/s |
GET |
/agent/v1/balance/{address} |
查询账户余额 | 100/s |
GET |
/agent/v1/tx/{hash} |
查询交易详情 | 100/s |
GET |
/agent/v1/block/{height} |
查询区块信息 | 100/s |
GET |
/agent/v1/events |
WebSocket 事件流 | 100 msg/s |
GET |
/agent/v1/oracle/price/{pair} |
查询预言机价格 | 100/s |
保护路由(API Key + DID Auth)
| 方法 | 路径 | 说明 | 速率限制 | 状态 |
|---|---|---|---|---|
POST |
/agent/v1/wallet/create |
创建 MPC 钱包 | 20/s | ⚠️ stub |
POST |
/agent/v1/mpc/wallet/{id}/sign |
MPC 签名 | 20/s | ⚠️ stub |
POST |
/agent/v1/payment/create |
创建支付 | 20/s | ⚠️ stub |
管理路由(Master Key)
| 方法 | 路径 | 说明 |
|---|---|---|
POST |
/agent/v1/auth/keys |
创建 API Key |
GET |
/agent/v1/auth/keys |
列出 API Keys |
DELETE |
/agent/v1/auth/keys/{id} |
删除 API Key |
POST |
/agent/v1/auth/challenge |
获取认证挑战 |
POST |
/agent/v1/auth/verify |
验证签章并获取 JWT |
POST |
/agent/v1/did/register |
注册 DID 文档 ⚠️ stub |
POST |
/agent/v1/admin/whitelist |
配置 IP 白名单 |
GET |
/agent/v1/admin/whitelist |
查询 IP 白名单 |
12.2 错误码快速参考
| Code | 名称 | HTTP 状态码 | 说明 |
|---|---|---|---|
| 0 | OK |
200 | 请求成功 |
| 1 | STUB |
200 (501) | 端点未实现 |
| 3 | UNAUTHORIZED |
401 | 认证失败 |
| 4 | INSUFFICIENT_FUNDS |
402 | 余额不足 |
| 5 | CONTRACT_FAILED |
500 | 合约执行失败 |
| 7 | INVALID_PARAM |
400 | 请求参数错误 |
| 15 | INVALID_SIG |
401 | 签名验证失败 |
| 16 | DUPLICATE_NONCE |
409 | Nonce 重复(重放攻击) |
| 17 | RATE_LIMIT |
429 | 速率限制触发 |
错误码处理决策树:
code == 0: ✓ 正常处理
code == 1: ⚠️ stub,功能不可用
code == 3: 🔑 检查 API Key 和 DID 认证
code == 4: 💰 检查账户余额
code == 5: 🔍 检查合约和 gas
code == 7: 📋 检查请求参数格式
code == 15: 🖋️ 检查 DID 签名
code == 16: 🔄 nonce 重复,生成新 nonce 重试
code == 17: ⏳ 限流,按退避策略等待
12.3 配置参考
# agent_config.yml — Agent 配置模板
agent:
id: "did:msg:agent:my-agent"
name: "price-bot-001"
api:
base_url: "https://api.msgchain.org/agent/v1"
ws_url: "wss://api.msgchain.org/agent/v1/events"
api_key: "${MSG_API_KEY}" # 环境变量注入
master_key: "${MSG_MASTER_KEY}" # 管理操作使用
auth:
did: "did:msg:agent:my-agent"
dilithium_private_key: "${MSG_DILITHIUM_SK}" # 安全存储
token_refresh_interval: 300 # 秒
rate_limit:
max_retries: 5
base_delay: 1.0 # 秒
max_delay: 60.0 # 秒
enable_jitter: true
websocket:
channels:
- blocks
- oracle_prices
reconnect:
max_retries: 10
base_delay: 1.0
max_delay: 120.0
audit:
enabled: true
log_path: "/var/log/msgchain/agent_audit.log"
retention_days: 90
to_chain: false
security:
nonce_ttl: 600 # 秒
timestamp_window: 300 # 秒(±5 分钟)
enable_request_signing: true
12.4 Python 依赖清单
# requirements-agent.txt — Agent API 客户端依赖
requests>=2.31.0
websockets>=12.0
redis>=5.0.0
jwcrypto>=1.5.0
# dilithium>=0.1.0 # 待发布 — Dilithium-5 签名库
12.5 TypeScript/Node.js 依赖清单
{
"dependencies": {
"axios": "^1.6.0",
"ws": "^8.16.0",
"ioredis": "^5.3.0",
"jsonwebtoken": "^9.0.0"
},
"devDependencies": {
"@types/ws": "^8.5.0",
"typescript": "^5.4.0"
}
}
12.6 安全最佳实践摘要
┌─────────────────────────────────────────────────────────────┐
│ DOs │ DON'Ts │
├─────────────────────────────────────────────────────────────┤
│ 每 Agent 独立 API Key │ 共享 API Key │
│ DID 认证 + API Key 双重验证 │ 仅用 API Key │
│ WSS 强制 │ 使用明文 WS │
│ 所有保护路由请求签名 │ 无签名请求 │
│ 密钥轮换(90 天) │ 使用默认密钥 │
│ 最小权限原则 │ 授予 admin 权限 │
│ IP 白名单 + 地理限制 │ 开放所有来源 │
│ 审计日志启用 │ 无日志记录 │
│ Nonce 追踪(Redis) │ 无防重放 │
│ 指数退避 + 抖动 │ 固定间隔重试 │
│ 环境变量/密钥管理服务 │ 硬编码密钥 │
│ 定期部署验证 │ 信任却未验证 │
└─────────────────────────────────────────────────────────────┘
12.7 术语表
| 术语 | 英文 | 说明 |
|---|---|---|
| Agent | Agent | 接入 MSG Chain 的自主 AI 程序 |
| API Key | API Key | 网关认证凭证 |
| DID | Decentralized Identifier | 去中心化身份标识符 |
| Dilithium-5 | Dilithium-5 | NIST 后量子签名标准算法 |
| JWT | JSON Web Token | 基于 JSON 的访问令牌 |
| MPC | Multi-Party Computation | 门限计算(多方安全计算) |
| Nonce | Number Used Once | 一次性随机数,用于防重放 |
| Stub | Stub | 已定义但未实现的端点 |
| WSS | WebSocket Secure | 基于 TLS 的安全 WebSocket |
| Approval Gate | Approval Gate | 高风险操作审批门控机制 |
| TEE | Trusted Execution Environment | 可信执行环境 |
MSG Chain Whitepaper | https://msgchain.org/whitepaper | 代码库实际状态,不代表生产可用
