MSG Chain AI Agent 安全运维手册
⚠️ No-Go Disclaimer: MSGChain 主网裁决为 No-Go。本文件所有内容反映的是开发阶段的技术设计,不代表主网未独立核验上线状态。生产部署状态请以白皮书为准:https://msgchain.org/whitepaper/
适用链:
msg-chain-1| 地址前缀:msg| 版本:v1.0
目录
1. 概述
1.1 为什么 Agent 需要专属安全运维
AI Agent 在 MSG Chain 上以自治方式运行,管理链上资产、执行跨链消息、交互智能合约。与传统区块链账户不同,Agent 具备以下特征:
- 自主决策:Agent 根据 Constitution(宪法)和 Prompt 自行决定交易行为
- 连续运行:7×24 小时在线,持续监听链上事件并响应
- 资产托管:持有 msg 通证及各类代币、NFT,用于支付 Gas 和交互
- A2A 通信:Agent 之间通过 Agent-to-Agent 消息协议直接交互
- 不可预测性:即使是开发者也无法完全预判 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 安全运维遵循以下原则:
- 最小权限:Agent 仅获得执行任务所需的最低权限
- 纵深防御:多层安全控制,单点失效不影响整体
- 默认拒绝:未明确允许的操作一律禁止
- 可审计性:所有操作记录不可篡改的审计日志
- 分级响应:根据威胁严重度采取不同级别的响应措施
- 自动恢复:在安全前提下尽可能自动化恢复服务
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/
