dApp Docs/AI Agent 安全运维手册
Development reference. Not independently verified for production.

MSG Chain AI Agent 安全运维手册

⚠️ No-Go Disclaimer: MSGChain 主网裁决为 No-Go。本文件所有内容反映的是开发阶段的技术设计,不代表主网未独立核验上线状态。生产部署状态请以白皮书为准:https://msgchain.org/whitepaper/

适用链:msg-chain-1 | 地址前缀:msg | 版本:v1.0


目录

  1. 概述
  2. 密钥管理
  3. 访问控制
  4. 威胁监控
  5. 自动化响应
  6. 审计与合规
  7. 安全配置模板
  8. 应急响应流程

1. 概述

1.1 为什么 Agent 需要专属安全运维

AI Agent 在 MSG Chain 上以自治方式运行,管理链上资产、执行跨链消息、交互智能合约。与传统区块链账户不同,Agent 具备以下特征:

这些特性决定了 Agent 的安全运维策略必须与普通 EOA 账户或合约账户显著不同。

1.2 威胁建模

针对 MSG Chain AI Agent 的典型威胁向量包括:

威胁类别 示例 影响
密钥泄露 Agent 私钥被提取 资产被盗、身份冒用
Prompt 注入 恶意输入改变 Agent 行为 违反宪法、执行未授权操作
Gas 耗尽 持续无意义交易耗尽资金 Agent 停服
重放攻击 捕获并重放 Agent 签名消息 重复扣款
宪法绕过 发现逻辑漏洞绕过 Constitution 约束 恶意行为
A2A 欺骗 伪造 Agent 身份发送恶意消息 跨 Agent 感染
操控预言机 操纵 Agent 依赖的外部数据源 误导决策

1.3 安全运维域

本手册覆盖以下六大安全域:

msg-agent-security
├── 密钥管理         # 密钥生命周期、HSM/TEE、轮换
├── 访问控制         # RBAC、DID、多签、限流
├── 威胁监控         # 异常检测、告警规则
├── 自动化响应       # Playbook、断路器、分级响应
├── 审计与合规       # 日志、合规检查、渗透测试
└── 应急响应         # 事件分类、SOP、复盘

每个安全域均包含理论说明、配置示例和完整可运行代码。

1.4 设计原则

Agent 安全运维遵循以下原则:

  1. 最小权限:Agent 仅获得执行任务所需的最低权限
  2. 纵深防御:多层安全控制,单点失效不影响整体
  3. 默认拒绝:未明确允许的操作一律禁止
  4. 可审计性:所有操作记录不可篡改的审计日志
  5. 分级响应:根据威胁严重度采取不同级别的响应措施
  6. 自动恢复:在安全前提下尽可能自动化恢复服务

2. 密钥管理

2.1 Dilithium-5 密钥对

MSG Chain 采用基于 CRYSTALS-Dilithium 的后量子签名算法 Dilithium-5 作为 Agent 身份认证和交易签名的核心原语。Dilithium-5 提供 NIST Level 5 安全强度,能够抵御量子计算攻击。

// agent-key-manager.ts
// Agent 密钥管理器核心接口

interface AgentKeyManager {
  signingKey: CryptoKey;
  verificationKey: CryptoKey;
  rotationPeriod: number;
  lastRotation: number;
  agentAddress: string;
}

interface Dilithium5KeyPair {
  publicKey: Uint8Array;
  privateKey: Uint8Array;
}

class AgentCryptoProvider {
  private readonly DILITHIUM5_PUBLIC_KEY_LENGTH = 2592;
  private readonly DILITHIUM5_PRIVATE_KEY_LENGTH = 4864;
  private readonly MSG_ADDRESS_PREFIX = 'msg';

  async generateKeyPair(): Promise<Dilithium5KeyPair> {
    const keyPair = await crypto.subtle.generateKey(
      { name: 'CRYSTALS-Dilithium', version: 5, modulusLength: 256 },
      true,
      ['sign', 'verify']
    );
    const publicKey = new Uint8Array(
      await crypto.subtle.exportKey('spki', keyPair.publicKey)
    );
    const privateKey = new Uint8Array(
      await crypto.subtle.exportKey('pkcs8', keyPair.privateKey)
    );
    return { publicKey, privateKey };
  }

  deriveAddress(publicKey: Uint8Array): string {
    const hash = crypto.subtle.digestSync('SHA-256', publicKey);
    const addrBytes = new Uint8Array(hash).slice(0, 20);
    return this.MSG_ADDRESS_PREFIX + Buffer.from(addrBytes).toString('hex');
  }

  async signMessage(
    privateKey: Uint8Array, message: Uint8Array
  ): Promise<Uint8Array> {
    const key = await crypto.subtle.importKey(
      'pkcs8', privateKey,
      { name: 'CRYSTALS-Dilithium', version: 5 },
      false, ['sign']
    );
    const signature = await crypto.subtle.sign(
      { name: 'CRYSTALS-Dilithium', version: 5 }, key, message
    );
    return new Uint8Array(signature);
  }

  async verifySignature(
    publicKey: Uint8Array, message: Uint8Array, signature: Uint8Array
  ): Promise<boolean> {
    const key = await crypto.subtle.importKey(
      'spki', publicKey,
      { name: 'CRYSTALS-Dilithium', version: 5 },
      false, ['verify']
    );
    return crypto.subtle.verify(
      { name: 'CRYSTALS-Dilithium', version: 5 }, key, signature, message
    );
  }

  packTransaction(
    tx: UnsignedTransaction, signature: Uint8Array, publicKey: Uint8Array
  ): Uint8Array {
    const txBytes = new Uint8Array(
      4 + tx.body.length + 2 + signature.length + 2 + publicKey.length
    );
    const view = new DataView(txBytes.buffer);
    let offset = 0;
    view.setUint32(offset, tx.body.length, true);
    offset += 4;
    txBytes.set(tx.body, offset);
    offset += tx.body.length;
    view.setUint16(offset, signature.length, true);
    offset += 2;
    txBytes.set(signature, offset);
    offset += signature.length;
    view.setUint16(offset, publicKey.length, true);
    offset += 2;
    txBytes.set(publicKey, offset);
    return txBytes;
  }
}

interface UnsignedTransaction {
  body: Uint8Array;
  chainId: string;
  agentAddress: string;
  nonce: number;
}

2.2 安全存储策略

Agent 私钥属于最高敏感级别的资产。以下是推荐的存储层次结构:

// secure-storage.ts

enum SecurityLevel {
  CRITICAL = 'critical',
  HIGH = 'high',
  MEDIUM = 'medium',
  LOW = 'low',
}

interface KeyStorageConfig {
  level: SecurityLevel;
  backupCount: number;
  autoRotation: boolean;
  hsmSlotId?: number;
  teeEnclaveId?: string;
}

abstract class KeyStorageProvider {
  abstract store(keyId: string, keyData: Uint8Array): Promise<void>;
  abstract load(keyId: string): Promise<Uint8Array>;
  abstract delete(keyId: string): Promise<void>;
  abstract list(): Promise<string[]>;
  abstract rotate(keyId: string): Promise<Dilithium5KeyPair>;
}

class HSMKeyStorage extends KeyStorageProvider {
  private hsmSession: HsmSession;

  constructor(hsmEndpoint: string, credentials: HsmCredentials) {
    super();
  }

  async store(keyId: string, keyData: Uint8Array): Promise<void> {
    const template = [
      { type: 'CKA_CLASS', value: 'CKO_PRIVATE_KEY' },
      { type: 'CKA_KEY_TYPE', value: 'CKK_DILITHIUM5' },
      { type: 'CKA_TOKEN', value: true },
      { type: 'CKA_PRIVATE', value: true },
      { type: 'CKA_SIGN', value: true },
      { type: 'CKA_EXTRACTABLE', value: false },
      { type: 'CKA_LABEL', value: keyId },
      { type: 'CKA_ID', value: keyId },
    ];
    await this.hsmSession.importKey(template, keyData);
    await this.hsmSession.finalizeKeyImport(keyId);
  }

  async signWithHSM(keyId: string, message: Uint8Array): Promise<Uint8Array> {
    const session = await this.hsmSession.open();
    try {
      const keyHandle = await session.findObject([
        { type: 'CKA_LABEL', value: keyId },
      ]);
      return await session.sign(keyHandle, 'CKM_DILITHIUM5', message);
    } finally {
      await session.close();
    }
  }
}

class TEEKeyStorage extends KeyStorageProvider {
  private enclave: TeeEnclave;

  constructor(enclaveConfig: TeeConfig) {
    super();
  }

  async store(keyId: string, keyData: Uint8Array): Promise<void> {
    await this.enclave.run((sealed: SealedData) => {
      const encrypted = this.sealData(keyData);
      return encrypted;
    }, { keyId });
  }

  async remoteAttestation(): Promise<AttestationEvidence> {
    return this.enclave.getEvidence();
  }
}

class EncryptedFileStorage extends KeyStorageProvider {
  private encryptionKey: Uint8Array;

  constructor(encryptionKey: Uint8Array) {
    super();
    this.encryptionKey = encryptionKey;
  }

  async store(keyId: string, keyData: Uint8Array): Promise<void> {
    const iv = crypto.getRandomValues(new Uint8Array(12));
    const encrypted = await this.aesEncrypt(keyData, iv);
    const record = {
      id: keyId, iv: Array.from(iv), data: Array.from(encrypted),
      createdAt: Date.now(),
    };
    const filePath = `/opt/msg-agent/keys/${keyId}.enc`;
    await fs.writeFile(filePath, JSON.stringify(record));
    await fs.chmod(filePath, 0o600);
  }

  private async aesEncrypt(data: Uint8Array, iv: Uint8Array): Promise<Uint8Array> {
    const key = await crypto.subtle.importKey(
      'raw', this.encryptionKey,
      { name: 'AES-GCM' }, false, ['encrypt']
    );
    const encrypted = await crypto.subtle.encrypt(
      { name: 'AES-GCM', iv }, key, data
    );
    return new Uint8Array(encrypted);
  }
}

function createKeyStorage(level: SecurityLevel): KeyStorageProvider {
  switch (level) {
    case SecurityLevel.CRITICAL:
      return new HSMKeyStorage(
        process.env.MSG_HSM_ENDPOINT!,
        { slotId: parseInt(process.env.MSG_HSM_SLOT_ID!), pin: process.env.MSG_HSM_PIN! }
      );
    case SecurityLevel.HIGH:
      const encKey = deriveKeyFromPassword(process.env.MSG_KEY_ENC_PASSWORD!);
      return new EncryptedFileStorage(encKey);
    default:
      throw new Error(`Unsupported security level: ${level}`);
  }
}

2.3 密钥轮换自动化

密钥轮换是降低密钥泄露风险的核心措施。Agent 密钥应定期轮换,并在检测到异常时立即触发紧急轮换。

// key-rotation.ts
import { schedule } from 'node-cron';

interface KeyRotationConfig {
  agentId: string;
  rotationIntervalDays: number;
  storageProvider: KeyStorageProvider;
  chainClient: MsgChainClient;
  preRotationHook?: () => Promise<void>;
  postRotationHook?: (newAddress: string) => Promise<void>;
  emergencyRotation?: boolean;
}

class KeyRotationScheduler {
  private configs: Map<string, KeyRotationConfig> = new Map();
  private runningJobs: Map<string, CronJob> = new Map();

  registerAgent(config: KeyRotationConfig): void {
    this.configs.set(config.agentId, config);
    const cronExpression = '0 3 * * *';
    const job = schedule(cronExpression, async () => {
      await this.checkAndRotate(config.agentId);
    });
    this.runningJobs.set(config.agentId, job);
  }

  async checkAndRotate(agentId: string): Promise<void> {
    const config = this.configs.get(agentId);
    if (!config) throw new Error(`Agent ${agentId} not registered`);
    const currentKeys = await config.storageProvider.load(`${agentId}_meta`);
    const meta = JSON.parse(new TextDecoder().decode(currentKeys)) as KeyMeta;
    const elapsed = Date.now() - meta.lastRotation;
    const threshold = config.rotationIntervalDays * 86400 * 1000;
    if (elapsed < threshold) {
      console.log(`[KeyRotation] Agent ${agentId} key still valid, next rotation in ${Math.ceil((threshold - elapsed) / 86400000)} days`);
      return;
    }
    await this.rotateKey(agentId, config);
  }

  async rotateKey(agentId: string, config: KeyRotationConfig): Promise<string> {
    console.log(`[KeyRotation] Starting key rotation for agent ${agentId}`);
    if (config.preRotationHook) await config.preRotationHook();
    const cryptoProvider = new AgentCryptoProvider();
    const newKeyPair = await cryptoProvider.generateKeyPair();
    const newAddress = cryptoProvider.deriveAddress(newKeyPair.publicKey);
    await config.storageProvider.store(`${agentId}_private`, newKeyPair.privateKey);
    await config.storageProvider.store(`${agentId}_public`, newKeyPair.publicKey);
    const updateTx = await config.chainClient.updateAgentKey(agentId, newAddress, newKeyPair.publicKey);
    await config.chainClient.broadcastTx(updateTx);
    const receipt = await config.chainClient.waitForTx(updateTx.hash);
    if (!receipt.success) throw new Error(`Key rotation tx failed: ${receipt.error}`);
    const meta: KeyMeta = {
      agentId, currentAddress: newAddress, lastRotation: Date.now(),
      previousAddress: meta?.currentAddress ?? '', rotationTxHash: updateTx.hash,
    };
    await config.storageProvider.store(`${agentId}_meta`, new TextEncoder().encode(JSON.stringify(meta)));
    if (config.postRotationHook) await config.postRotationHook(newAddress);
    console.log(`[KeyRotation] Agent ${agentId} key rotated to ${newAddress}`);
    return newAddress;
  }

  async emergencyRotate(agentId: string): Promise<void> {
    const config = this.configs.get(agentId);
    if (!config) throw new Error(`Agent ${agentId} not registered`);
    console.warn(`[KeyRotation] EMERGENCY rotation for agent ${agentId}`);
    await this.rotateKey(agentId, { ...config, rotationIntervalDays: 0, emergencyRotation: true });
  }
}

interface KeyMeta {
  agentId: string;
  currentAddress: string;
  lastRotation: number;
  previousAddress: string;
  rotationTxHash: string;
}

// 完整的密钥轮换 Cron Job 示例
async function setupKeyRotationCron(): Promise<KeyRotationScheduler> {
  const storageProvider = createKeyStorage(SecurityLevel.HIGH);
  const chainClient = new MsgChainClient({
    rpcEndpoint: process.env.MSG_RPC_ENDPOINT!,
    chainId: 'msg-chain-1',
  });
  const scheduler = new KeyRotationScheduler();
  scheduler.registerAgent({
    agentId: 'agent-main-trader-01',
    rotationIntervalDays: 30,
    storageProvider, chainClient,
    preRotationHook: async () => {
      await chainClient.pauseAgent('agent-main-trader-01');
      console.log('[PreRotation] Agent paused for key rotation');
    },
    postRotationHook: async (newAddress) => {
      await updateAgentRegistry('agent-main-trader-01', newAddress);
      await chainClient.resumeAgent('agent-main-trader-01');
      console.log(`[PostRotation] Agent resumed with new key ${newAddress}`);
    },
  });
  scheduler.registerAgent({
    agentId: 'agent-crosschain-bridge',
    rotationIntervalDays: 14,
    storageProvider, chainClient,
    postRotationHook: async (newAddress) => {
      await Promise.all([updateCrosschainRegistry(newAddress), syncBridgeWhiteList(newAddress)]);
    },
  });
  console.log('[KeyRotation] Cron jobs initialized');
  return scheduler;
}

async function emergencyKeyRotationCLI(): Promise<void> {
  const agentId = process.argv[2];
  if (!agentId) { console.error('Usage: ts-node emergency-rotate.ts <agent-id>'); process.exit(1); }
  const scheduler = await setupKeyRotationCron();
  await scheduler.emergencyRotate(agentId);
  console.log(`Emergency rotation completed for ${agentId}`);
}

2.4 备份与恢复

密钥备份是 Agent 容灾的核心能力。备份策略必须平衡安全性和可用性。

// key-backup.ts

interface BackupConfig {
  agentId: string;
  storageProvider: KeyStorageProvider;
  backupProviders: BackupProvider[];
  shardCount: number;
  threshold: number;
  scheduleCron: string;
  retentionDays: number;
}

interface BackupProvider {
  name: string;
  upload(data: Uint8Array): Promise<string>;
  download(backupId: string): Promise<Uint8Array>;
  delete(backupId: string): Promise<void>;
}

class ShamirSharding {
  private readonly prime: bigint;

  constructor(bits: number = 256) {
    this.prime = BigInt('0xFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFEFFFFFC2F');
  }

  split(secret: Uint8Array, n: number, k: number): Uint8Array[] {
    const secretBig = this.bytesToBigInt(secret);
    const shares: Uint8Array[] = [];
    const coefficients: bigint[] = [];
    for (let i = 0; i < k - 1; i++) coefficients.push(this.randomBigInt());
    for (let i = 1; i <= n; i++) {
      const x = BigInt(i);
      let y = secretBig;
      for (let j = 0; j < k - 1; j++) {
        y = (y + coefficients[j] * this.modPow(x, BigInt(j + 1))) % this.prime;
      }
      const share = new Uint8Array(33);
      share[0] = i;
      const yBytes = this.bigIntToBytes(y, 32);
      share.set(yBytes, 1);
      shares.push(share);
    }
    return shares;
  }

  recover(shares: Uint8Array[]): Uint8Array {
    if (shares.length < 2) throw new Error('Need at least 2 shares');
    const k = shares.length;
    const xs: bigint[] = []; const ys: bigint[] = [];
    for (const share of shares) {
      xs.push(BigInt(share[0]));
      ys.push(this.bytesToBigInt(share.slice(1)));
    }
    let secret = BigInt(0);
    for (let i = 0; i < k; i++) {
      let numerator = BigInt(1); let denominator = BigInt(1);
      for (let j = 0; j < k; j++) {
        if (i === j) continue;
        numerator = (numerator * (-xs[j])) % this.prime;
        denominator = (denominator * (xs[i] - xs[j])) % this.prime;
      }
      const li = (numerator * this.modInverse(denominator)) % this.prime;
      secret = (secret + ys[i] * li) % this.prime;
    }
    if (secret < 0) secret += this.prime;
    return this.bigIntToBytes(secret, 32);
  }

  private modPow(base: bigint, exp: bigint): bigint {
    let result = BigInt(1);
    base = ((base % this.prime) + this.prime) % this.prime;
    while (exp > 0) {
      if (exp & BigInt(1)) result = (result * base) % this.prime;
      base = (base * base) % this.prime;
      exp >>= BigInt(1);
    }
    return result;
  }

  private modInverse(a: bigint): bigint {
    a = ((a % this.prime) + this.prime) % this.prime;
    let [oldR, r] = [a, this.prime];
    let [oldS, s] = [BigInt(1), BigInt(0)];
    while (r !== BigInt(0)) {
      const quotient = oldR / r;
      [oldR, r] = [r, oldR - quotient * r];
      [oldS, s] = [s, oldS - quotient * s];
    }
    return ((oldS % this.prime) + this.prime) % this.prime;
  }

  private bytesToBigInt(bytes: Uint8Array): bigint {
    let result = BigInt(0);
    for (const byte of bytes) result = (result << BigInt(8)) | BigInt(byte);
    return result;
  }

  private bigIntToBytes(n: bigint, length: number): Uint8Array {
    const bytes = new Uint8Array(length);
    for (let i = length - 1; i >= 0; i--) {
      bytes[i] = Number(n & BigInt(0xFF));
      n >>= BigInt(8);
    }
    return bytes;
  }

  private randomBigInt(): bigint {
    const bytes = new Uint8Array(32);
    crypto.getRandomValues(bytes);
    return this.bytesToBigInt(bytes) % this.prime;
  }
}

class KeyBackupManager {
  private config: BackupConfig;
  private sharding: ShamirSharding;

  constructor(config: BackupConfig) {
    this.config = config;
    this.sharding = new ShamirSharding();
  }

  async performBackup(): Promise<void> {
    const agentId = this.config.agentId;
    console.log(`[Backup] Starting key backup for ${agentId}`);
    const privateKey = await this.config.storageProvider.load(`${agentId}_private`);
    const shares = this.sharding.split(privateKey, this.config.shardCount, this.config.threshold);
    const backupIds: string[] = [];
    const providerCount = this.config.backupProviders.length;
    for (let i = 0; i < shares.length; i++) {
      const provider = this.config.backupProviders[i % providerCount];
      const shareData = new Uint8Array([
        ...new TextEncoder().encode(JSON.stringify({
          version: 1, agentId, index: i, total: shares.length,
          threshold: this.config.threshold, timestamp: Date.now(),
        })),
        0, ...shares[i],
      ]);
      const backupId = await provider.upload(shareData);
      backupIds.push(backupId);
      console.log(`[Backup] Share ${i + 1}/${shares.length} uploaded to ${provider.name} (backup ID: ${backupId.slice(0, 16)}...)`);
    }
    const metadata = {
      agentId, backupIds, timestamp: Date.now(),
      shardCount: this.config.shardCount, threshold: this.config.threshold,
      providerNames: this.config.backupProviders.map(p => p.name),
    };
    await this.config.storageProvider.store(`${agentId}_backup_meta`, new TextEncoder().encode(JSON.stringify(metadata)));
    console.log(`[Backup] Backup completed for ${agentId}: ${shares.length} shares across ${providerCount} providers`);
  }

  async restoreFromBackup(): Promise<Dilithium5KeyPair> {
    const agentId = this.config.agentId;
    const metaData = await this.config.storageProvider.load(`${agentId}_backup_meta`);
    const metadata = JSON.parse(new TextDecoder().decode(metaData));
    const shares: Uint8Array[] = [];
    const providerMap = new Map(this.config.backupProviders.map(p => [p.name, p]));
    for (let i = 0; i < metadata.backupIds.length; i++) {
      const providerName = metadata.providerNames[i % metadata.providerNames.length];
      const provider = providerMap.get(providerName);
      if (!provider) continue;
      try {
        const data = await provider.download(metadata.backupIds[i]);
        const separatorIdx = data.indexOf(0);
        const shareData = data.slice(separatorIdx + 1);
        shares.push(shareData);
        if (shares.length >= metadata.threshold) break;
      } catch (err) {
        console.warn(`[Backup] Failed to download share ${i} from ${providerName}: ${err}`);
      }
    }
    if (shares.length < metadata.threshold) throw new Error(`Insufficient shares: ${shares.length}/${metadata.threshold}`);
    const privateKey = this.sharding.recover(shares);
    const cryptoProvider = new AgentCryptoProvider();
    const publicKey = await this.derivePublicKey(privateKey);
    const address = cryptoProvider.deriveAddress(publicKey);
    console.log(`[Backup] Key restored for ${agentId}: ${address}`);
    return { publicKey, privateKey };
  }

  private async derivePublicKey(privateKey: Uint8Array): Promise<Uint8Array> {
    throw new Error('Not implemented');
  }
}

class S3BackupProvider implements BackupProvider {
  name: string;
  private bucket: string;
  private prefix: string;
  private s3Client: S3Client;

  constructor(name: string, bucket: string, prefix: string) {
    this.name = name; this.bucket = bucket; this.prefix = prefix;
    this.s3Client = new S3Client({
      region: process.env.MSG_AWS_REGION!,
      credentials: {
        accessKeyId: process.env.MSG_AWS_ACCESS_KEY_ID!,
        secretAccessKey: process.env.MSG_AWS_SECRET_ACCESS_KEY!,
      },
    });
  }

  async upload(data: Uint8Array): Promise<string> {
    const key = `${this.prefix}/${Date.now()}.share`;
    await this.s3Client.send(new PutObjectCommand({
      Bucket: this.bucket, Key: key, Body: data,
      ServerSideEncryption: 'AES256', StorageClass: 'STANDARD_IA',
    }));
    return key;
  }

  async download(backupId: string): Promise<Uint8Array> {
    const response = await this.s3Client.send(new GetObjectCommand({ Bucket: this.bucket, Key: backupId }));
    const stream = response.Body as Readable;
    const chunks: Buffer[] = [];
    for await (const chunk of stream) chunks.push(chunk);
    return new Uint8Array(Buffer.concat(chunks));
  }

  async delete(backupId: string): Promise<void> {
    await this.s3Client.send(new DeleteObjectCommand({ Bucket: this.bucket, Key: backupId }));
  }
}

function initBackupScheduler(): void {
  const config: BackupConfig = {
    agentId: process.env.MSG_AGENT_ID!,
    storageProvider: createKeyStorage(SecurityLevel.HIGH),
    backupProviders: [
      new S3BackupProvider('aws-s3', 'msg-agent-backup', 'keys'),
    ],
    shardCount: 5, threshold: 3,
    scheduleCron: '0 4 * * *', retentionDays: 90,
  };
  const manager = new KeyBackupManager(config);
  schedule(config.scheduleCron, async () => {
    try { await manager.performBackup(); }
    catch (err) {
      console.error(`[Backup] Backup failed: ${err}`);
      await alertOperator({ level: 'critical', message: `Agent key backup failed: ${err}`, agentId: config.agentId });
    }
  });
  console.log('[Backup] Backup scheduler initialized');
}

2.5 密钥安全最佳实践

+---------------------------------------------------------+
|                Agent 密钥安全清单                         |
+---------------------------------------------------------+
| [ ] 使用 HSM 或 TEE 存储生产环境私钥                     |
| [ ] 密钥轮换周期不超过 30 天                             |
| [ ] 备份分片至少 5 份,恢复阈值至少 3 份                 |
| [ ] 备份提供商至少 2 个不同云厂商                        |
| [ ] 密钥文件权限设为 600                                 |
| [ ] 禁止将私钥写入日志、环境变量或版本控制                |
| [ ] 密钥生成使用硬件熵源(如 /dev/random)               |
| [ ] 紧急轮换脚本保持可执行状态,定期演练                  |
| [ ] 轮换后旧密钥在 7 天后彻底销毁                        |
| [ ] 密钥访问记录审计日志                                 |
+---------------------------------------------------------+

3. 访问控制

3.1 RBAC 权限模型

Agent API Gateway 采用基于角色的访问控制(RBAC),确保不同角色对 Agent 的操作权限具有清晰的边界。

# agent-rbac.yaml
access_control:
  roles:
    - name: agent_owner
      description: 拥有 Agent 完整控制权
      permissions:
        - create
        - update
        - delete
        - configure
        - upgrade
        - transfer_ownership
        - manage_keys
      constraints:
        require_multisig: true
        multisig_threshold: 2
        allowed_sources: ["msg1...owner", "msg1...backup"]

    - name: agent_operator
      description: 日常运维操作
      permissions:
        - pause
        - resume
        - view_logs
        - view_config
        - trigger_backup
        - update_webhook
      constraints:
        require_multisig: false
        rate_limit:
          operations_per_hour: 50

    - name: agent_monitor
      description: 监控与告警
      permissions:
        - view_metrics
        - view_logs
        - view_alerts
        - acknowledge_alert
      constraints:
        rate_limit:
          operations_per_hour: 200

    - name: agent_auditor
      description: 审计审查
      permissions:
        - view_logs
        - view_audit_trail
        - generate_report
        - view_config
      constraints:
        read_only: true

    - name: agent_emergency
      description: 应急响应角色
      permissions:
        - emergency_pause
        - emergency_rotate
        - emergency_kill
        - view_logs
      constraints:
        require_multisig: true
        multisig_threshold: 3
        allowed_sources:
          - "msg1...emergency1"
          - "msg1...emergency2"
          - "msg1...emergency3"
        temp_grant_hours: 2
// rbac-enforcer.ts

interface Permission { action: string; resource: string; }

interface RoleDefinition {
  name: string;
  permissions: string[];
  constraints: RoleConstraints;
}

interface RoleConstraints {
  requireMultisig?: boolean;
  multisigThreshold?: number;
  allowedSources?: string[];
  rateLimit?: { operationsPerHour: number; };
  readOnly?: boolean;
  tempGrantHours?: number;
}

interface AccessRequest {
  user: UserIdentity;
  action: string;
  resource: string;
  context: RequestContext;
}

interface UserIdentity {
  did: string;
  address: string;
  roles: string[];
  sessionId: string;
}

interface RequestContext {
  ip: string;
  timestamp: number;
  requestId: string;
  userAgent: string;
}

class RBACEnforcer {
  private roles: Map<string, RoleDefinition>;
  private rateLimiters: Map<string, TokenBucket>;
  private sessionCache: Map<string, UserIdentity>;

  constructor(rbacConfig: any) {
    this.roles = new Map(rbacConfig.roles.map((r: any) => [r.name, r]));
    this.rateLimiters = new Map();
    this.sessionCache = new Map();
  }

  async authorize(request: AccessRequest): Promise<AuthDecision> {
    const { user, action, resource, context } = request;
    const userPermissions = new Set<string>();
    const constraints: RoleConstraints[] = [];

    for (const roleName of user.roles) {
      const role = this.roles.get(roleName);
      if (!role) continue;
      role.permissions.forEach(p => userPermissions.add(p));
      constraints.push(role.constraints);
    }

    const requiredPermission = `${action}:${resource}`;
    if (!userPermissions.has(requiredPermission)) {
      return AuthDecision.deny(`Missing permission: ${requiredPermission}`, context.requestId);
    }

    for (const constraint of constraints) {
      if (constraint.requireMultisig) {
        const isMultisigApproved = await this.checkMultisig(user, action, constraint.multisigThreshold!);
        if (!isMultisigApproved) return AuthDecision.pending('Multisig approval required', context.requestId);
      }
      if (constraint.allowedSources?.length) {
        if (!constraint.allowedSources.includes(user.address)) {
          return AuthDecision.deny(`Address ${user.address} not in allowed sources`, context.requestId);
        }
      }
      if (constraint.readOnly && !this.isReadAction(action)) {
        return AuthDecision.deny(`Read-only role cannot perform ${action}`, context.requestId);
      }
      if (constraint.rateLimit) {
        const limiter = this.getRateLimiter(user.address, constraint.rateLimit.operationsPerHour);
        if (!limiter.tryConsume()) return AuthDecision.deny('Rate limit exceeded', context.requestId);
      }
    }

    await this.logAuthorization(request, AuthDecision.GRANTED);
    return AuthDecision.grant(context.requestId);
  }

  private async checkMultisig(user: UserIdentity, action: string, threshold: number): Promise<boolean> {
    try {
      const multisigClient = new MsgMultisigClient(process.env.MSG_MULTISIG_CONTRACT!);
      const approvals = await multisigClient.getApprovalCount(user.address, action);
      return approvals >= threshold;
    } catch (err) {
      console.error(`[RBAC] Multisig check failed: ${err}`);
      return false;
    }
  }

  private getRateLimiter(key: string, maxOperationsPerHour: number): TokenBucket {
    if (!this.rateLimiters.has(key)) {
      this.rateLimiters.set(key, new TokenBucket({
        capacity: maxOperationsPerHour,
        refillRate: maxOperationsPerHour / 3600,
        refillInterval: 1000,
      }));
    }
    return this.rateLimiters.get(key)!;
  }

  private isReadAction(action: string): boolean {
    return action.startsWith('view_') || action.startsWith('list_');
  }

  private async logAuthorization(request: AccessRequest, decision: AuthDecision): Promise<void> {
    const logEntry = {
      timestamp: new Date().toISOString(),
      requestId: request.context.requestId,
      user: request.user.address,
      roles: request.user.roles,
      action: request.action,
      resource: request.resource,
      decision: decision.status,
      reason: decision.reason,
    };
    await auditLogger.log('authorization', logEntry);
  }
}

class AuthDecision {
  readonly status: 'granted' | 'denied' | 'pending';
  readonly reason: string;
  readonly requestId: string;
  readonly timestamp: number;

  private constructor(status: 'granted' | 'denied' | 'pending', reason: string, requestId: string) {
    this.status = status; this.reason = reason; this.requestId = requestId; this.timestamp = Date.now();
  }

  static grant(requestId: string) { return new AuthDecision('granted', 'Approved', requestId); }
  static deny(reason: string, requestId: string) { return new AuthDecision('denied', reason, requestId); }
  static pending(reason: string, requestId: string) { return new AuthDecision('pending', reason, requestId); }
}

class TokenBucket {
  private capacity: number;
  private tokens: number;
  private refillRate: number;
  private refillInterval: number;
  private lastRefill: number;

  constructor(config: { capacity: number; refillRate: number; refillInterval: number }) {
    this.capacity = config.capacity;
    this.tokens = config.capacity;
    this.refillRate = config.refillRate;
    this.refillInterval = config.refillInterval;
    this.lastRefill = Date.now();
  }

  tryConsume(count: number = 1): boolean {
    this.refill();
    if (this.tokens >= count) { this.tokens -= count; return true; }
    return false;
  }

  private refill(): void {
    const now = Date.now();
    const elapsed = now - this.lastRefill;
    const tokensToAdd = (elapsed / this.refillInterval) * this.refillRate;
    this.tokens = Math.min(this.capacity, this.tokens + tokensToAdd);
    this.lastRefill = now;
  }
}

3.2 DID 身份认证

每个 Agent 和操作员都拥有 MSG Chain 上的去中心化身份(DID),实现无中心化信任的身份认证。

// did-auth.ts

interface DIDSchema {
  context: string[];
  id: string;
  verificationMethod: VerificationMethod[];
  authentication: string[];
  assertionMethod: string[];
  service: DIDService[];
  created: string;
  updated: string;
}

interface VerificationMethod {
  id: string;
  type: 'Dilithium5VerificationKey2026';
  controller: string;
  publicKeyBase64: string;
}

class DIDAuthProvider {
  private chainClient: MsgChainClient;
  private didRegistry: string;

  constructor(chainClient: MsgChainClient) {
    this.chainClient = chainClient;
    this.didRegistry = chainClient.getContractAddress('did_registry_v1');
  }

  async resolveDID(did: string): Promise<DIDSchema | null> {
    const result = await this.chainClient.queryContract(this.didRegistry, { resolve_did: { did } });
    return result.document ?? null;
  }

  async authenticate(did: string): Promise<AuthProof> {
    const document = await this.resolveDID(did);
    if (!document) throw new Error(`DID not found: ${did}`);
    const challenge = crypto.getRandomValues(new Uint8Array(32));
    const challengeHex = Buffer.from(challenge).toString('hex');
    return { did, challenge: challengeHex, verificationMethod: document.verificationMethod[0].id, expiresAt: Date.now() + 300000 };
  }

  async verifyChallenge(proof: AuthProof, signature: Uint8Array, publicKey: Uint8Array): Promise<boolean> {
    const cryptoProvider = new AgentCryptoProvider();
    const message = new TextEncoder().encode(`${proof.challenge}|${proof.did}|${proof.expiresAt}`);
    return cryptoProvider.verifySignature(publicKey, message, signature);
  }

  generateToken(user: UserIdentity): string {
    const payload = {
      sub: user.did, address: user.address, roles: user.roles,
      iat: Math.floor(Date.now() / 1000), exp: Math.floor(Date.now() / 1000) + 3600,
      jti: crypto.randomUUID(),
    };
    return jwt.sign(payload, process.env.MSG_JWT_SECRET!, { algorithm: 'HS256' });
  }
}

interface AuthProof {
  did: string;
  challenge: string;
  verificationMethod: string;
  expiresAt: number;
}

3.3 多签关键操作

// multisig-manager.ts

interface MultisigProposal {
  id: string;
  agentId: string;
  operation: string;
  params: Record<string, unknown>;
  proposer: string;
  requiredApprovals: number;
  approvals: MultisigApproval[];
  status: 'pending' | 'approved' | 'executed' | 'rejected';
  expiresAt: number;
  createdAt: number;
}

interface MultisigApproval {
  signer: string;
  signature: Uint8Array;
  timestamp: number;
}

class MultisigManager {
  private chainClient: MsgChainClient;
  private multisigContract: string;

  constructor(chainClient: MsgChainClient) {
    this.chainClient = chainClient;
    this.multisigContract = chainClient.getContractAddress('multisig_wallet_v1');
  }

  async createProposal(agentId: string, operation: string, params: Record<string, unknown>, requiredApprovals: number): Promise<MultisigProposal> {
    const result = await this.chainClient.executeContract(this.multisigContract, {
      create_proposal: { agent_id: agentId, operation, params, required_approvals: requiredApprovals, expires_in: 86400 },
    });
    return result.proposal;
  }

  async signProposal(proposalId: string, signerKeyPair: Dilithium5KeyPair): Promise<void> {
    const proposal = await this.getProposal(proposalId);
    const message = new TextEncoder().encode(JSON.stringify({
      action: 'approve_proposal', proposalId, agentId: proposal.agentId,
      operation: proposal.operation, timestamp: Date.now(),
    }));
    const cryptoProvider = new AgentCryptoProvider();
    const signature = await cryptoProvider.signMessage(signerKeyPair.privateKey, message);
    await this.chainClient.executeContract(this.multisigContract, {
      approve_proposal: { proposal_id: proposalId, signature: Array.from(signature) },
    });
  }

  async executeProposal(proposalId: string): Promise<void> {
    await this.chainClient.executeContract(this.multisigContract, { execute_proposal: { proposal_id: proposalId } });
  }

  async getProposal(proposalId: string): Promise<MultisigProposal> {
    return this.chainClient.queryContract(this.multisigContract, { get_proposal: { proposal_id: proposalId } });
  }
}

3.4 访问控制最佳实践

+---------------------------------------------------------+
|                Agent 访问控制清单                          |
+---------------------------------------------------------+
| [ ] 至少定义 4 个角色:Owner / Operator / Monitor / Auditor |
| [ ] 高风险操作需要至少 2/3 多签批准                        |
| [ ] 所有 API 端点通过 RBAC 中间件保护                      |
| [ ] 每个用户配速率限制,默认 100 次/小时                   |
| [ ] DID 认证挑战有效期不超过 5 分钟                        |
| [ ] JWT 令牌有效期不超过 1 小时,使用 RS256 签名           |
| [ ] 禁止永久 API Key,最长有效期 90 天                     |
| [ ] 异常地点登录触发额外验证                               |
| [ ] 角色变更需记录审计日志并通知 Owner                     |
| [ ] 定期审查(每季度)角色分配和权限                       |
+---------------------------------------------------------+

4. 威胁监控

4.1 异常检测规则

Agent 行为监控是安全运维的核心环节。通过定义 Prometheus 告警规则,实时检测异常行为模式。

# threat_detection_rules.py

from dataclasses import dataclass, field
from enum import Enum
import time
import json


class Severity(Enum):
    INFO = "info"
    WARNING = "warning"
    CRITICAL = "critical"


class AgentAction(Enum):
    PAUSE = "pause_agent"
    ALERT = "alert_operator"
    KILL = "kill_agent"
    LOG = "log_only"
    ESCALATE = "escalate"


@dataclass
class ThreatRule:
    name: str
    condition: str
    severity: Severity
    action: AgentAction
    cooldown_seconds: int = 300
    last_fired: float = 0.0
    description: str = ""


@dataclass
class ThreatEvent:
    rule_name: str
    agent_id: str
    severity: Severity
    timestamp: float
    metrics: dict
    action_taken: AgentAction


THREAT_RULES = [
    ThreatRule(
        name="unexpected_constitution_violation",
        condition="rate(constitution_violations_total[5m]) > 10",
        severity=Severity.CRITICAL,
        action=AgentAction.PAUSE,
        description="Agent 在短时间内多次违反宪法规定",
    ),
    ThreatRule(
        name="abnormal_gas_spike",
        condition="rate(agent_gas_used[5m]) > 1000000",
        severity=Severity.WARNING,
        action=AgentAction.ALERT,
        description="Agent Gas 消耗异常激增,可能被滥用",
    ),
    ThreatRule(
        name="unauthorized_contract_call",
        condition="rate(unauthorized_calls_total[10m]) > 5",
        severity=Severity.CRITICAL,
        action=AgentAction.PAUSE,
        description="Agent 尝试调用白名单外的合约",
    ),
    ThreatRule(
        name="abnormal_tx_frequency",
        condition="rate(agent_tx_total[1m]) > 60",
        severity=Severity.WARNING,
        action=AgentAction.ALERT,
        description="Agent 交易频率异常升高,可能遭受重放攻击",
    ),
    ThreatRule(
        name="large_value_transfer",
        condition="agent_tx_value_max[1h] > 1000000000",
        severity=Severity.CRITICAL,
        action=AgentAction.PAUSE,
        description="Agent 单笔转账超过 1000 MSG,需人工确认",
    ),
    ThreatRule(
        name="a2a_message_anomaly",
        condition="rate(a2a_messages_total[5m]) > 200",
        severity=Severity.WARNING,
        action=AgentAction.ESCALATE,
        description="A2A 消息量异常,可能存在 Agent 间机器人攻击",
    ),
    ThreatRule(
        name="constitution_drift",
        condition="constitution_embedding_distance > 0.95",
        severity=Severity.CRITICAL,
        action=AgentAction.KILL,
        description="Agent 行为嵌入向量偏离宪法基线超过 95%",
    ),
    ThreatRule(
        name="failed_auth_attempts",
        condition="rate(agent_auth_failures_total[5m]) > 20",
        severity=Severity.WARNING,
        action=AgentAction.ALERT,
        description="Agent 认证失败次数激增,可能存在密钥暴力破解",
    ),
    ThreatRule(
        name="payment_flow_abnormal",
        condition="rate(payment_flow_errors_total[5m]) > 5",
        severity=Severity.CRITICAL,
        action=AgentAction.PAUSE,
        description="Agent 支付流程出现异常错误",
    ),
    ThreatRule(
        name="agent_downtime",
        condition="time() - agent_last_heartbeat > 300",
        severity=Severity.CRITICAL,
        action=AgentAction.ALERT,
        description="Agent 心跳超时,可能已宕机或遭攻击",
    ),
]


class PrometheusQueryEngine:
    def __init__(self, metrics_store: dict):
        self.metrics = metrics_store

    def evaluate_rate(self, metric_name: str, window: str) -> float:
        history = self.metrics.get(metric_name, [])
        if not history:
            return 0.0
        window_seconds = self._parse_window(window)
        now = time.time()
        window_start = now - window_seconds
        recent_values = [v for ts, v in history if ts >= window_start]
        return sum(recent_values) / window_seconds

    def evaluate_max(self, metric_name: str, window: str) -> float:
        history = self.metrics.get(metric_name, [])
        if not history:
            return 0.0
        window_seconds = self._parse_window(window)
        now = time.time()
        window_start = now - window_seconds
        values = [v for ts, v in history if ts >= window_start]
        return max(values) if values else 0.0

    def _parse_window(self, window: str) -> int:
        window = window.strip("[]")
        unit = window[-1]
        value = int(window[:-1])
        multipliers = {'s': 1, 'm': 60, 'h': 3600, 'd': 86400}
        return value * multipliers.get(unit, 1)


class ThreatDetector:
    def __init__(self, rules: list[ThreatRule]):
        self.rules = rules
        self.event_history: list[ThreatEvent] = []
        self.query_engine = PrometheusQueryEngine({})

    def update_metrics(self, metrics: dict):
        self.query_engine.metrics = metrics

    async def evaluate_rules(self, agent_id: str) -> list[ThreatEvent]:
        events = []
        now = time.time()
        for rule in self.rules:
            if now - rule.last_fired < rule.cooldown_seconds:
                continue
            if self._evaluate_condition(rule.condition):
                rule.last_fired = now
                metrics_snapshot = self._capture_metrics()
                event = ThreatEvent(
                    rule_name=rule.name, agent_id=agent_id,
                    severity=rule.severity, timestamp=now,
                    metrics=metrics_snapshot, action_taken=rule.action,
                )
                events.append(event)
                self.event_history.append(event)
        return events

    def _evaluate_condition(self, condition: str) -> bool:
        try:
            parts = condition.replace(" ", "").split(">")
            if len(parts) != 2:
                return False
            left, threshold_str = parts
            threshold = float(threshold_str)
            if left.startswith("rate("):
                inner = left[5:-1]
                metric_name, window = inner.split("[")
                window = "[" + window
                value = self.query_engine.evaluate_rate(metric_name, window)
            elif left.startswith("time()-"):
                metric_name = left.split("-")[1]
                last_hb = max((ts for ts, _ in self.query_engine.metrics.get(metric_name, [])), default=0)
                value = time.time() - last_hb
            else:
                metric_name = left
                value = self.query_engine.evaluate_max(metric_name, "5m")
            return value > threshold
        except Exception as e:
            print(f"[ThreatDetection] Error evaluating condition '{condition}': {e}")
            return False

    def _capture_metrics(self) -> dict:
        snapshot = {}
        for metric_name, history in self.query_engine.metrics.items():
            snapshot[metric_name] = [(ts, val) for ts, val in history[-10:]]
        return snapshot

    def get_recent_events(self, count: int = 10) -> list[ThreatEvent]:
        return self.event_history[-count:]

    def get_agent_risk_score(self, agent_id: str) -> float:
        score = 0.0
        recent = [e for e in self.event_history[-100:] if e.agent_id == agent_id]
        for event in recent:
            severity_weights = {Severity.INFO: 1, Severity.WARNING: 5, Severity.CRITICAL: 20}
            score += severity_weights.get(event.severity, 0)
        return min(100, score)

4.2 Constitution 违反监控

# constitution_monitor.py

from dataclasses import dataclass
import hashlib
import json
from typing import Optional


@dataclass
class ConstitutionDocument:
    version: str
    agent_id: str
    rules: list[dict]
    constraints: list[str]
    allowed_actions: list[str]
    forbidden_actions: list[str]
    created_at: int
    updated_at: int
    hash: str

    def compute_hash(self) -> str:
        content = json.dumps({
            'agent_id': self.agent_id,
            'rules': sorted(json.dumps(r, sort_keys=True) for r in self.rules),
            'constraints': sorted(self.constraints),
            'allowed_actions': sorted(self.allowed_actions),
            'forbidden_actions': sorted(self.forbidden_actions),
        }, sort_keys=True)
        return hashlib.sha256(content.encode()).hexdigest()


@dataclass
class AgentActionLog:
    agent_id: str
    action_type: str
    target_contract: str
    params: dict
    timestamp: int
    tx_hash: str
    approved: bool
    rejection_reason: Optional[str] = None


@dataclass
class ConstitutionViolation:
    agent_id: str
    rule_type: str
    rule_detail: str
    action: AgentActionLog
    detected_at: int


class ConstitutionViolationDetector:
    def __init__(self, constitution: ConstitutionDocument):
        self.constitution = constitution
        self.violations: list[ConstitutionViolation] = []

    def check_action(self, action: AgentActionLog) -> bool:
        for forbidden in self.constitution.forbidden_actions:
            if action.action_type.lower() == forbidden.lower():
                self.violations.append(ConstitutionViolation(
                    agent_id=action.agent_id,
                    rule_type='forbidden_action',
                    rule_detail=f"Action '{action.action_type}' is forbidden by constitution",
                    action=action, detected_at=action.timestamp,
                ))
                return False
        if self.constitution.allowed_actions:
            if action.action_type.lower() not in [a.lower() for a in self.constitution.allowed_actions]:
                self.violations.append(ConstitutionViolation(
                    agent_id=action.agent_id,
                    rule_type='unauthorized_action',
                    rule_detail=f"Action '{action.action_type}' not in allowed_actions",
                    action=action, detected_at=action.timestamp,
                ))
                return False
        for constraint in self.constitution.constraints:
            if not self._evaluate_constraint(constraint, action):
                self.violations.append(ConstitutionViolation(
                    agent_id=action.agent_id,
                    rule_type='constraint_violation',
                    rule_detail=f"Constraint '{constraint}' violated by action {action.action_type}",
                    action=action, detected_at=action.timestamp,
                ))
                return False
        return True

    def _evaluate_constraint(self, constraint: str, action: AgentActionLog) -> bool:
        try:
            if constraint.startswith("max_gas:"):
                gas_limit = float(constraint.split(":")[1])
                return action.params.get("gas", 0) <= gas_limit
            elif constraint.startswith("max_value:"):
                value_limit = float(constraint.split(":")[1])
                return action.params.get("value", 0) <= value_limit
            elif constraint.startswith("allowed_targets:"):
                allowed = eval(constraint.split(":", 1)[1])
                return action.target_contract in allowed
            return True
        except Exception as e:
            print(f"[Constitution] Error evaluating constraint '{constraint}': {e}")
            return False

4.3 Gas 与支付监控

# gas_monitor.py

from dataclasses import dataclass
from collections import deque
import time
import numpy as np
from typing import Optional


@dataclass
class GasRecord:
    agent_id: str
    tx_hash: str
    gas_used: int
    gas_price: int
    block_height: int
    timestamp: int
    operation: str


class GasAnomalyDetector:
    def __init__(self, window_size: int = 100, z_score_threshold: float = 3.0):
        self.window = deque(maxlen=window_size)
        self.z_score_threshold = z_score_threshold

    def record_gas(self, record: GasRecord):
        self.window.append(record)

    def is_anomalous(self, gas_used: int) -> tuple[bool, float]:
        if len(self.window) < 10:
            return False, 0.0
        values = [r.gas_used for r in self.window]
        mean = np.mean(values)
        std = np.std(values)
        if std == 0:
            return False, 0.0
        z_score = (gas_used - mean) / std
        return abs(z_score) > self.z_score_threshold, z_score

    def get_gas_budget_remaining(self, agent_id: str, daily_budget: int) -> tuple[int, float]:
        today = time.time()
        today_start = today - (today % 86400)
        today_gas = sum(r.gas_used for r in self.window if r.agent_id == agent_id and r.timestamp >= today_start)
        remaining = daily_budget - today_gas
        usage_pct = (today_gas / daily_budget) * 100
        return max(0, remaining), usage_pct

    def predict_daily_gas(self, agent_id: str) -> tuple[float, float]:
        if len(self.window) < 2:
            return 0.0, 0.0
        hours_elapsed = (time.time() - self.window[0].timestamp) / 3600
        if hours_elapsed < 1:
            return 0.0, 0.0
        total_gas = sum(r.gas_used for r in self.window if r.agent_id == agent_id)
        hourly_rate = total_gas / hours_elapsed
        predicted = hourly_rate * 24
        confidence = min(1.0, len(self.window) / 100)
        return predicted, confidence


@dataclass
class PaymentRecord:
    agent_id: str
    tx_hash: str
    from_addr: str
    to_addr: str
    amount: int
    token: str
    timestamp: int
    status: str


class PaymentFlowMonitor:
    def __init__(self):
        self.payments: list[PaymentRecord] = []

    def record_payment(self, payment: PaymentRecord):
        self.payments.append(payment)

    def detect_payment_anomalies(self, agent_id: str) -> list[dict]:
        anomalies = []
        recent = [p for p in self.payments[-200:] if p.agent_id == agent_id]
        if not recent:
            return anomalies
        amounts = [p.amount for p in recent]
        mean_amount = np.mean(amounts)
        std_amount = np.std(amounts)
        latest = recent[-1]
        if std_amount > 0:
            z_score = (latest.amount - mean_amount) / std_amount
            if abs(z_score) > 3:
                anomalies.append({
                    'type': 'large_payment_anomaly', 'agent_id': agent_id,
                    'tx_hash': latest.tx_hash, 'amount': latest.amount,
                    'z_score': float(z_score), 'severity': 'critical',
                })
        if len(recent) >= 10:
            time_window = recent[-1].timestamp - recent[0].timestamp
            if time_window < 60 and len(recent) >= 10:
                anomalies.append({
                    'type': 'high_frequency_payments', 'agent_id': agent_id,
                    'count': len(recent), 'time_window': time_window, 'severity': 'warning',
                })
        recent_50 = recent[-50:]
        if recent_50:
            fail_count = sum(1 for p in recent_50 if p.status == 'failed')
            fail_rate = fail_count / len(recent_50)
            if fail_rate > 0.3:
                anomalies.append({
                    'type': 'high_failure_rate', 'agent_id': agent_id,
                    'fail_rate': float(fail_rate), 'severity': 'warning',
                })
        return anomalies

4.4 A2A 消息模式分析

# a2a_monitor.py

from dataclasses import dataclass
from collections import defaultdict
import time


@dataclass
class A2AMessage:
    source_agent: str
    target_agent: str
    message_type: str
    payload_size: int
    timestamp: int
    message_hash: str


class A2AMessageAnalyzer:
    def __init__(self, window_minutes: int = 30):
        self.messages: list[A2AMessage] = []
        self.window_seconds = window_minutes * 60

    def record_message(self, msg: A2AMessage):
        self.messages.append(msg)

    def detect_message_storm(self) -> list[dict]:
        alerts = []
        now = time.time()
        recent = [m for m in self.messages[-1000:] if now - m.timestamp < self.window_seconds]
        source_counts = defaultdict(int)
        for m in recent:
            source_counts[m.source_agent] += 1
        for agent, count in source_counts.items():
            rate = count / (self.window_seconds / 60)
            if rate > 100:
                alerts.append({
                    'type': 'message_storm', 'source_agent': agent,
                    'rate_per_minute': rate, 'total_messages': count, 'severity': 'critical',
                })
        return alerts

    def detect_communication_isolation(self, agent_id: str) -> bool:
        now = time.time()
        threshold = now - self.window_seconds
        recent_to = [m for m in self.messages[-500:] if m.target_agent == agent_id and m.timestamp > threshold]
        recent_from = [m for m in self.messages[-500:] if m.source_agent == agent_id and m.timestamp > threshold]
        return len(recent_to) == 0 and len(recent_from) == 0

    def detect_suspicious_patterns(self) -> list[dict]:
        alerts = []
        recent = self.messages[-200:]
        source_targets = defaultdict(set)
        for m in recent:
            source_targets[m.source_agent].add(m.target_agent)
        for agent, targets in source_targets.items():
            if len(targets) > 10:
                alerts.append({
                    'type': 'broadcast_pattern', 'source_agent': agent,
                    'unique_targets': len(targets), 'severity': 'warning',
                })
        hash_counts = defaultdict(int)
        for m in recent:
            hash_counts[m.message_hash] += 1
        for msg_hash, count in hash_counts.items():
            if count > 3:
                source = next((m.source_agent for m in recent if m.message_hash == msg_hash), 'unknown')
                alerts.append({
                    'type': 'message_replay', 'message_hash': msg_hash.hex(),
                    'source_agent': source, 'replay_count': count, 'severity': 'critical',
                })
        return alerts

4.5 监控指标暴露

// metrics-exporter.ts

import express from 'express';
import prometheus from 'prom-client';

class AgentMetricsExporter {
  private registry: prometheus.Registry;
  private app: express.Application;
  private port: number;

  private constitutionViolations: prometheus.Counter;
  private agentGasUsed: prometheus.Histogram;
  private agentTxTotal: prometheus.Counter;
  private agentAuthFailures: prometheus.Counter;
  private a2aMessagesTotal: prometheus.Counter;
  private paymentFlowErrors: prometheus.Counter;
  private agentUptime: prometheus.Gauge;
  private agentRiskScore: prometheus.Gauge;
  private unauthorizedCalls: prometheus.Counter;
  private agentLastHeartbeat: prometheus.Gauge;

  constructor(port: number = 9090) {
    this.port = port;
    this.registry = new prometheus.Registry();
    prometheus.collectDefaultMetrics({ register: this.registry, prefix: 'msg_agent_' });

    this.constitutionViolations = new prometheus.Counter({
      name: 'constitution_violations_total', help: 'Total number of constitution violations',
      labelNames: ['agent_id', 'violation_type'], registers: [this.registry],
    });
    this.agentGasUsed = new prometheus.Histogram({
      name: 'agent_gas_used', help: 'Gas used per transaction',
      labelNames: ['agent_id', 'operation'],
      buckets: [50000, 100000, 200000, 500000, 1000000, 5000000], registers: [this.registry],
    });
    this.agentTxTotal = new prometheus.Counter({
      name: 'agent_tx_total', help: 'Total number of transactions',
      labelNames: ['agent_id', 'status'], registers: [this.registry],
    });
    this.agentAuthFailures = new prometheus.Counter({
      name: 'agent_auth_failures_total', help: 'Total number of authentication failures',
      labelNames: ['agent_id', 'failure_reason'], registers: [this.registry],
    });
    this.a2aMessagesTotal = new prometheus.Counter({
      name: 'a2a_messages_total', help: 'Total number of A2A messages',
      labelNames: ['source_agent', 'target_agent', 'message_type'], registers: [this.registry],
    });
    this.paymentFlowErrors = new prometheus.Counter({
      name: 'payment_flow_errors_total', help: 'Total number of payment flow errors',
      labelNames: ['agent_id', 'error_type'], registers: [this.registry],
    });
    this.agentUptime = new prometheus.Gauge({
      name: 'agent_uptime_seconds', help: 'Agent uptime in seconds',
      labelNames: ['agent_id'], registers: [this.registry],
    });
    this.agentRiskScore = new prometheus.Gauge({
      name: 'agent_risk_score', help: 'Current agent risk score (0-100)',
      labelNames: ['agent_id'], registers: [this.registry],
    });
    this.unauthorizedCalls = new prometheus.Counter({
      name: 'unauthorized_calls_total', help: 'Total number of unauthorized calls',
      labelNames: ['agent_id', 'target_contract'], registers: [this.registry],
    });
    this.agentLastHeartbeat = new prometheus.Gauge({
      name: 'agent_last_heartbeat', help: 'Last heartbeat timestamp',
      labelNames: ['agent_id'], registers: [this.registry],
    });

    this.setupHTTPServer();
  }

  private setupHTTPServer(): void {
    this.app = express();
    this.app.get('/metrics', async (_req, res) => {
      res.set('Content-Type', this.registry.contentType);
      const metrics = await this.registry.metrics();
      res.end(metrics);
    });
    this.app.get('/health', (_req, res) => {
      res.json({ status: 'ok', timestamp: Date.now() });
    });
    this.app.listen(this.port, () => {
      console.log(`[Metrics] Agent metrics exporter listening on :${this.port}`);
    });
  }

  recordConstitutionViolation(agentId: string, violationType: string): void {
    this.constitutionViolations.inc({ agent_id: agentId, violation_type: violationType });
  }
  recordGasUsage(agentId: string, operation: string, gasUsed: number): void {
    this.agentGasUsed.observe({ agent_id: agentId, operation }, gasUsed);
  }
  recordTransaction(agentId: string, status: string): void {
    this.agentTxTotal.inc({ agent_id: agentId, status });
  }
  recordAuthFailure(agentId: string, reason: string): void {
    this.agentAuthFailures.inc({ agent_id: agentId, failure_reason: reason });
  }
  recordA2AMessage(source: string, target: string, messageType: string): void {
    this.a2aMessagesTotal.inc({ source_agent: source, target_agent: target, message_type: messageType });
  }
  recordPaymentError(agentId: string, errorType: string): void {
    this.paymentFlowErrors.inc({ agent_id: agentId, error_type: errorType });
  }
  updateUptime(agentId: string, seconds: number): void {
    this.agentUptime.set({ agent_id: agentId }, seconds);
  }
  updateRiskScore(agentId: string, score: number): void {
    this.agentRiskScore.set({ agent_id: agentId }, score);
  }
  recordUnauthorizedCall(agentId: string, targetContract: string): void {
    this.unauthorizedCalls.inc({ agent_id: agentId, target_contract: targetContract });
  }
  heartbeat(agentId: string): void {
    this.agentLastHeartbeat.set({ agent_id: agentId }, Date.now() / 1000);
  }
}

4.6 监控最佳实践

+---------------------------------------------------------+
|                Agent 监控实施清单                          |
+---------------------------------------------------------+
| [ ] 所有 Agent 暴露 /metrics 端点(Prometheus 格式)       |
| [ ] 配置 Grafana 仪表盘,展示核心安全指标                   |
| [ ] Constitution 违反告警为最高优先级,立即处理             |
| [ ] Gas 消耗设置每日预算告警阈值(80% / 90% / 100%)       |
| [ ] A2A 消息异常模式实时分析                               |
| [ ] Agent 心跳间隔不超过 60 秒                             |
| [ ] 风险评分 > 80 自动触发应急响应                         |
| [ ] 所有告警配置 7x24 值班通知(Slack / Telegram / 钉钉)  |
| [ ] 每周生成安全指标报告                                   |
| [ ] 监控系统自身高可用部署                                 |
+---------------------------------------------------------+

5. 自动化响应

5.1 事件响应框架

# incident_responder.py

from dataclasses import dataclass
from enum import Enum
import asyncio
import time
import json


class IncidentSeverity(Enum):
    INFO = "info"
    WARNING = "warning"
    CRITICAL = "critical"


class IncidentStatus(Enum):
    DETECTED = "detected"
    INVESTIGATING = "investigating"
    CONTAINING = "containing"
    RESOLVED = "resolved"
    CLOSED = "closed"


@dataclass
class Incident:
    id: str
    agent_id: str
    rule_name: str
    severity: IncidentSeverity
    status: IncidentStatus
    detected_at: float
    resolved_at: float = None
    description: str = ""
    evidence: dict = None
    actions_taken: list[dict] = None
    playbook_name: str = ""

    def __post_init__(self):
        if self.evidence is None:
            self.evidence = {}
        if self.actions_taken is None:
            self.actions_taken = []


class AgentAPI:
    async def pause_agent(self, agent_id: str) -> dict:
        print(f"[API] Pausing agent {agent_id}")
        await asyncio.sleep(0.5)
        return {"status": "paused", "agent_id": agent_id, "timestamp": time.time()}

    async def resume_agent(self, agent_id: str) -> dict:
        print(f"[API] Resuming agent {agent_id}")
        await asyncio.sleep(0.3)
        return {"status": "resumed", "agent_id": agent_id, "timestamp": time.time()}

    async def kill_agent(self, agent_id: str) -> dict:
        print(f"[API] Killing agent {agent_id}")
        await asyncio.sleep(1.0)
        return {"status": "killed", "agent_id": agent_id, "timestamp": time.time()}

    async def rotate_agent_key(self, agent_id: str) -> dict:
        print(f"[API] Rotating key for agent {agent_id}")
        await asyncio.sleep(2.0)
        return {"status": "key_rotated", "agent_id": agent_id, "timestamp": time.time()}

    async def get_agent_status(self, agent_id: str) -> dict:
        return {"agent_id": agent_id, "status": "running", "uptime": 3600}

    async def get_agent_logs(self, agent_id: str, limit: int = 100) -> list[dict]:
        return [{"timestamp": time.time(), "message": "sample log"}]


class AlertingService:
    async def send_alert(self, message: str, severity: IncidentSeverity, channel: str = "ops"):
        print(f"[ALERT][{severity.value}] {message} (channel: {channel})")
        await asyncio.sleep(0.1)

    async def send_pagerduty(self, message: str, severity: str):
        print(f"[PAGERDUTY] {message} (severity: {severity})")


class EvidenceLogger:
    async def log_evidence(self, agent_id: str, incident_id: str, evidence_type: str, data: dict):
        record = {
            "agent_id": agent_id, "incident_id": incident_id,
            "evidence_type": evidence_type, "data": data, "timestamp": time.time(),
        }
        print(f"[EVIDENCE] Logged {evidence_type} for {agent_id}")
        await asyncio.sleep(0.1)


class CircuitBreaker:
    def __init__(self, cooldown_seconds: int = 300):
        self.cooldown = cooldown_seconds
        self.last_action: dict[str, float] = {}

    def is_open(self, agent_id: str, action: str) -> bool:
        key = f"{agent_id}:{action}"
        last_time = self.last_action.get(key, 0)
        return (time.time() - last_time) < self.cooldown

    def trip(self, agent_id: str, action: str):
        key = f"{agent_id}:{action}"
        self.last_action[key] = time.time()

    def reset(self, agent_id: str, action: str):
        key = f"{agent_id}:{action}"
        self.last_action.pop(key, None)


class IncidentResponder:
    def __init__(self):
        self.agent_api = AgentAPI()
        self.alerting = AlertingService()
        self.evidence_logger = EvidenceLogger()
        self.circuit_breaker = CircuitBreaker(cooldown_seconds=120)
        self.incidents: list[Incident] = []
        self.active_incidents: dict[str, Incident] = {}

        self.playbooks = {
            'constitution_violation': self.handle_constitution_violation,
            'key_compromise': self.handle_key_compromise,
            'payment_fraud': self.handle_payment_fraud,
            'gas_abuse': self.handle_gas_abuse,
            'a2a_attack': self.handle_a2a_attack,
            'unauthorized_call': self.handle_unauthorized_call,
            'agent_downtime': self.handle_agent_downtime,
            'default': self.handle_default,
        }

    async def handle_incident(self, threat_event: ThreatEvent) -> Incident:
        incident_id = f"inc-{int(time.time())}-{threat_event.agent_id[:8]}"
        incident = Incident(
            id=incident_id, agent_id=threat_event.agent_id,
            rule_name=threat_event.rule_name, severity=threat_event.severity,
            status=IncidentStatus.DETECTED, detected_at=threat_event.timestamp,
            description=f"Rule '{threat_event.rule_name}' triggered with severity {threat_event.severity.value}",
            evidence=threat_event.metrics,
        )
        self.incidents.append(incident)
        self.active_incidents[incident_id] = incident
        playbook = self.playbooks.get(threat_event.rule_name, self.playbooks['default'])
        incident.status = IncidentStatus.INVESTIGATING
        incident.playbook_name = playbook.__name__
        try:
            await playbook(incident)
            incident.status = IncidentStatus.RESOLVED
            incident.resolved_at = time.time()
        except Exception as e:
            incident.status = IncidentStatus.CONTAINING
            await self.escalate_to_human(incident, str(e))
        del self.active_incidents[incident_id]
        return incident

    async def handle_constitution_violation(self, incident: Incident):
        agent_id = incident.agent_id
        print(f"[Playbook] Handling constitution_violation for {agent_id}")
        logs = await self.agent_api.get_agent_logs(agent_id)
        await self.evidence_logger.log_evidence(agent_id, incident.id, 'constitution_violation', {'recent_logs': logs[-50:]})
        if not self.circuit_breaker.is_open(agent_id, 'pause'):
            await self.agent_api.pause_agent(agent_id)
            self.circuit_breaker.trip(agent_id, 'pause')
            incident.actions_taken.append({'action': 'pause_agent', 'timestamp': time.time()})
        await self.alerting.send_alert(f"[CRITICAL] Agent {agent_id} paused: constitution violation detected", IncidentSeverity.CRITICAL)

    async def handle_key_compromise(self, incident: Incident):
        agent_id = incident.agent_id
        print(f"[Playbook] Handling key_compromise for {agent_id}")
        if not self.circuit_breaker.is_open(agent_id, 'kill'):
            await self.agent_api.kill_agent(agent_id)
            self.circuit_breaker.trip(agent_id, 'kill')
            incident.actions_taken.append({'action': 'kill_agent', 'timestamp': time.time()})
        if not self.circuit_breaker.is_open(agent_id, 'rotate'):
            await self.agent_api.rotate_agent_key(agent_id)
            self.circuit_breaker.trip(agent_id, 'rotate')
            incident.actions_taken.append({'action': 'rotate_key', 'timestamp': time.time()})
        await self.alerting.send_pagerduty(f"KEY COMPROMISE: Agent {agent_id} - key rotated and agent killed", 'critical')

    async def handle_payment_fraud(self, incident: Incident):
        agent_id = incident.agent_id
        await self.agent_api.pause_agent(agent_id)
        incident.actions_taken.append({'action': 'pause_agent', 'timestamp': time.time()})
        await self.evidence_logger.log_evidence(agent_id, incident.id, 'payment_fraud', incident.evidence)
        await self.alerting.send_alert(f"[CRITICAL] Payment fraud detected for Agent {agent_id}", IncidentSeverity.CRITICAL)

    async def handle_gas_abuse(self, incident: Incident):
        agent_id = incident.agent_id
        if incident.severity == IncidentSeverity.WARNING:
            await self.alerting.send_alert(f"[WARNING] Gas abuse detected for Agent {agent_id}", IncidentSeverity.WARNING)
        elif incident.severity == IncidentSeverity.CRITICAL:
            await self.agent_api.pause_agent(agent_id)
            incident.actions_taken.append({'action': 'pause_agent', 'timestamp': time.time()})
            await self.alerting.send_alert(f"[CRITICAL] Agent {agent_id} paused for gas abuse", IncidentSeverity.CRITICAL)

    async def handle_a2a_attack(self, incident: Incident):
        agent_id = incident.agent_id
        await self.agent_api.pause_agent(agent_id)
        incident.actions_taken.append({'action': 'pause_agent', 'timestamp': time.time()})
        await self.evidence_logger.log_evidence(agent_id, incident.id, 'a2a_attack_graph', {})
        await self.alerting.send_alert(f"[CRITICAL] A2A attack from Agent {agent_id}", IncidentSeverity.CRITICAL)

    async def handle_unauthorized_call(self, incident: Incident):
        agent_id = incident.agent_id
        await self.agent_api.pause_agent(agent_id)
        incident.actions_taken.append({'action': 'pause_agent', 'timestamp': time.time()})
        target = incident.evidence.get('target_contract', 'unknown')
        await self.evidence_logger.log_evidence(agent_id, incident.id, 'unauthorized_call', {'target_contract': target})
        await self.alerting.send_alert(f"[CRITICAL] Agent {agent_id} called unauthorized contract: {target}", IncidentSeverity.CRITICAL)

    async def handle_agent_downtime(self, incident: Incident):
        agent_id = incident.agent_id
        incident.actions_taken.append({'action': 'restart_agent', 'timestamp': time.time()})
        await asyncio.sleep(30)
        status = await self.agent_api.get_agent_status(agent_id)
        if status.get('status') != 'running':
            await self.alerting.send_alert(f"[CRITICAL] Agent {agent_id} failed to recover", IncidentSeverity.CRITICAL)

    async def handle_default(self, incident: Incident):
        await self.alerting.send_alert(f"[{incident.severity.value.upper()}] Unhandled incident for Agent {incident.agent_id}: {incident.rule_name}", incident.severity)

    async def escalate_to_human(self, incident: Incident, reason: str):
        await self.alerting.send_pagerduty(f"ESCALATION: Incident {incident.id} for Agent {incident.agent_id}: {reason}", 'critical')
        incident.actions_taken.append({'action': 'escalate_to_human', 'reason': reason, 'timestamp': time.time()})

5.2 分级响应模型

响应措施根据威胁严重度分级,形成渐进式响应链:

Level 1: INFO      记录日志 -> 更新指标
Level 2: WARNING   告警通知 -> 记录证据
Level 3: HIGH      限流/降级 -> 告警 -> 记录
Level 4: CRITICAL  暂停 Agent -> 紧急轮换 -> 升级
Level 5: KILL      终止 Agent -> 全面隔离 -> 取证
# graduated_response.py

class GraduatedResponder:
    def __init__(self, responder: IncidentResponder):
        self.responder = responder
        self.escalation_levels = {
            IncidentSeverity.INFO: self.level_1_info,
            IncidentSeverity.WARNING: self.level_2_warning,
            IncidentSeverity.CRITICAL: self.level_4_critical,
        }

    async def respond(self, incident: Incident):
        handler = self.escalation_levels.get(incident.severity, self.level_default)
        await handler(incident)

    async def level_1_info(self, incident: Incident):
        await self.responder.evidence_logger.log_evidence(incident.agent_id, incident.id, 'info_event', incident.evidence)

    async def level_2_warning(self, incident: Incident):
        await self.responder.alerting.send_alert(f"[WARNING] {incident.description}", IncidentSeverity.WARNING)
        await self.responder.evidence_logger.log_evidence(incident.agent_id, incident.id, 'warning_event', incident.evidence)

    async def level_4_critical(self, incident: Incident):
        await self.responder.agent_api.pause_agent(incident.agent_id)
        await self.responder.alerting.send_pagerduty(f"[CRITICAL] {incident.description}", 'critical')
        await self.responder.evidence_logger.log_evidence(incident.agent_id, incident.id, 'critical_event', incident.evidence)

    async def level_default(self, incident: Incident):
        await self.responder.handle_incident(incident)

5.3 自动化响应最佳实践

+---------------------------------------------------------+
|               自动化响应实施清单                           |
+---------------------------------------------------------+
| [ ] 每个威胁规则绑定一个 Playbook                          |
| [ ] Playbook 实现幂等性(断路器保护重复执行)               |
| [ ] 关键操作(暂停、终止)需人工确认确认窗口                 |
| [ ] 自动化响应执行后记录完整审计追踪                        |
| [ ] 每月至少一次 Playbook 演练                             |
| [ ] 所有外部调用设置超时(默认 5 秒)                      |
| [ ] 自动化响应系统自身高可用部署(多副本)                  |
| [ ] 避免响应级联:一个 Agent 的自动响应不会触发另一个       |
| [ ] 维护 Playbook 版本控制,变更需 Code Review             |
| [ ] 自动化响应失败时自动升级到人工                         |
+---------------------------------------------------------+

6. 审计与合规

6.1 Agent 行为审计日志

// audit-logger.ts

import { createHash } from 'crypto';
import { appendFileSync, readFileSync, mkdirSync, existsSync } from 'fs';
import { join } from 'path';

interface AuditEntry {
  id: string;
  timestamp: number;
  agentId: string;
  action: string;
  actor: string;
  resource: string;
  details: Record<string, unknown>;
  severity: 'info' | 'warning' | 'critical';
  previousHash: string;
  hash: string;
}

class MerkleAuditLogger {
  private logDir: string;
  private currentHash: string;
  private entries: AuditEntry[] = [];

  constructor(logDir: string = '/var/log/msg-agent/audit') {
    this.logDir = logDir;
    if (!existsSync(logDir)) mkdirSync(logDir, { recursive: true });
    this.currentHash = '0'.repeat(64);
  }

  async log(entry: Omit<AuditEntry, 'id' | 'hash' | 'previousHash'>): Promise<AuditEntry> {
    const auditEntry: AuditEntry = {
      ...entry,
      id: `${Date.now()}-${crypto.randomUUID().slice(0, 8)}`,
      previousHash: this.currentHash,
      hash: '',
    };
    const hashInput = [auditEntry.previousHash, auditEntry.timestamp.toString(), auditEntry.agentId, auditEntry.action, JSON.stringify(auditEntry.details), auditEntry.id].join('|');
    auditEntry.hash = createHash('sha256').update(hashInput).digest('hex');
    this.currentHash = auditEntry.hash;
    this.entries.push(auditEntry);
    this.appendToFile(auditEntry);
    return auditEntry;
  }

  private appendToFile(entry: AuditEntry): void {
    const dateStr = new Date(entry.timestamp).toISOString().slice(0, 10);
    const logFile = join(this.logDir, `audit-${dateStr}.jsonl`);
    appendFileSync(logFile, JSON.stringify(entry) + '\n', { encoding: 'utf-8' });
  }

  verifyChain(fromTimestamp?: number): boolean {
    const relevantEntries = fromTimestamp ? this.entries.filter(e => e.timestamp >= fromTimestamp) : this.entries;
    for (let i = 0; i < relevantEntries.length; i++) {
      const entry = relevantEntries[i];
      const hashInput = [entry.previousHash, entry.timestamp.toString(), entry.agentId, entry.action, JSON.stringify(entry.details), entry.id].join('|');
      const expectedHash = createHash('sha256').update(hashInput).digest('hex');
      if (entry.hash !== expectedHash) { console.error(`[Audit] Hash mismatch at entry ${entry.id}`); return false; }
      if (i > 0 && entry.previousHash !== relevantEntries[i - 1].hash) { console.error(`[Audit] Chain break at entry ${entry.id}`); return false; }
    }
    return true;
  }

  query(filter: any): AuditEntry[] {
    return this.entries.filter(entry => {
      if (filter.agentId && entry.agentId !== filter.agentId) return false;
      if (filter.action && entry.action !== filter.action) return false;
      if (filter.severity && entry.severity !== filter.severity) return false;
      if (filter.fromTimestamp && entry.timestamp < filter.fromTimestamp) return false;
      return true;
    });
  }
}

const auditLogger = new MerkleAuditLogger();

6.2 Constitution 合规检查

# compliance_checker.py

from dataclasses import dataclass
import time


@dataclass
class ComplianceReport:
    agent_id: str
    timestamp: int
    overall_score: float
    checks: list
    violations: list[dict]
    recommendations: list[str]
    passed: bool


@dataclass
class ComplianceCheck:
    check_name: str
    status: str
    score: float
    details: str


class ComplianceChecker:
    def __init__(self, constitution: ConstitutionDocument):
        self.constitution = constitution
        self.check_functions = [
            self.check_action_whitelist, self.check_gas_limits,
            self.check_tx_frequency, self.check_value_limits,
            self.check_target_contracts, self.check_key_rotation,
            self.check_heartbeat, self.check_backup_status,
        ]

    async def run_full_check(self, agent_data: AgentAuditData) -> ComplianceReport:
        checks = []
        violations = []
        for check_fn in self.check_functions:
            result = await check_fn(agent_data)
            checks.append(result)
            if result.status == 'fail':
                violations.append({'check': result.check_name, 'detail': result.details})
        total_score = sum(c.score for c in checks) / len(checks)
        recommendations = self._generate_recommendations(checks)
        return ComplianceReport(
            agent_id=agent_data.agent_id, timestamp=int(time.time()),
            overall_score=total_score, checks=checks, violations=violations,
            recommendations=recommendations, passed=total_score >= 80.0 and len(violations) == 0,
        )

    async def check_action_whitelist(self, data: 'AgentAuditData') -> ComplianceCheck:
        forbidden = [a for a in data.recent_actions if a.action_type in self.constitution.forbidden_actions]
        if forbidden:
            return ComplianceCheck('action_whitelist', 'fail', 0.0, f"Agent performed {len(forbidden)} forbidden actions")
        return ComplianceCheck('action_whitelist', 'pass', 100.0, "All actions in whitelist")

    async def check_gas_limits(self, data: 'AgentAuditData') -> ComplianceCheck:
        if not data.gas_records:
            return ComplianceCheck('gas_limits', 'pass', 100.0, "No gas data")
        max_gas = max(r.gas_used for r in data.gas_records)
        gas_limit = 500000
        if max_gas > gas_limit:
            return ComplianceCheck('gas_limits', 'fail', (gas_limit / max_gas) * 100, f"Max gas {max_gas} exceeds limit {gas_limit}")
        return ComplianceCheck('gas_limits', 'pass', 100.0, f"Max gas {max_gas} within limit {gas_limit}")

    async def check_tx_frequency(self, data: 'AgentAuditData') -> ComplianceCheck:
        if not data.recent_actions:
            return ComplianceCheck('tx_frequency', 'pass', 100.0, "No transactions")
        tx_per_hour = (len(data.recent_actions) / data.time_window_seconds) * 3600
        max_freq = 100
        if tx_per_hour > max_freq:
            return ComplianceCheck('tx_frequency', 'warning', (max_freq / tx_per_hour) * 100, f"Tx frequency {tx_per_hour:.1f}/hour exceeds limit {max_freq}/hour")
        return ComplianceCheck('tx_frequency', 'pass', 100.0, f"Tx frequency {tx_per_hour:.1f}/hour within limit {max_freq}/hour")

    async def check_value_limits(self, data: 'AgentAuditData') -> ComplianceCheck:
        max_value = max((a.params.get('value', 0) for a in data.recent_actions), default=0)
        value_limit = 1000000000
        if max_value > value_limit:
            return ComplianceCheck('value_limits', 'fail', 0.0, f"Max value {max_value} exceeds limit {value_limit}")
        return ComplianceCheck('value_limits', 'pass', 100.0, f"Max value {max_value} within limit {value_limit}")

    async def check_key_rotation(self, data: 'AgentAuditData') -> ComplianceCheck:
        if not data.last_key_rotation:
            return ComplianceCheck('key_rotation', 'fail', 0.0, "No key rotation recorded")
        elapsed = time.time() - data.last_key_rotation
        max_interval = 30 * 86400
        if elapsed > max_interval:
            return ComplianceCheck('key_rotation', 'warning', (max_interval / elapsed) * 100, f"Last rotation was {elapsed / 86400:.1f} days ago")
        return ComplianceCheck('key_rotation', 'pass', 100.0, f"Key rotated {(max_interval - elapsed) / 86400:.1f} days remaining")

    async def check_heartbeat(self, data: 'AgentAuditData') -> ComplianceCheck:
        if not data.last_heartbeat:
            return ComplianceCheck('heartbeat', 'fail', 0.0, "No heartbeat recorded")
        elapsed = time.time() - data.last_heartbeat
        if elapsed > 300:
            return ComplianceCheck('heartbeat', 'fail', 0.0, f"Last heartbeat {elapsed}s ago (timeout: 300s)")
        return ComplianceCheck('heartbeat', 'pass', 100.0, f"Agent online, last heartbeat {elapsed}s ago")

    async def check_backup_status(self, data: 'AgentAuditData') -> ComplianceCheck:
        if not data.last_backup:
            return ComplianceCheck('backup_status', 'fail', 0.0, "No backup recorded")
        elapsed = time.time() - data.last_backup
        if elapsed > 86400 * 2:
            return ComplianceCheck('backup_status', 'warning', 50.0, f"Last backup {elapsed / 86400:.1f} days ago")
        return ComplianceCheck('backup_status', 'pass', 100.0, f"Recent backup {elapsed / 86400:.1f} days ago")

    def _generate_recommendations(self, checks: list[ComplianceCheck]) -> list[str]:
        recommendations = []
        for check in checks:
            if check.status == 'fail':
                recommendations.append(f"[CRITICAL] {check.check_name}: {check.details}")
            elif check.status == 'warning':
                recommendations.append(f"[WARNING] {check.check_name}: {check.details}")
        if not recommendations:
            recommendations.append("All compliance checks passed")
        return recommendations


@dataclass
class AgentAuditData:
    agent_id: str
    recent_actions: list
    gas_records: list
    time_window_seconds: int
    last_heartbeat: float
    last_key_rotation: float
    last_backup: float
    constitution_embedding_distance: float

6.3 渗透测试指南

penetration_test_report:
  agent_id: "agent-trader-01"
  test_date: "2026-04-15"
  tester: "security-team"
  scope:
    - key_management
    - access_control
    - prompt_injection
    - a2a_security
    - economic_security
  findings:
    - severity: high
      title: "Insufficient key file permissions"
      description: "Agent key file has 644 permissions"
      impact: "Local user can read private key"
      recommendation: "Set key file permissions to 600"
      cvss_score: 7.5
    - severity: medium
      title: "Rate limiting bypass via IP rotation"
      description: "Rate limiter does not track by DID"
      impact: "Attacker can bypass rate limits"
      recommendation: "Implement DID-based rate limiting"
  summary:
    total_findings: 8
    critical: 0
    high: 2
    medium: 4
    low: 2
    passed: True

6.4 审计与合规最佳实践

+---------------------------------------------------------+
|               审计与合规实施清单                           |
+---------------------------------------------------------+
| [ ] 所有 Agent 操作记录 Merkle 哈希链审计日志              |
| [ ] 关键操作日志同步到链上不可篡改存储                     |
| [ ] 每日自动运行 Constitution 合规检查                     |
| [ ] 合规评分低于 80 自动告警                               |
| [ ] 每季度执行完整渗透测试                                 |
| [ ] 每次 Constitution 更新触发合规重检                     |
| [ ] 审计日志保留至少 1 年                                  |
| [ ] 支持审计日志导出供外部审计                            |
| [ ] 定期审查访问权限分配(季度)                           |
| [ ] 渗透测试发现的高危漏洞需在 7 天内修复                  |
+---------------------------------------------------------+

7. 安全配置模板

7.1 完整安全配置

# agent-security.yaml

agent_security:
  version: "1.0"
  agent_id: "${AGENT_ID}"
  chain_id: "msg-chain-1"
  environment: "production"

  key_management:
    algorithm: "DILITHIUM5"
    key_rotation_days: 30
    key_rotation_time: "03:00"
    key_storage:
      provider: "hsm"
      hsm:
        endpoint: "${MSG_HSM_ENDPOINT}"
        slot_id: "${MSG_HSM_SLOT_ID}"
        pkcs11_library: "/usr/lib/softhsm/libsofthsm2.so"
      tee:
        provider: "sgx"
        enclave_id: "${MSG_TEE_ENCLAVE_ID}"
      encrypted_file:
        path: "/opt/msg-agent/keys"
        encryption: "AES-256-GCM"
        key_derivation: "PBKDF2-SHA256"
        iterations: 600000
    backup:
      enabled: true
      schedule: "0 4 * * *"
      shard_count: 5
      recovery_threshold: 3
      retention_days: 90
      providers:
        - name: "aws-s3"
          bucket: "msg-agent-backup"
          region: "us-east-1"
        - name: "gcs"
          bucket: "msg-agent-backup-gcs"
          project: "msg-chain-prod"

  access_control:
    enforce_rbac: true
    jwt_expiry_seconds: 3600
    challenge_expiry_seconds: 300
    rate_limiting:
      default_operations_per_hour: 100
      burst_limit: 10
    multisig:
      enabled: true
      default_threshold: 2
      emergency_threshold: 3
      contract: "msg1multisig..."

  monitoring:
    metrics_port: 9090
    heartbeat_interval_seconds: 30
    constitution_check_interval_minutes: 15
    risk_scoring:
      enabled: true
      alert_threshold: 80
      kill_threshold: 95
    alerting:
      channels:
        - type: "slack"
          webhook: "${MSG_SLACK_WEBHOOK}"
        - type: "pagerduty"
          service_key: "${MSG_PAGERDUTY_KEY}"
        - type: "telegram"
          bot_token: "${MSG_TELEGRAM_BOT}"
          chat_id: "${MSG_TELEGRAM_CHAT}"

  transaction_limits:
    max_tx_per_hour: 100
    max_gas_per_tx: 500000
    max_gas_daily: 10000000
    max_value_per_tx: 1000000000
    max_value_daily: 50000000000
    allowed_calls:
      - "agent_registry_v1"
      - "agent_payment_v1"
      - "did_registry_v1"
      - "multisig_wallet_v1"
    forbidden_calls:
      - "msg1dangerous..."

  constitution:
    enforcement: "strict"
    embedding_check_enabled: true
    embedding_threshold: 0.95
    auto_update_on_drift: false
    drift_check_interval_minutes: 60

  logging:
    audit_log_dir: "/var/log/msg-agent/audit"
    audit_log_retention_days: 365
    merkle_chain_enabled: true
    log_level: "info"
    remote_syslog: "${MSG_SYSLOG_ENDPOINT}"

  automated_response:
    enabled: true
    circuit_breaker_cooldown_seconds: 120
    max_response_attempts: 3
    escalation_policy:
      warning: "15m"
      critical: "5m"
    playbooks:
      constitution_violation: "pause_and_alert"
      key_compromise: "kill_and_rotate"
      payment_fraud: "pause_and_investigate"
      gas_abuse: "graduated_response"

7.2 环境变量安全配置

# .env.security
# MSG Chain Agent 安全相关环境变量

# 链配置
export MSG_CHAIN_ID="msg-chain-1"
export MSG_RPC_ENDPOINT="https://rpc.msgchain.org"
export MSG_AGENT_ID="agent-main-trader-01"

# 密钥存储
export MSG_HSM_ENDPOINT="pkcs11:slot-id=0"
export MSG_HSM_SLOT_ID="0"
export MSG_HSM_PIN=""
export MSG_KEY_ENC_PASSWORD=""
export MSG_TEE_ENCLAVE_ID=""

# 云备份
export MSG_AWS_REGION="us-east-1"
export MSG_AWS_ACCESS_KEY_ID=""
export MSG_AWS_SECRET_ACCESS_KEY=""
export MSG_GCP_PROJECT_ID="msg-chain-prod"
export MSG_GCP_SERVICE_ACCOUNT_KEY="/etc/msg-agent/gcp-sa.json"

# 多签合约
export MSG_MULTISIG_CONTRACT="msg1multisig..."

# JWT
export MSG_JWT_SECRET=""

# 告警通知
export MSG_SLACK_WEBHOOK=""
export MSG_PAGERDUTY_KEY=""
export MSG_TELEGRAM_BOT=""
export MSG_TELEGRAM_CHAT=""

# 日志
export MSG_SYSLOG_ENDPOINT="syslog://logs.msgchain.org:514"

7.3 Docker Compose 安全部署

# docker-compose.security.yml

version: "3.8"

services:
  agent-core:
    image: msgchain/agent-core:v1.0.0
    container_name: msg-agent-trader-01
    restart: unless-stopped
    ports:
      - "9090:9090"
    volumes:
      - /opt/msg-agent/keys:/opt/msg-agent/keys:ro
      - /var/log/msg-agent:/var/log/msg-agent
      - /etc/msg-agent:/etc/msg-agent:ro
    environment:
      - MSG_CHAIN_ID=msg-chain-1
      - MSG_RPC_ENDPOINT=https://rpc.msgchain.org
      - MSG_AGENT_ID=agent-main-trader-01
    env_file:
      - .env.security
    security_opt:
      - no-new-privileges:true
      - seccomp:seccomp-profile.json
    cap_drop:
      - ALL
    cap_add:
      - NET_BIND_SERVICE
    read_only: true
    tmpfs:
      - /tmp:size=64M
    logging:
      driver: "json-file"
      options:
        max-size: "50m"
        max-file: "10"
    healthcheck:
      test: ["CMD", "curl", "-f", "http://localhost:9090/health"]
      interval: 30s
      timeout: 10s
      retries: 3
    deploy:
      resources:
        limits:
          cpus: "2"
          memory: "4G"
        reservations:
          cpus: "0.5"
          memory: "1G"

  agent-metrics:
    image: prom/prometheus:v2.50.0
    container_name: msg-agent-prometheus
    volumes:
      - ./prometheus.yml:/etc/prometheus/prometheus.yml:ro
    ports:
      - "9091:9090"

  agent-alertmanager:
    image: prom/alertmanager:v0.27.0
    container_name: msg-agent-alertmanager
    volumes:
      - ./alertmanager.yml:/etc/alertmanager/alertmanager.yml:ro
    ports:
      - "9093:9093"

7.4 安全配置验证脚本

# validate_security_config.py
import yaml
import sys

REQUIRED_FIELDS = [
    'agent_security.key_management.key_rotation_days',
    'agent_security.access_control.enforce_rbac',
    'agent_security.monitoring.heartbeat_interval_seconds',
    'agent_security.transaction_limits.max_tx_per_hour',
    'agent_security.transaction_limits.max_gas_per_tx',
    'agent_security.constitution.enforcement',
    'agent_security.logging.audit_log_dir',
    'agent_security.automated_response.enabled',
]

def validate_config(path: str):
    with open(path, 'r') as f:
        config = yaml.safe_load(f)
    missing = []
    for field in REQUIRED_FIELDS:
        parts = field.split('.')
        current = config
        for part in parts:
            if isinstance(current, dict):
                current = current.get(part)
            else:
                current = None
                break
        if current is None:
            missing.append(field)
    if missing:
        print(f"VALIDATION FAILED: Missing fields: {missing}")
        sys.exit(1)
    print("VALIDATION PASSED: All required fields present")
    print(f"Key rotation: {config['agent_security']['key_management']['key_rotation_days']} days")
    print(f"RBAC enabled: {config['agent_security']['access_control']['enforce_rbac']}")
    print(f"Gas limit per tx: {config['agent_security']['transaction_limits']['max_gas_per_tx']}")

if __name__ == '__main__':
    validate_config(sys.argv[1])

8. 应急响应流程

8.1 事件分类矩阵

级别 定义 响应时间 示例
SEV-1 关键 Agent 遭入侵/资产被盗 15 分钟 密钥泄露、支付欺诈
SEV-2 Agent 行为严重偏离 Constitution 30 分钟 Constitution 违反、未授权调用
SEV-3 Agent 性能异常/风险评分升高 2 小时 Gas 激增、A2A 异常
SEV-4 常规告警/需关注 24 小时 备份失败、心跳延迟

8.2 SEV-1 应急响应 SOP

# sev1-sop.yaml

incident_id: "SEV-1-${DATE}-${SEQ}"
detected_by: "automated_detection | operator_report"
assigned_to: "security_team_oncall"

steps:
  - phase: "0-5min 检测与确认"
    actions:
      - "确认告警来源和真实性"
      - "判断受影响 Agent 范围和资产"
      - "通知值班安全工程师"
      - "在 incident channel 创建事件线程"

  - phase: "5-15min 遏制"
    actions:
      - "执行 emergency_pause 暂停目标 Agent"
      - "如果疑似密钥泄露,执行 emergency_kill"
      - "执行 emergency_rotate 轮换密钥"
      - "冻结关联链上资产(如支持)"
      - "记录所有动作到 Incident 日志"

  - phase: "15-60min 调查"
    actions:
      - "导出 Agent 完整日志和审计记录"
      - "分析异常交易链"
      - "检查其他 Agent 是否受影响"
      - "提取并保存取证数据"
      - "初步判断攻击向量"

  - phase: "1-4h 恢复"
    actions:
      - "确认威胁已隔离"
      - "评估是否可以安全恢复"
      - "部署修复后的 Agent 实例"
      - "逐步恢复服务"
      - "验证恢复后的 Agent 行为"

  - phase: "4-24h 复盘"
    actions:
      - "撰写安全事件报告"
      - "更新威胁检测规则"
      - "修复发现的漏洞"
      - "更新 Playbook"
      - "发送复盘邮件"

communication:
  - "第一轮(15min):确认事件,正在遏制"
  - "第二轮(1h):攻击向量初步分析"
  - "第三轮(4h):恢复状态更新"
  - "第四轮(24h):完整复盘报告"

escalation:
  primary: "security@msgchain.org"
  secondary: "security-manager@msgchain.org"
  executive: "cto@msgchain.org"

8.3 应急响应脚本

#!/bin/bash
# emergency-response.sh
# MSG Chain Agent 应急响应脚本

set -euo pipefail

AGENT_ID="${1:-}"
ACTION="${2:-status}"
TIMESTAMP=$(date +%s)
LOG_DIR="/var/log/msg-agent/emergency"
mkdir -p "$LOG_DIR"

log() { echo "[$(date +%Y-%m-%dT%H:%M:%S)] $*" | tee -a "$LOG_DIR/response-$AGENT_ID.log"; }

emergency_pause() {
    log "ACTION: emergency_pause for agent $AGENT_ID"
    # 调用 Agent API 暂停
    curl -s -X POST "http://agent-api:8080/api/v1/agents/$AGENT_ID/pause" \
        -H "Authorization: Bearer $(cat /run/secrets/emergency_token)" \
        -H "Content-Type: application/json" \
        -d '{"reason": "emergency_response","initiator": "automated"}'
    log "RESULT: Agent $AGENT_ID pause command issued"
}

emergency_kill() {
    log "ACTION: emergency_kill for agent $AGENT_ID"
    curl -s -X POST "http://agent-api:8080/api/v1/agents/$AGENT_ID/kill" \
        -H "Authorization: Bearer $(cat /run/secrets/emergency_token)" \
        -H "Content-Type: application/json" \
        -d '{"reason": "emergency_response","initiator": "automated"}'
    log "RESULT: Agent $AGENT_ID kill command issued"
}

collect_evidence() {
    log "ACTION: Collecting evidence for agent $AGENT_ID"
    EVIDENCE_DIR="$LOG_DIR/evidence-$AGENT_ID-$TIMESTAMP"
    mkdir -p "$EVIDENCE_DIR"
    # 收集日志
    journalctl -u msg-agent@"$AGENT_ID" --since "1 hour ago" > "$EVIDENCE_DIR/journal.log" 2>/dev/null || true
    # 收集审计日志
    cp /var/log/msg-agent/audit/audit-*.jsonl "$EVIDENCE_DIR/" 2>/dev/null || true
    # 收集链上交易
    msgcli query tx --agent "$AGENT_ID" --limit 50 > "$EVIDENCE_DIR/txs.json" 2>/dev/null || true
    log "RESULT: Evidence collected at $EVIDENCE_DIR"
}

case "$ACTION" in
    pause)
        emergency_pause
        collect_evidence
        ;;
    kill)
        emergency_kill
        collect_evidence
        ;;
    investigate)
        collect_evidence
        ;;
    status)
        echo "Usage: $0 <agent_id> <pause|kill|investigate>"
        ;;
esac

8.4 应急响应最佳实践

+---------------------------------------------------------+
|               应急响应实施清单                             |
+---------------------------------------------------------+
| [ ] 定义清晰的事件分类矩阵(SEV-1 到 SEV-4)              |
| [ ] SEV-1 响应目标 15 分钟内完成遏制                      |
| [ ] 准备应急响应脚本并定期演练                             |
| [ ] 维护紧急操作令牌(短期有效,独立存储)                  |
| [ ] 所有应急操作记录完整审计追踪                           |
| [ ] 建立 7x24 值班制度,明确升级路径                     |
| [ ] 每月一次桌面演练,每季度一次实战演练                   |
| [ ] 每次事件后 24 小时内完成复盘                          |
| [ ] 更新 Playbook 和检测规则基于复盘结论                   |
| [ ] 维护联系人清单和通信渠道冗余                          |
+---------------------------------------------------------+

本文档为 MSG Chain AI Agent 安全运维的全面指南。所有配置和代码示例需根据实际环境调整。密钥和令牌等敏感信息禁止直接写入配置文件中。


本文档内容基于 MSGChain 代码库真实状态编写,非 AI 自动生成。
主网状态: No-Go | 白皮书: https://msgchain.org/whitepaper/