dApp Docs/Agent API 网关安全与运维指南
Development reference. Not independently verified for production.

MSG Chain Agent API 网关安全与运维指南

数据来源:MSG Chain 代码库核实

主网状态: No-Go — 当前 MSGChain 主网裁决为 No-Go,以下内容反映代码实际状态,不代表生产可用。

面向 AI Agent 开发者的安全接入手册
协议版本: Agent API v1 | 链 ID: msg-chain-1 | Bech32 前缀: msg
状态: 部分实现 — 标记为 stub 的端点尚未投产


目录

  1. 概述
  2. API密钥管理
  3. DID身份认证
  4. 速率限制策略
  5. IP白名单与地理限制
  6. WebSocket安全
  7. 请求签名与防重放
  8. 审计日志
  9. 生产部署检查清单
  10. 安全事件响应
  11. 边界与限制
  12. 附录

1. 概述

1.1 Agent API 架构

MSG Chain Agent API 是一组专为 AI Agent 设计的 RESTful + WebSocket 接口,使 Agent 能够:

网关架构采用三层隔离模型:

  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 状态说明

本文档描述的某些功能处于"部分实现"状态:

所有代码示例中的密钥和端点均为示例,不可用于生产。


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 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>"

撤销后建议操作:

  1. 确认 Agent 已切换至新密钥
  2. 监控旧密钥是否有异常访问(理论上不应有)
  3. 如怀疑泄漏,立即触发安全事件响应

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)

检测方法:

响应步骤:

  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)

检测方法:

响应步骤:

  1. 确认告警 → 检查审计日志 → 确认攻击来源 IP
  2. 将攻击 IP 加入黑名单
  3. 如涉及保护路由,检查是否存在非授权操作
  4. 检查并更新 IP 白名单
  5. 必要时吊销受影响的 API Key 和 DID

事件类型 C: 速率限制滥用

严重程度: 中 (Medium)

检测方法:

响应步骤:

  1. 识别滥用 Agent 的 API Key
  2. 降低该 Key 的速率限制(临时)
  3. 联系 Agent 运维人员
  4. 如恶意攻击,吊销 Key 并加入黑名单

事件类型 D: 重放攻击

严重程度: 严重 (Critical)

检测方法:

响应步骤:

  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 已知问题

  1. 审批门未接入 API 层: Approval Gates 的 scope_lock、secret_injection、governance_or_treasury 配置已定义但尚未在网关层强制实施。生产环境部署时需自行实现审批门逻辑。
  2. 链上审计追踪可选: AuditLogger 的 _record_to_chain 方法调用未实现。需要部署单独的审计合约并配置后启用。
  3. 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 | 代码库实际状态,不代表生产可用