AI Agent 可观测性与监控运维指南
⚠️ No-Go Disclaimer: MSGChain 主网裁决为 No-Go。本文件所有内容反映的是开发阶段的技术设计,不代表主网未独立核验上线状态。生产部署状态请以白皮书为准:https://msgchain.org/whitepaper/
目录
- 概述
- 指标(Metrics)体系
- 日志(Logging)体系
- 链路追踪(Tracing)
- 告警规则
- Grafana 仪表盘
- 告警通知渠道
- 部署架构
- 附录
1. 概述
1.1 为什么 AI Agent 需要可观测性
AI Agent 运行在 MSG Chain 上,涉及链上交易、A2A(Agent-to-Agent)通信、支付结算等多个环节。与传统微服务不同,Agent 的行为具有非确定性和自治性,传统的监控手段不足以覆盖以下风险:
- 非预期行为:LLM 驱动的 Agent 可能产生意料之外的决策路径
- 资金安全:Agent 自主发起的支付交易需要全程可审计
- 宪法违规:Agent 可能偏离其宪法约束(Constitution),需要实时检测
- 性能退化:链上交互延迟、Gas 消耗异常、A2A 通信超时
- 故障定位:多 Agent 协作场景下,根因分析复杂度指数级增长
1.2 三大支柱
可观测性建立在三个核心支柱之上:
| 支柱 | 数据形态 | 核心工具 | 解决的问题 |
|---|---|---|---|
| 指标 (Metrics) | 聚合数值 | Prometheus | 趋势分析、阈值告警 |
| 日志 (Logging) | 离散事件 | Loki + structlog | 细粒度排查、审计 |
| 链路追踪 (Tracing) | 请求链路 | Tempo/Jaeger | 分布式性能分析 |
1.3 适用读者
- MSG Chain 上的 AI Agent 开发者
- 区块链基础设施运维工程师
- SRE 和平台工程团队
- Web3 项目技术负责人
1.4 前置依赖
- MSG Chain 节点运行正常(RPC 端点可达)
- Docker Compose 或 Kubernetes 集群
- Go 1.24.6+ / Python 3.11+ / Node.js 18+
- Prometheus + Grafana 部署经验
2. 指标(Metrics)体系
2.1 Agent 指标定义
Agent 指标以 Prometheus 标准格式暴露,所有指标遵循 msg_agent_* 命名约定。
// agent-metrics.ts
import { Gauge, Counter, Histogram, Registry } from 'prom-client';
import { MSG_CHAIN_CONFIG } from './config';
export const registry = new Registry();
// ===== Agent 生命周期指标 =====
export const agent_uptime_seconds = new Gauge({
name: 'msg_agent_uptime_seconds',
help: 'Agent 运行时长(秒)',
labelNames: ['agent_id', 'agent_name', 'version'],
registers: [registry],
});
export const agent_start_timestamp = new Gauge({
name: 'msg_agent_start_timestamp',
help: 'Agent 启动时间戳',
labelNames: ['agent_id'],
registers: [registry],
});
// ===== 交易指标 =====
export const agent_transactions_total = new Counter({
name: 'msg_agent_transactions_total',
help: 'Agent 发起的链上交易总数',
labelNames: ['agent_id', 'tx_type', 'status'],
registers: [registry],
});
export const agent_transactions_duration_seconds = new Histogram({
name: 'msg_agent_transactions_duration_seconds',
help: '交易确认耗时分布(秒)',
labelNames: ['agent_id', 'tx_type'],
buckets: [0.5, 1, 2, 5, 10, 30, 60, 120],
registers: [registry],
});
export const agent_gas_spent = new Histogram({
name: 'msg_agent_gas_spent',
help: '单笔交易 Gas 消耗分布',
labelNames: ['agent_id', 'tx_type'],
buckets: [50000, 100000, 200000, 500000, 1000000, 5000000],
registers: [registry],
});
export const agent_gas_price_gauge = new Gauge({
name: 'msg_agent_gas_price_gauge',
help: '当前使用的 Gas 价格(amsg)',
labelNames: ['agent_id'],
registers: [registry],
});
// ===== A2A 通信指标 =====
export const agent_a2a_messages_sent = new Counter({
name: 'msg_agent_a2a_messages_sent',
help: 'Agent 发送的 A2A 消息总数',
labelNames: ['sender_id', 'receiver_id', 'protocol', 'message_type'],
registers: [registry],
});
export const agent_a2a_messages_received = new Counter({
name: 'msg_agent_a2a_messages_received',
help: 'Agent 接收的 A2A 消息总数',
labelNames: ['receiver_id', 'sender_id', 'protocol'],
registers: [registry],
});
export const agent_a2a_latency_seconds = new Histogram({
name: 'msg_agent_a2a_latency_seconds',
help: 'A2A 消息端到端延迟分布(秒)',
labelNames: ['sender_id', 'receiver_id'],
buckets: [0.1, 0.5, 1, 2, 5, 10, 30],
registers: [registry],
});
export const agent_a2a_message_size_bytes = new Histogram({
name: 'msg_agent_a2a_message_size_bytes',
help: 'A2A 消息大小分布(字节)',
labelNames: ['message_type'],
buckets: [256, 1024, 4096, 16384, 65536, 262144, 1048576],
registers: [registry],
});
// ===== 支付指标 =====
export const agent_payment_sessions_total = new Counter({
name: 'msg_agent_payment_sessions_total',
help: 'Agent 支付会话总数',
labelNames: ['agent_id', 'payment_method', 'currency'],
registers: [registry],
});
export const agent_payment_success_rate = new Gauge({
name: 'msg_agent_payment_success_rate',
help: 'Agent 支付成功率(最近窗口)',
labelNames: ['agent_id', 'payment_method'],
registers: [registry],
});
export const agent_payment_amount_total = new Counter({
name: 'msg_agent_payment_amount_total',
help: 'Agent 支付总金额(amsg)',
labelNames: ['agent_id', 'currency'],
registers: [registry],
});
export const agent_payment_settlement_latency = new Histogram({
name: 'msg_agent_payment_settlement_latency',
help: '支付结算延迟分布(秒)',
labelNames: ['agent_id', 'payment_method'],
buckets: [1, 5, 10, 30, 60, 120, 300],
registers: [registry],
});
// ===== 宪法合规指标 =====
export const agent_constitution_checks_total = new Counter({
name: 'msg_agent_constitution_checks_total',
help: '宪法合规检查总次数',
labelNames: ['agent_id', 'result'],
registers: [registry],
});
export const agent_constitution_violations_total = new Counter({
name: 'msg_agent_constitution_violations_total',
help: '宪法违规总次数',
labelNames: ['agent_id', 'constitution_article', 'severity'],
registers: [registry],
});
export const agent_constitution_check_duration = new Histogram({
name: 'msg_agent_constitution_check_duration',
help: '宪法检查耗时分布(毫秒)',
labelNames: ['agent_id'],
buckets: [10, 50, 100, 500, 1000, 5000],
registers: [registry],
});
// ===== LLM 调用指标 =====
export const agent_llm_calls_total = new Counter({
name: 'msg_agent_llm_calls_total',
help: 'LLM 调用总次数',
labelNames: ['agent_id', 'model', 'provider'],
registers: [registry],
});
export const agent_llm_latency_seconds = new Histogram({
name: 'msg_agent_llm_latency_seconds',
help: 'LLM 调用延迟分布(秒)',
labelNames: ['agent_id', 'model'],
buckets: [0.5, 1, 2, 5, 10, 30, 60, 120],
registers: [registry],
});
export const agent_llm_tokens_total = new Counter({
name: 'msg_agent_llm_tokens_total',
help: 'LLM Token 消耗总数',
labelNames: ['agent_id', 'model', 'token_type'],
registers: [registry],
});
export const agent_llm_cost_total = new Counter({
name: 'msg_agent_llm_cost_total',
help: 'LLM 调用累计成本(USD)',
labelNames: ['agent_id', 'provider'],
registers: [registry],
});
// ===== Agent 资源指标 =====
export const agent_memory_usage_bytes = new Gauge({
name: 'msg_agent_memory_usage_bytes',
help: 'Agent 进程内存使用量(字节)',
labelNames: ['agent_id', 'memory_type'],
registers: [registry],
});
export const agent_cpu_usage_ratio = new Gauge({
name: 'msg_agent_cpu_usage_ratio',
help: 'Agent 进程 CPU 使用率(0-1)',
labelNames: ['agent_id'],
registers: [registry],
});
export const agent_goroutine_count = new Gauge({
name: 'msg_agent_goroutine_count',
help: 'Agent Go 协程数',
labelNames: ['agent_id'],
registers: [registry],
});
export const agent_open_connections = new Gauge({
name: 'msg_agent_open_connections',
help: 'Agent 当前打开连接数',
labelNames: ['agent_id', 'connection_type'],
registers: [registry],
});
// ===== 业务指标 =====
export const agent_tasks_completed = new Counter({
name: 'msg_agent_tasks_completed',
help: 'Agent 完成任务总数',
labelNames: ['agent_id', 'task_type', 'outcome'],
registers: [registry],
});
export const agent_tasks_in_flight = new Gauge({
name: 'msg_agent_tasks_in_flight',
help: 'Agent 当前执行中的任务数',
labelNames: ['agent_id'],
registers: [registry],
});
export const agent_task_queue_depth = new Gauge({
name: 'msg_agent_task_queue_depth',
help: 'Agent 任务队列深度',
labelNames: ['agent_id', 'priority'],
registers: [registry],
});
export const agent_task_duration_seconds = new Histogram({
name: 'msg_agent_task_duration_seconds',
help: 'Agent 任务执行耗时分布(秒)',
labelNames: ['agent_id', 'task_type'],
buckets: [1, 5, 10, 30, 60, 300, 600, 1800],
registers: [registry],
});
// ===== 链连接指标 =====
export const agent_chain_connection_status = new Gauge({
name: 'msg_agent_chain_connection_status',
help: '链连接状态(1=已连接, 0=断开)',
labelNames: ['agent_id', 'chain_id'],
registers: [registry],
});
export const agent_chain_block_height = new Gauge({
name: 'msg_agent_chain_block_height',
help: 'Agent 观测到的最近区块高度',
labelNames: ['agent_id', 'chain_id'],
registers: [registry],
});
export const agent_chain_sync_status = new Gauge({
name: 'msg_agent_chain_sync_status',
help: '链同步状态(1=已同步, 0=同步中)',
labelNames: ['agent_id', 'chain_id'],
registers: [registry],
});
2.2 链级指标
Agent 需要监控 MSG Chain 的基础健康状况,以下是链级指标定义:
// chain-metrics.ts
import { Gauge, Counter, Histogram, Registry } from 'prom-client';
export const chainRegistry = new Registry();
export const chain_block_height = new Gauge({
name: 'msg_chain_block_height',
help: '当前区块高度',
labelNames: ['chain_id', 'node'],
registers: [chainRegistry],
});
export const chain_block_time_seconds = new Gauge({
name: 'msg_chain_block_time_seconds',
help: '最新出块间隔(秒)',
labelNames: ['chain_id'],
registers: [chainRegistry],
});
export const chain_tx_pool_size = new Gauge({
name: 'msg_chain_tx_pool_size',
help: '交易池待处理交易数',
labelNames: ['chain_id', 'node'],
registers: [chainRegistry],
});
export const chain_gas_price_wei = new Gauge({
name: 'msg_chain_gas_price_wei',
help: '当前网络 Gas 价格(amsg)',
labelNames: ['chain_id'],
registers: [chainRegistry],
});
export const chain_peer_count = new Gauge({
name: 'msg_chain_peer_count',
help: '节点对等连接数',
labelNames: ['node_id'],
registers: [chainRegistry],
});
export const chain_tx_daily_total = new Counter({
name: 'msg_chain_tx_daily_total',
help: '日交易总数',
labelNames: ['chain_id'],
registers: [chainRegistry],
});
export const chain_active_agents = new Gauge({
name: 'msg_chain_active_agents',
help: '链上活跃 Agent 数量',
labelNames: ['chain_id'],
registers: [chainRegistry],
});
export const chain_staked_amount = new Gauge({
name: 'msg_chain_staked_amount',
help: '链上总质押量(amsg)',
labelNames: ['chain_id', 'validator'],
registers: [chainRegistry],
});
export const chain_avg_block_gas = new Gauge({
name: 'msg_chain_avg_block_gas',
help: '区块平均 Gas 消耗',
labelNames: ['chain_id'],
registers: [chainRegistry],
});
2.3 Prometheus Metrics 暴露端点
Agent 通过 HTTP 端点暴露指标,供 Prometheus 采集:
// metrics-exporter.ts
import express from 'express';
import { registry } from './agent-metrics';
import { chainRegistry } from './chain-metrics';
import { authenticate } from './middleware/auth';
const app = express();
const PORT = 9464;
// 认证中间件
app.use('/metrics', authenticate);
// Agent 指标端点
app.get('/metrics/agent', async (req, res) => {
try {
res.set('Content-Type', registry.contentType);
const metrics = await registry.metrics();
res.end(metrics);
} catch (err) {
res.status(500).json({ error: 'metrics collection failed' });
}
});
// 链级指标端点
app.get('/metrics/chain', async (req, res) => {
try {
res.set('Content-Type', chainRegistry.contentType);
const metrics = await chainRegistry.metrics();
res.end(metrics);
} catch (err) {
res.status(500).json({ error: 'chain metrics collection failed' });
}
});
// 合并端点
app.get('/metrics', async (req, res) => {
try {
const merged = Registry.merge([registry, chainRegistry]);
res.set('Content-Type', merged.contentType);
const metrics = await merged.metrics();
res.end(metrics);
} catch (err) {
res.status(500).json({ error: 'metrics collection failed' });
}
});
app.listen(PORT, () => {
console.log(`Metrics exporter listening on port ${PORT}`);
});
2.4 完整示例:Agent API Gateway Metrics Exporter
此示例展示了一个开发参考级别的 Agent API Gateway 指标导出器,整合了 HTTP 请求指标、链交互指标和业务指标。
// agent-gateway-metrics.ts
import { Gauge, Counter, Histogram, Summary, Registry } from 'prom-client';
import express from 'express';
import axios from 'axios';
import { MSG_RPC_ENDPOINT, AGENT_CONTRACT_ADDRESS } from './config';
const gatewayRegistry = new Registry();
// ===== HTTP 请求指标 =====
const http_requests_total = new Counter({
name: 'msg_gateway_http_requests_total',
help: 'HTTP 请求总数',
labelNames: ['method', 'path', 'status_code'],
registers: [gatewayRegistry],
});
const http_request_duration_seconds = new Histogram({
name: 'msg_gateway_http_request_duration_seconds',
help: 'HTTP 请求延迟分布',
labelNames: ['method', 'path'],
buckets: [0.01, 0.05, 0.1, 0.5, 1, 2, 5, 10],
registers: [gatewayRegistry],
});
const http_request_size_bytes = new Summary({
name: 'msg_gateway_http_request_size_bytes',
help: 'HTTP 请求大小',
labelNames: ['method'],
registers: [gatewayRegistry],
});
const http_response_size_bytes = new Summary({
name: 'msg_gateway_http_response_size_bytes',
help: 'HTTP 响应大小',
labelNames: ['method'],
registers: [gatewayRegistry],
});
// ===== 链交互指标 =====
const chain_rpc_calls_total = new Counter({
name: 'msg_gateway_chain_rpc_calls_total',
help: '链 RPC 调用总数',
labelNames: ['rpc_method', 'status'],
registers: [gatewayRegistry],
});
const chain_rpc_latency_seconds = new Histogram({
name: 'msg_gateway_chain_rpc_latency_seconds',
help: '链 RPC 调延迟分布',
labelNames: ['rpc_method'],
buckets: [0.05, 0.1, 0.25, 0.5, 1, 2, 5],
registers: [gatewayRegistry],
});
const chain_last_block_height = new Gauge({
name: 'msg_gateway_chain_last_block_height',
help: '网关观测到的最后区块高度',
labelNames: [],
registers: [gatewayRegistry],
});
// ===== Agent 业务指标 =====
const agent_query_total = new Counter({
name: 'msg_gateway_agent_query_total',
help: 'Agent 查询请求总数',
labelNames: ['agent_address', 'query_type', 'outcome'],
registers: [gatewayRegistry],
});
const agent_balance_amsg = new Gauge({
name: 'msg_gateway_agent_balance_amsg',
help: 'Agent 合约余额(amsg)',
labelNames: ['agent_address'],
registers: [gatewayRegistry],
});
const agent_rate_limit_remaining = new Gauge({
name: 'msg_gateway_agent_rate_limit_remaining',
help: '剩余速率限制配额',
labelNames: ['agent_address', 'limit_type'],
registers: [gatewayRegistry],
});
const active_sessions = new Gauge({
name: 'msg_gateway_active_sessions',
help: '当前活跃会话数',
labelNames: [],
registers: [gatewayRegistry],
});
// ===== 中间件 =====
function metricsMiddleware(req: express.Request, res: express.Response, next: express.NextFunction) {
const start = Date.now();
const originalEnd = res.end;
res.end = function (this: express.Response, ...args: any[]) {
const duration = (Date.now() - start) / 1000;
const path = req.route?.path || req.path;
const method = req.method;
http_requests_total.inc({ method, path, status_code: res.statusCode });
http_request_duration_seconds.observe({ method, path }, duration);
if (req.headers['content-length']) {
http_request_size_bytes.observe({ method }, parseInt(req.headers['content-length']));
}
if (res.getHeader('content-length')) {
http_response_size_bytes.observe({ method }, parseInt(res.getHeader('content-length') as string));
}
return originalEnd.apply(this, args as any);
}.bind(res);
next();
}
// ===== 链状态采集器 =====
async function collectChainMetrics() {
try {
const start = Date.now();
const response = await axios.post(MSG_RPC_ENDPOINT, {
jsonrpc: '2.0',
method: 'eth_blockNumber',
params: [],
id: 1,
});
const duration = (Date.now() - start) / 1000;
chain_rpc_calls_total.inc({ rpc_method: 'eth_blockNumber', status: 'success' });
chain_rpc_latency_seconds.observe({ rpc_method: 'eth_blockNumber' }, duration);
if (response.data && response.data.result) {
const height = parseInt(response.data.result, 16);
chain_last_block_height.set(height);
}
} catch (err) {
chain_rpc_calls_total.inc({ rpc_method: 'eth_blockNumber', status: 'error' });
console.error('Failed to fetch block height:', err);
}
}
async function collectBalanceMetrics(agentAddresses: string[]) {
for (const address of agentAddresses) {
try {
const start = Date.now();
const response = await axios.post(MSG_RPC_ENDPOINT, {
jsonrpc: '2.0',
method: 'eth_getBalance',
params: [address, 'latest'],
id: 1,
});
chain_rpc_calls_total.inc({ rpc_method: 'eth_getBalance', status: 'success' });
chain_rpc_latency_seconds.observe({ rpc_method: 'eth_getBalance' }, (Date.now() - start) / 1000);
if (response.data?.result) {
const balanceWei = BigInt(response.data.result);
const balanceAmsg = Number(balanceWei) / 1e18;
agent_balance_amsg.set({ agent_address: address }, balanceAmsg);
}
} catch (err) {
chain_rpc_calls_total.inc({ rpc_method: 'eth_getBalance', status: 'error' });
}
}
}
// ===== 定时采集 =====
let collectInterval: NodeJS.Timeout;
export function startMetricsCollection(agentAddresses: string[], intervalMs = 15000) {
collectChainMetrics();
collectBalanceMetrics(agentAddresses);
collectInterval = setInterval(() => {
collectChainMetrics();
collectBalanceMetrics(agentAddresses);
}, intervalMs);
}
export function stopMetricsCollection() {
if (collectInterval) clearInterval(collectInterval);
}
// ===== HTTP 服务器 =====
export function startMetricsServer(port = 9464) {
const app = express();
app.use(metricsMiddleware);
app.get('/health', (req, res) => {
res.json({ status: 'ok', timestamp: Date.now() });
});
app.get('/metrics', async (req, res) => {
try {
res.set('Content-Type', gatewayRegistry.contentType);
res.end(await gatewayRegistry.metrics());
} catch (err) {
res.status(500).json({ error: 'metrics error' });
}
});
app.get('/metrics/:agentAddress', async (req, res) => {
try {
const agentAddress = req.params.agentAddress;
const metrics = await gatewayRegistry.getMetricsAsJSON();
const filtered = metrics.map(m => ({
...m,
values: m.values?.filter((v: any) => v.labels?.agent_address === agentAddress),
})).filter(m => m.values?.length > 0);
res.json(filtered);
} catch (err) {
res.status(500).json({ error: 'filtered metrics error' });
}
});
app.listen(port, () => {
console.log(`[Metrics] Gateway exporter listening on :${port}`);
});
}
// ===== 主入口 =====
if (require.main === module) {
const agentAddresses = [
'msg1agent000000000000000000000000000000001',
'msg1agent000000000000000000000000000000002',
'msg1agent000000000000000000000000000000003',
];
startMetricsCollection(agentAddresses, 15000);
startMetricsServer(9464);
}
2.5 Prometheus 采集配置
# prometheus.yml
global:
scrape_interval: 15s
evaluation_interval: 15s
scrape_timeout: 10s
scrape_configs:
# Agent 网关指标
- job_name: 'agent-gateway'
static_configs:
- targets:
- 'agent-gateway-01:9464'
- 'agent-gateway-02:9464'
- 'agent-gateway-03:9464'
labels:
environment: 'production'
chain: 'msg-chain-1'
# 独立 Agent 实例
- job_name: 'agent-instance'
kubernetes_sd_configs:
- role: pod
selectors:
- role: pod
label: 'app.kubernetes.io/component=ai-agent'
relabel_configs:
- source_labels: [__meta_kubernetes_pod_label_agent_id]
target_label: agent_id
- source_labels: [__meta_kubernetes_pod_annotation_agent_name]
target_label: agent_name
metric_relabel_configs:
- source_labels: [__name__]
regex: 'msg_agent_.*'
action: keep
# MSG Chain 节点
- job_name: 'msg-chain-node'
static_configs:
- targets:
- 'validator-01:26660'
- 'validator-02:26660'
- 'rpc-node-01:26660'
labels:
chain_id: 'msg-chain-1'
# Cosmos SDK 模块指标
- job_name: 'cosmos-modules'
static_configs:
- targets:
- 'validator-01:26660'
metrics_path: '/cosmos/metrics'
# LLM 提供商指标
- job_name: 'llm-provider'
static_configs:
- targets:
- 'llm-proxy:9465'
3. 日志(Logging)体系
3.1 结构化日志基础
使用 structlog 实现结构化日志,统一格式,方便 Loki 解析:
# log_config.py
import structlog
import logging
import sys
import uuid
from datetime import datetime, timezone
from typing import Optional
from opentelemetry import trace
def add_correlation_id(logger, method_name, event_dict):
"""添加关联 ID,优先使用 trace ID"""
span = trace.get_current_span()
ctx = span.get_span_context() if span else None
if ctx and ctx.trace_id:
event_dict['trace_id'] = format(ctx.trace_id, '032x')
event_dict['span_id'] = format(ctx.span_id, '016x')
elif 'correlation_id' not in event_dict:
event_dict['correlation_id'] = uuid.uuid4().hex[:16]
return event_dict
def add_timestamp(logger, method_name, event_dict):
event_dict['timestamp'] = datetime.now(timezone.utc).isoformat()
return event_dict
def add_agent_context(logger, method_name, event_dict):
"""从线程/协程本地存储中读取 Agent 上下文"""
from contextvars import ContextVar
agent_id_var: ContextVar[Optional[str]] = ContextVar('agent_id', default=None)
session_id_var: ContextVar[Optional[str]] = ContextVar('session_id', default=None)
agent_id = agent_id_var.get()
session_id = session_id_var.get()
if agent_id:
event_dict['agent_id'] = agent_id
if session_id:
event_dict['session_id'] = session_id
return event_dict
def setup_logging(
environment: str = 'production',
log_level: str = 'INFO',
json_output: bool = True,
):
"""配置结构化日志"""
timestamper = structlog.processors.TimeStamper(fmt='iso')
shared_processors = [
add_correlation_id,
add_timestamp,
add_agent_context,
structlog.stdlib.add_log_level,
structlog.stdlib.add_logger_name,
structlog.stdlib.PositionalArgumentsArgsFormatter(),
]
if json_output:
processors = shared_processors + [
structlog.processors.StackInfoRenderer(),
structlog.processors.format_exc_info,
structlog.processors.UnicodeDecoder(),
structlog.processors.JSONRenderer(),
]
else:
processors = shared_processors + [
structlog.stdlib.ProcessorFormatter.wrap_for_formatter,
]
structlog.configure(
processors=processors,
wrapper_class=structlog.stdlib.BoundLogger,
context_class=dict,
logger_factory=structlog.stdlib.LoggerFactory(),
cache_logger_on_first_use=True,
)
handler = logging.StreamHandler(sys.stdout)
handler.setLevel(getattr(logging, log_level.upper(), logging.INFO))
if json_output:
handler.setFormatter(
structlog.stdlib.ProcessorFormatter(
processors=[
structlog.stdlib.ProcessorFormatter.remove_processors_meta,
structlog.processors.JSONRenderer(),
],
)
)
else:
handler.setFormatter(
structlog.stdlib.ProcessorFormatter(
processors=[
structlog.stdlib.ProcessorFormatter.remove_processors_meta,
structlog.dev.ConsoleRenderer(),
],
)
)
root_logger = logging.getLogger()
root_logger.addHandler(handler)
root_logger.setLevel(getattr(logging, log_level.upper(), logging.INFO))
logger = structlog.get_logger()
logger.info('logging_initialized', environment=environment, log_level=log_level)
return logger
3.2 Agent 核心日志模块
# agent_logger.py
import structlog
from typing import Any, Dict, Optional
from contextvars import ContextVar
from dataclasses import dataclass, field
from enum import Enum
from datetime import datetime, timezone
from log_config import setup_logging
class AgentEventType(str, Enum):
"""Agent 事件类型枚举"""
LIFECYCLE = 'lifecycle'
TRANSACTION = 'transaction'
A2A_MESSAGE = 'a2a_message'
PAYMENT = 'payment'
CONSTITUTION = 'constitution'
LLM_CALL = 'llm_call'
TASK = 'task'
ERROR = 'error'
SECURITY = 'security'
@dataclass
class AgentLogContext:
"""Agent 日志上下文"""
agent_id: str
agent_name: str
agent_version: str
agent_address: str
chain_id: str = 'msg-chain-1'
session_id: Optional[str] = None
task_id: Optional[str] = None
parent_agent_id: Optional[str] = None
extra: Dict[str, Any] = field(default_factory=dict)
class AgentLogger:
"""Agent 日志记录器,封装 structlog"""
def __init__(self, context: AgentLogContext):
self.context = context
self._agent_id_var = ContextVar('agent_id', default=None)
self._session_id_var = ContextVar('session_id', default=None)
self._task_id_var = ContextVar('task_id', default=None)
self._log = structlog.get_logger()
def _bind_context(self) -> Dict[str, Any]:
return {
'agent_id': self.context.agent_id,
'agent_name': self.context.agent_name,
'agent_version': self.context.agent_version,
'agent_address': self.context.agent_address,
'chain_id': self.context.chain_id,
'session_id': self._session_id_var.get() or self.context.session_id,
'task_id': self._task_id_var.get() or self.context.task_id,
'parent_agent_id': self.context.parent_agent_id,
**self.context.extra,
}
def set_session(self, session_id: str):
self._session_id_var.set(session_id)
def set_task(self, task_id: str):
self._task_id_var.set(task_id)
def lifecycle(self, action: str, status: str, **kwargs):
self._log.info('agent_lifecycle', event_type=AgentEventType.LIFECYCLE.value,
action=action, status=status, **self._bind_context(), **kwargs)
def transaction(
self,
tx_hash: str,
tx_type: str,
status: str,
gas_used: Optional[int] = None,
gas_price: Optional[int] = None,
error: Optional[str] = None,
**kwargs
):
self._log.info('agent_transaction', event_type=AgentEventType.TRANSACTION.value,
tx_hash=tx_hash, tx_type=tx_type, status=status,
gas_used=gas_used, gas_price=gas_price, error=error,
**self._bind_context(), **kwargs)
def a2a_message(
self,
message_id: str,
sender_id: str,
receiver_id: str,
protocol: str,
message_type: str,
size_bytes: int,
outcome: str = 'sent',
**kwargs
):
self._log.info('agent_a2a', event_type=AgentEventType.A2A_MESSAGE.value,
message_id=message_id, sender_id=sender_id,
receiver_id=receiver_id, protocol=protocol,
message_type=message_type, size_bytes=size_bytes,
outcome=outcome, **self._bind_context(), **kwargs)
def payment(
self,
session_id: str,
payment_method: str,
amount_amsg: float,
currency: str = 'amsg',
status: str = 'pending',
counterparty: Optional[str] = None,
error: Optional[str] = None,
**kwargs
):
self._log.info('agent_payment', event_type=AgentEventType.PAYMENT.value,
session_id=session_id, payment_method=payment_method,
amount_amsg=amount_amsg, currency=currency, status=status,
counterparty=counterparty, error=error,
**self._bind_context(), **kwargs)
def constitution_check(
self,
check_id: str,
articles: list,
result: str = 'pass',
violations: Optional[list] = None,
duration_ms: Optional[float] = None,
**kwargs
):
level = 'warning' if result == 'violation' else 'info'
log_method = getattr(self._log, level)
log_method('agent_constitution', event_type=AgentEventType.CONSTITUTION.value,
check_id=check_id, articles=articles, result=result,
violations=violations, duration_ms=duration_ms,
**self._bind_context(), **kwargs)
def llm_call(
self,
model: str,
provider: str,
prompt_tokens: int,
completion_tokens: int,
latency_ms: float,
cost_usd: float,
success: bool = True,
error: Optional[str] = None,
**kwargs
):
self._log.info('agent_llm', event_type=AgentEventType.LLM_CALL.value,
model=model, provider=provider,
prompt_tokens=prompt_tokens,
completion_tokens=completion_tokens,
total_tokens=prompt_tokens + completion_tokens,
latency_ms=latency_ms, cost_usd=cost_usd,
success=success, error=error,
**self._bind_context(), **kwargs)
def task_event(
self,
task_id: str,
task_type: str,
status: str,
duration_s: Optional[float] = None,
outcome: Optional[str] = None,
error: Optional[str] = None,
**kwargs
):
self._log.info('agent_task', event_type=AgentEventType.TASK.value,
task_id=task_id, task_type=task_type, status=status,
duration_s=duration_s, outcome=outcome, error=error,
**self._bind_context(), **kwargs)
def security_event(
self,
event: str,
severity: str = 'info',
source_ip: Optional[str] = None,
details: Optional[Dict] = None,
**kwargs
):
log_method = getattr(self._log, 'warning' if severity == 'high' else 'info')
log_method('agent_security', event_type=AgentEventType.SECURITY.value,
event=event, severity=severity, source_ip=source_ip,
details=details, **self._bind_context(), **kwargs)
def error(
self,
message: str,
error_type: str,
details: Optional[Dict] = None,
exception: Optional[Exception] = None,
**kwargs
):
self._log.error('agent_error', event_type=AgentEventType.ERROR.value,
message=message, error_type=error_type, details=details,
exc_info=exception, **self._bind_context(), **kwargs)
def debug(self, event: str, **kwargs):
self._log.debug(event, **self._bind_context(), **kwargs)
def info(self, event: str, **kwargs):
self._log.info(event, **self._bind_context(), **kwargs)
def warning(self, event: str, **kwargs):
self._log.warning(event, **self._bind_context(), **kwargs)
# 全局 Agent 日志上下文
_agent_log_context_var: ContextVar[Optional[AgentLogContext]] = ContextVar('agent_log_context', default=None)
def get_current_agent_logger() -> Optional[AgentLogger]:
"""获取当前线程/协程的 Agent 日志记录器"""
ctx = _agent_log_context_var.get()
if ctx:
return AgentLogger(ctx)
return None
def set_current_agent_context(context: AgentLogContext):
"""设置当前线程/协程的 Agent 日志上下文"""
_agent_log_context_var.set(context)
3.3 日志过滤与采样
# log_filter.py
import re
import structlog
from typing import Dict, Any
class SensitiveDataFilter:
"""过滤敏感数据,防止密钥和地址泄露到日志中"""
SENSITIVE_PATTERNS = [
(r'(msg1[qpzry9x8gf2tvdw0s3jn54khce6mua7l]{38})', lambda m: m.group(1)[:6] + '...' + m.group(1)[-4:]),
(r'(secret_key["\']?\s*[:=]\s*["\'])([^"\']+)', lambda m: m.group(1) + '***REDACTED***'),
(r'(private_key["\']?\s*[:=]\s*["\'])([^"\']+)', lambda m: m.group(1) + '***REDACTED***'),
(r'(mnemonic["\']?\s*[:=]\s*["\'])([^"\']+)', lambda m: m.group(1) + '***REDACTED***'),
(r'(password["\']?\s*[:=]\s*["\'])([^"\']+)', lambda m: m.group(1) + '***REDACTED***'),
(r'(api_key["\']?\s*[:=]\s*["\'])([^"\']+)', lambda m: m.group(1) + '***REDACTED***'),
(r'(0x[a-fA-F0-9]{64})', lambda m: m.group(0)[:10] + '...' + m.group(0)[-6:]),
]
def __init__(self):
self.processor = structlog.get_config()['processors']
def __call__(self, logger, method_name, event_dict: Dict[str, Any]) -> Dict[str, Any]:
for key, value in list(event_dict.items()):
if isinstance(value, str):
for pattern, repl in self.SENSITIVE_PATTERNS:
event_dict[key] = re.sub(pattern, repl, value)
elif isinstance(value, dict):
self._sanitize_dict(value)
return event_dict
def _sanitize_dict(self, d: Dict[str, Any]):
for key, value in list(d.items()):
if isinstance(value, str):
for pattern, repl in self.SENSITIVE_PATTERNS:
d[key] = re.sub(pattern, repl, value)
elif isinstance(value, dict):
self._sanitize_dict(value)
class LogLevelFilter:
"""基于规则动态调整日志级别"""
def __init__(self):
self.noisy_patterns = [
re.compile(r'heartbeat'),
re.compile(r'health_check'),
re.compile(r'keepalive'),
re.compile(r'metrics_collection'),
]
def __call__(self, logger, method_name, event_dict: Dict[str, Any]) -> Dict[str, Any]:
event = str(event_dict.get('event', ''))
is_noisy = any(p.search(event) for p in self.noisy_patterns)
if is_noisy and method_name in ('debug', 'info'):
return None # 丢弃该日志
return event_dict
class SamplingFilter:
"""高吞吐日志采样"""
def __init__(self, sample_rate: float = 0.1, max_per_second: int = 100):
self.sample_rate = sample_rate
self.max_per_second = max_per_second
self.counter = 0
def __call__(self, logger, method_name, event_dict: Dict[str, Any]) -> Dict[str, Any]:
event_type = event_dict.get('event_type')
# 关键事件类型不采样
if event_type in ('security', 'error', 'constitution_violation'):
return event_dict
self.counter += 1
if self.counter > self.max_per_second:
should_sample = hash(str(event_dict)) % 100 < (self.sample_rate * 100)
if not should_sample:
return None # 丢弃采样外日志
return event_dict
3.4 日志聚合 Loki 配置
# loki-config.yaml
auth_enabled: false
server:
http_listen_port: 3100
grpc_listen_port: 9095
ingester:
lifecycler:
ring:
kvstore:
store: inmemory
replication_factor: 1
chunk_idle_period: 15m
chunk_block_size: 262144
chunk_retain_period: 5m
max_transfer_retries: 3
ingester_client:
remote_timeout: 10s
grpc_client_config:
max_recv_msg_size: 67108864
schema_config:
configs:
- from: 2025-01-01
store: boltdb-shipper
object_store: filesystem
schema: v11
index:
prefix: index_
period: 24h
storage_config:
boltdb_shipper:
active_index_directory: /data/loki/index
cache_location: /data/loki/index_cache
cache_ttl: 24h
shared_store: filesystem
filesystem:
directory: /data/loki/chunks
limits_config:
enforce_metric_name: false
reject_old_samples: true
reject_old_samples_max_age: 168h
max_entries_limit_per_query: 5000
max_line_size: 256000
ingestion_rate_mb: 10
ingestion_burst_size_mb: 20
per_stream_rate_limit: 3MB
per_stream_burst_limit: 5MB
compactor:
working_directory: /data/loki/compactor
shared_store: filesystem
compaction_interval: 10m
retention_enabled: true
retention_delete_delay: 2h
retention_delete_worker_count: 150
chunk_store_config:
max_look_back_period: 336h # 14 天
table_manager:
retention_deletes_enabled: true
retention_period: 720h # 30 天
ruler:
storage:
type: local
local:
directory: /data/loki/rules
rule_path: /data/loki/rules-temp
alertmanager_url: http://alertmanager:9093
ring:
kvstore:
store: inmemory
enable_api: true
3.5 Promtail 采集配置
# promtail-config.yaml
clients:
- url: http://loki:3100/loki/api/v1/push
scrape_configs:
- job_name: agent-containers
docker_sd_configs:
- host: unix:///var/run/docker.sock
refresh_interval: 5s
filters:
- name: label
values: ['com.docker.compose.project=agent-fleet']
relabel_configs:
- source_labels: ['__meta_docker_container_name']
target_label: container_name
- source_labels: ['__meta_docker_container_label_agent_id']
target_label: agent_id
- source_labels: ['__meta_docker_container_label_agent_name']
target_label: agent_name
- source_labels: ['__meta_docker_container_log_stream']
target_label: log_stream
pipeline_stages:
- json:
expressions:
event: event
event_type: event_type
agent_id: agent_id
level: level
timestamp: timestamp
trace_id: trace_id
session_id: session_id
tx_hash: tx_hash
- timestamp:
source: timestamp
format: RFC3339
- labels:
event_type:
agent_id:
level:
- job_name: node-logs
static_configs:
- targets: [localhost]
labels:
job: msg-chain-node
__path__: /var/log/msg-chain/*.log
pipeline_stages:
- regex:
expression: '^(?P<timestamp>\d{4}-\d{2}-\d{2}T\d{2}:\d{2}:\d{2}.\d+Z)\s+(?P<level>[A-Z]+)\s+(?P<module>[^\s]+)\s+(?P<message>.+)$'
- timestamp:
source: timestamp
format: RFC3339Nano
- labels:
level:
module:
- job_name: system-logs
static_configs:
- targets: [localhost]
labels:
job: system
__path__: /var/log/syslog
3.6 日志查询示例
以下是在 Loki 中查询 Agent 相关日志的 LogQL 示例:
# 查询指定 Agent 的所有错误日志
{agent_id="msg1agent000000000000000000000000000001", level="error"} |= ""
# 查询交易失败的日志
{event_type="transaction"} |= "failed"
# 查询宪法违规日志
{event_type="constitution"} |= "violation"
# 按会话 ID 关联日志
{session_id="sess_abc123"}
# 查询 gas 消耗异常
{event_type="transaction"} | json | gas_used > 1000000
# 按时间范围聚合
{event_type="payment"} | json
| rate by (status) [5m]
# 查询特定 Trace 的完整调用链
{trace_id="abc123def456"} |= ""
# 查询 A2A 通信延迟 > 5s 的日志
{event_type="a2a_message"} | json | latency_ms > 5000
# 查询最近的 Agent 生命周期事件
{agent_id="msg1agent*", event_type="lifecycle"} |= "" | topk 100 by (timestamp)
4. 链路追踪(Tracing)
4.1 OpenTelemetry 配置
# tracing_config.py
from opentelemetry import trace
from opentelemetry.sdk.trace import TracerProvider
from opentelemetry.sdk.trace.export import BatchSpanProcessor, SimpleSpanProcessor
from opentelemetry.exporter.otlp.proto.grpc.exporter import OTLPSpanExporter
from opentelemetry.sdk.resources import Resource, SERVICE_NAME, DEPLOYMENT_ENVIRONMENT
from opentelemetry.instrumentation.requests import RequestsInstrumentor
from opentelemetry.instrumentation.aiohttp_client import AioHttpClientInstrumentor
from opentelemetry.instrumentation.grpc import GrpcInstrumentorClient
from opentelemetry.propagators.jaeger import JaegerPropagator
from opentelemetry.propagators.composite import CompositePropagator
from opentelemetry.trace.propagation.tracecontext import TraceContextTextMapPropagator
from typing import Optional
def setup_tracing(
service_name: str,
environment: str = 'production',
agent_id: Optional[str] = None,
otlp_endpoint: str = 'http://tempo:4317',
sample_rate: float = 1.0,
use_console: bool = False,
) -> trace.Tracer:
"""配置 OpenTelemetry 追踪"""
attributes = {
SERVICE_NAME: service_name,
DEPLOYMENT_ENVIRONMENT: environment,
'chain.id': 'msg-chain-1',
'chain.bech32_prefix': 'msg',
}
if agent_id:
attributes['agent.id'] = agent_id
resource = Resource.create(attributes)
tracer_provider = TracerProvider(
resource=resource,
)
# OTLP gRPC 导出器
otlp_exporter = OTLPSpanExporter(
endpoint=otlp_endpoint,
insecure=True,
timeout=10,
)
if use_console:
from opentelemetry.sdk.trace.export import ConsoleSpanExporter
span_processor = SimpleSpanProcessor(ConsoleSpanExporter())
else:
span_processor = BatchSpanProcessor(
otlp_exporter,
max_queue_size=2048,
max_export_batch_size=512,
schedule_delay_millis=5000,
export_timeout_millis=30000,
)
tracer_provider.add_span_processor(span_processor)
trace.set_tracer_provider(tracer_provider)
# 设置传播器
propagator = CompositePropagator([
TraceContextTextMapPropagator(),
JaegerPropagator(),
])
trace.set_tracer_provider(tracer_provider)
# 自动埋点
RequestsInstrumentor().instrument()
AioHttpClientInstrumentor().instrument()
GrpcInstrumentorClient().instrument()
tracer = tracer_provider.get_tracer(service_name, '1.0.0')
return tracer
def get_agent_tracer(agent_id: str, agent_name: str) -> trace.Tracer:
"""获取 Agent 专属 Tracer"""
return setup_tracing(
service_name=f'agent-{agent_name}',
agent_id=agent_id,
environment='production',
)
4.2 Agent 关键操作埋点
# agent_tracing.py
import json
import time
from typing import Optional, Dict, Any
from opentelemetry import trace
from opentelemetry.trace import SpanKind, Status, StatusCode
from opentelemetry.propagate import inject, extract
from tracing_config import setup_tracing
class AgentTracer:
"""Agent 业务追踪器"""
def __init__(self, agent_id: str, agent_name: str):
self.agent_id = agent_id
self.agent_name = agent_name
self.tracer = setup_tracing(
service_name=f'agent-{agent_name}',
agent_id=agent_id,
)
def start_task_span(self, task_id: str, task_type: str, input_data: Optional[Dict] = None):
"""启动任务追踪 Span"""
span = self.tracer.start_span(
name=f'task.{task_type}',
kind=SpanKind.SERVER,
attributes={
'agent.id': self.agent_id,
'agent.name': self.agent_name,
'task.id': task_id,
'task.type': task_type,
'task.input': json.dumps(input_data, default=str)[:2000] if input_data else '',
},
)
return span
def end_task_span(self, span, outcome: str, error: Optional[str] = None):
"""结束任务追踪 Span"""
if error:
span.set_status(Status(StatusCode.ERROR, error))
span.set_attribute('task.error', error)
else:
span.set_status(Status(StatusCode.OK))
span.set_attribute('task.outcome', outcome)
span.end()
def trace_transaction(self, tx_hash: str, tx_type: str, tx_data: Dict[str, Any]):
"""追踪链上交易"""
with self.tracer.start_span(
name=f'transaction.{tx_type}',
kind=SpanKind.CLIENT,
attributes={
'agent.id': self.agent_id,
'tx.hash': tx_hash,
'tx.type': tx_type,
'tx.data': json.dumps(tx_data, default=str)[:2000],
},
) as span:
return span
def trace_a2a_message(
self,
message_id: str,
sender_id: str,
receiver_id: str,
protocol: str,
message_type: str,
payload: Optional[Dict] = None,
):
"""追踪 A2A 消息传递"""
with self.tracer.start_span(
name=f'a2a.{protocol}.{message_type}',
kind=SpanKind.PRODUCER if sender_id == self.agent_id else SpanKind.CONSUMER,
attributes={
'agent.id': self.agent_id,
'message.id': message_id,
'a2a.sender': sender_id,
'a2a.receiver': receiver_id,
'a2a.protocol': protocol,
'a2a.message_type': message_type,
'a2a.payload_size': len(json.dumps(payload or {}, default=str)),
},
) as span:
# 注入追踪上下文到消息头
carrier: Dict[str, str] = {}
inject(carrier)
span.set_attribute('a2a.trace_carrier', json.dumps(carrier))
return span
def trace_payment(
self,
session_id: str,
payment_method: str,
amount: float,
currency: str = 'amsg',
):
"""追踪支付流程"""
with self.tracer.start_span(
name='payment.process',
kind=SpanKind.INTERNAL,
attributes={
'agent.id': self.agent_id,
'payment.session_id': session_id,
'payment.method': payment_method,
'payment.amount': amount,
'payment.currency': currency,
},
) as span:
return span
def trace_constitution_check(self, articles: list, input_text: str):
"""追踪宪法合规检查"""
with self.tracer.start_span(
name='constitution.check',
kind=SpanKind.INTERNAL,
attributes={
'agent.id': self.agent_id,
'constitution.articles': ','.join(articles),
'constitution.input_length': len(input_text),
},
) as span:
return span
def trace_llm_call(
self,
model: str,
provider: str,
prompt: str,
max_tokens: int = 4096,
temperature: float = 0.7,
):
"""追踪 LLM 调用"""
with self.tracer.start_span(
name=f'llm.{provider}.{model}',
kind=SpanKind.CLIENT,
attributes={
'agent.id': self.agent_id,
'llm.model': model,
'llm.provider': provider,
'llm.prompt_length': len(prompt),
'llm.max_tokens': max_tokens,
'llm.temperature': temperature,
},
) as span:
return span
def extract_context_from_message(self, carrier: Dict[str, str]):
"""从传入消息中提取追踪上下文"""
ctx = extract(carrier)
return ctx
def inject_context_to_message(self) -> Dict[str, str]:
"""将当前追踪上下文注入到传出消息"""
carrier: Dict[str, str] = {}
inject(carrier)
return carrier
4.3 A2A 通信追踪传播
# a2a_tracing.py
from opentelemetry import trace
from opentelemetry.trace import SpanKind, Status, StatusCode
from opentelemetry.propagate import inject, extract
from opentelemetry.propagators.textmap import TextMapPropagator
from typing import Dict, Optional, Any
import json
import uuid
class A2ATracePropagator:
"""A2A 协议追踪上下文传播器"""
TRACE_HEADER = 'X-Agent-Trace-Context'
W3C_TRACE_PARENT = 'traceparent'
W3C_TRACE_STATE = 'tracestate'
@staticmethod
def inject_headers() -> Dict[str, str]:
"""生成需要注入到 A2A 消息头的追踪 Header"""
headers: Dict[str, str] = {}
inject(headers)
# 同时注入自定义 Agent Trace Header
span = trace.get_current_span()
if span:
ctx = span.get_span_context()
if ctx.is_valid:
headers[A2ATracePropagator.TRACE_HEADER] = json.dumps({
'trace_id': format(ctx.trace_id, '032x'),
'span_id': format(ctx.span_id, '016x'),
'trace_flags': ctx.trace_flags,
})
return headers
@staticmethod
def extract_headers(headers: Dict[str, str]):
"""从 A2A 消息头中提取追踪上下文"""
ctx = extract(headers)
return ctx
@staticmethod
def create_a2a_span(
message_id: str,
sender_id: str,
receiver_id: str,
protocol: str,
message_type: str,
parent_context=None,
):
"""创建 A2A 消息 Span"""
tracer = trace.get_tracer(__name__)
headers = A2ATracePropagator.inject_headers()
span = tracer.start_span(
name=f'a2a.{protocol}.{message_type}',
kind=SpanKind.PRODUCER,
attributes={
'a2a.message_id': message_id,
'a2a.sender': sender_id,
'a2a.receiver': receiver_id,
'a2a.protocol': protocol,
'a2a.message_type': message_type,
'a2a.trace_headers': json.dumps(headers),
},
)
return span
class A2AMessageMonitor:
"""A2A 消息监控器,整合追踪和指标"""
def __init__(self, agent_id: str, tracer: trace.Tracer):
self.agent_id = agent_id
self.tracer = tracer
self.metrics = {} # 实际会引入 Prometheus 客户端
def on_message_send(
self,
receiver_id: str,
protocol: str,
message_type: str,
payload: Dict[str, Any],
) -> Dict[str, str]:
"""发送 A2A 消息时的监控埋点"""
message_id = uuid.uuid4().hex
start_time = time.time()
span = self.tracer.start_span(
name=f'a2a.{protocol}.{message_type}.send',
kind=SpanKind.PRODUCER,
attributes={
'agent_id': self.agent_id,
'receiver_id': receiver_id,
'protocol': protocol,
'message_type': message_type,
'message_id': message_id,
'payload_size': len(json.dumps(payload, default=str)),
},
)
# 注入追踪上下文
headers = A2ATracePropagator.inject_headers()
span.set_attribute('trace_headers', json.dumps(headers))
message_metadata = {
'message_id': message_id,
'trace_headers': headers,
'sender_id': self.agent_id,
'timestamp': time.time(),
}
span.end()
return message_metadata
def on_message_recv(
self,
sender_id: str,
protocol: str,
message_type: str,
metadata: Dict[str, Any],
):
"""接收 A2A 消息时的监控埋点"""
trace_headers = metadata.get('trace_headers', {})
ctx = A2ATracePropagator.extract_headers(trace_headers)
span = self.tracer.start_span(
name=f'a2a.{protocol}.{message_type}.recv',
kind=SpanKind.CONSUMER,
context=ctx,
attributes={
'agent_id': self.agent_id,
'sender_id': sender_id,
'protocol': protocol,
'message_type': message_type,
'message_id': metadata.get('message_id', ''),
},
)
# 计算延迟
send_timestamp = metadata.get('timestamp', 0)
if send_timestamp:
latency = time.time() - send_timestamp
span.set_attribute('a2a.latency_seconds', latency)
span.end()
return span
4.4 Jaeger/Tempo 存储配置
# tempo-config.yaml
server:
http_listen_port: 3200
grpc_listen_port: 4317
distributor:
receivers:
otlp:
protocols:
grpc:
endpoint: 0.0.0.0:4317
http:
endpoint: 0.0.0.0:4318
jaeger:
protocols:
thrift_http:
endpoint: 0.0.0.0:14268
thrift_binary:
endpoint: 0.0.0.0:6832
thrift_compact:
endpoint: 0.0.0.0:6831
zipkin:
endpoint: 0.0.0.0:9411
ingester:
max_block_duration: 30m
trace_idle_period: 30s
lifecycler:
ring:
kvstore:
store: inmemory
replication_factor: 1
compactor:
compaction:
block_retention: 336h # 14 天
storage:
trace:
backend: s3
s3:
bucket: tempo-traces
endpoint: minio:9000
access_key: tempo
secret_key: temposecret
insecure: true
pool:
max_workers: 100
queue_depth: 10000
querier:
max_concurrent_queries: 10
search:
max_duration: 168h # 7 天
query_frontend:
max_outstanding_per_tenant: 100
search:
max_spans_per_span_set: 10000
metrics_generator:
registry:
external_labels:
source: tempo
cluster: msg-chain
storage:
path: /data/tempo/generator
remote_write:
- url: http://prometheus:9090/api/v1/write
send_exemplars: true
overrides:
max_traces_per_user: 10000
max_bytes_per_trace: 5000000
ingestion_rate_strategy: local
ingestion_rate_limit_bytes: 50000000
ingestion_burst_size_bytes: 100000000
4.5 Trace 查询与分析
# trace_query.py
from opentelemetry.trace import SpanKind
from typing import Dict, List, Optional
import requests
import json
class TempoClient:
"""Tempo 查询客户端"""
def __init__(self, base_url: str = 'http://tempo:3200'):
self.base_url = base_url
def query_trace(self, trace_id: str) -> Optional[Dict]:
"""按 Trace ID 查询完整链路"""
url = f'{self.base_url}/api/traces/{trace_id}'
resp = requests.get(url)
if resp.status_code == 200:
return resp.json()
return None
def search_traces(
self,
service_name: str = None,
tags: Dict[str, str] = None,
min_duration: str = None,
max_duration: str = None,
start_time: int = None,
end_time: int = None,
limit: int = 20,
) -> List[Dict]:
"""按条件搜索 Trace"""
params = {'limit': limit}
if service_name:
params['service'] = service_name
if min_duration:
params['minDuration'] = min_duration
if max_duration:
params['maxDuration'] = max_duration
if start_time:
params['start'] = start_time
if end_time:
params['end'] = end_time
if tags:
for k, v in tags.items():
params[f'tag_{k}'] = v
url = f'{self.base_url}/api/search'
resp = requests.get(url, params=params)
if resp.status_code == 200:
return resp.json().get('traces', [])
return []
def get_service_graph(self) -> Dict:
"""获取服务依赖图"""
url = f'{self.base_url}/api/services/graph'
resp = requests.get(url)
if resp.status_code == 200:
return resp.json()
return {}
def analyze_agent_traces(
self,
agent_id: str,
span_type: str = 'task',
limit: int = 50,
) -> Dict:
"""分析 Agent 特定类型的追踪数据"""
traces = self.search_traces(
tags={'agent.id': agent_id},
limit=limit,
)
results = {
'agent_id': agent_id,
'total_traces': len(traces),
'span_type': span_type,
'traces': traces,
'stats': self._compute_stats(traces, span_type),
}
return results
def _compute_stats(self, traces: List[Dict], span_type: str) -> Dict:
"""计算追踪统计"""
durations = []
spans_counts = []
errors = 0
for trace in traces:
for span in trace.get('spans', []):
if span_type in span.get('name', ''):
duration = span.get('duration', 0) / 1_000_000 # ns to ms
durations.append(duration)
spans_counts.append(len(trace.get('spans', [])))
if span.get('status', {}).get('code') == 2: # ERROR
errors += 1
if not durations:
return {}
import statistics
return {
'count': len(durations),
'avg_duration_ms': statistics.mean(durations),
'median_duration_ms': statistics.median(durations),
'p95_duration_ms': sorted(durations)[int(len(durations) * 0.95)],
'p99_duration_ms': sorted(durations)[int(len(durations) * 0.99)],
'avg_spans_per_trace': statistics.mean(spans_counts) if spans_counts else 0,
'error_count': errors,
'error_rate': errors / len(traces) if traces else 0,
'min_duration_ms': min(durations),
'max_duration_ms': max(durations),
}
def get_agent_dependency_graph(self, agent_id: str) -> Dict:
"""获取 Agent 的依赖关系图"""
spans_by_service = {}
traces = self.search_traces(
tags={'agent.id': agent_id},
limit=200,
)
for trace in traces:
for span in trace.get('spans', []):
svc = span.get('serviceName', 'unknown')
if svc not in spans_by_service:
spans_by_service[svc] = {'count': 0, 'durations': []}
spans_by_service[svc]['count'] += 1
spans_by_service[svc]['durations'].append(
span.get('duration', 0) / 1_000_000
)
graph = {'nodes': [], 'edges': []}
for svc, data in spans_by_service.items():
graph['nodes'].append({
'id': svc,
'count': data['count'],
'avg_duration_ms': sum(data['durations']) / len(data['durations']),
})
return graph
5. 告警规则
5.1 关键告警规则
# alerts/agent-critical.yml
groups:
- name: agent-critical
interval: 30s
limit: 10
rules:
# ===== Agent 存活告警 =====
- alert: AgentDown
expr: up{job="agent-instance"} == 0
for: 1m
labels:
severity: critical
team: agent-sre
annotations:
summary: 'Agent {{ $labels.agent_id }} 已下线超过 1 分钟'
description: 'Agent {{ $labels.agent_name }} ({{ $labels.agent_id }}) 在 {{ $labels.instance }} 上已停止响应。最后检测时间: {{ $value | humanizeTimestamp }}'
runbook_url: 'https://wiki.internal/runbooks/agent-down'
- alert: AgentHighMemoryUsage
expr: msg_agent_memory_usage_bytes{job="agent-instance"} / on(agent_id) msg_agent_memory_limit_bytes > 0.85
for: 5m
labels:
severity: critical
team: agent-sre
annotations:
summary: 'Agent {{ $labels.agent_id }} 内存超限 85%'
description: 'Agent {{ $labels.agent_name }} 内存使用率 {{ $value | humanizePercentage }},当前 {{ $value | humanize1024 }}'
- alert: AgentCpuThrottling
expr: rate(process_cpu_seconds_total{job="agent-instance"}[5m]) > 0.9
for: 5m
labels:
severity: critical
team: agent-sre
annotations:
summary: 'Agent {{ $labels.agent_id }} CPU 使用率 > 90%'
# ===== 支付告警 =====
- alert: AgentPaymentHighFailureRate
expr: |
rate(msg_agent_transactions_total{job="agent-instance", tx_type="payment", status="failed"}[15m])
/
rate(msg_agent_transactions_total{job="agent-instance", tx_type="payment"}[15m])
> 0.1
for: 5m
labels:
severity: critical
team: agent-payments
annotations:
summary: 'Agent {{ $labels.agent_id }} 支付失败率 > 10%'
description: 'Agent 支付失败率在过去 15 分钟内达到 {{ $value | humanizePercentage }},请立即检查链状态和余额'
- alert: AgentPaymentStuck
expr: msg_agent_payment_sessions_total{status="pending"} > 5
for: 10m
labels:
severity: critical
team: agent-payments
annotations:
summary: 'Agent {{ $labels.agent_id }} 有 {{ $value }} 笔支付卡住超过 10 分钟'
- alert: AgentLowBalance
expr: msg_agent_balance_amsg < 1
for: 1m
labels:
severity: critical
team: agent-sre
annotations:
summary: 'Agent {{ $labels.agent_id }} 余额低于 1 amsg'
description: 'Agent 地址余额不足,可能无法支付 Gas 费用。当前余额: {{ $value }} amsg'
# ===== 宪法违规告警 =====
- alert: AgentConstitutionViolation
expr: rate(msg_agent_constitution_violations_total[5m]) > 0
for: 0m
labels:
severity: critical
team: agent-security
annotations:
summary: 'Agent {{ $labels.agent_id }} 宪法违规({{ $labels.severity }})'
description: 'Agent 违反宪法第 {{ $labels.constitution_article }} 条,严重级别: {{ $labels.severity }}。建议立即暂停该 Agent'
- alert: AgentMultipleConstitutionViolations
expr: sum by (agent_id) (increase(msg_agent_constitution_violations_total[1h])) > 3
for: 2m
labels:
severity: critical
team: agent-security
annotations:
summary: 'Agent {{ $labels.agent_id }} 1 小时内多次宪法违规'
description: 'Agent {{ $labels.agent_id }} 在最近 1 小时内触发了 {{ $value }} 次宪法违规'
# ===== A2A 通信告警 =====
- alert: AgentA2AHighLatency
expr: histogram_quantile(0.95, rate(msg_agent_a2a_latency_seconds_bucket[5m])) > 10
for: 5m
labels:
severity: critical
team: agent-sre
annotations:
summary: 'A2A 通信 P95 延迟 > 10s'
description: 'Agent {{ $labels.sender_id }} -> {{ $labels.receiver_id }} P95 延迟达到 {{ $value }}s'
- alert: AgentA2AMessageLoss
expr: |
rate(msg_agent_a2a_messages_sent[5m]) - rate(msg_agent_a2a_messages_received[5m])
> 5
for: 5m
labels:
severity: critical
team: agent-sre
annotations:
summary: 'A2A 消息丢失率异常'
description: '消息发送与接收之差 > 5,可能存在消息丢失'
# ===== 链连接告警 =====
- alert: AgentChainDisconnected
expr: msg_agent_chain_connection_status == 0
for: 30s
labels:
severity: critical
team: agent-sre
annotations:
summary: 'Agent {{ $labels.agent_id }} 与链断开连接'
description: 'Agent 无法连接到 MSG Chain ({{ $labels.chain_id }}),请检查 RPC 节点状态'
- alert: AgentChainSyncLagging
expr: abs(msg_chain_block_height - msg_agent_chain_block_height) > 10
for: 1m
labels:
severity: warning
team: agent-sre
annotations:
summary: 'Agent {{ $labels.agent_id }} 区块同步落后'
description: 'Agent 区块高度落后链上 {{ $value }} 个区块'
5.2 警告告警规则
# alerts/agent-warning.yml
groups:
- name: agent-warning
interval: 30s
limit: 20
rules:
- alert: AgentHighGasSpending
expr: histogram_quantile(0.95, rate(msg_agent_gas_spent_bucket[1h])) > 500000
for: 10m
labels:
severity: warning
team: agent-sre
annotations:
summary: 'Agent {{ $labels.agent_id }} Gas 消耗偏高'
description: 'Agent P95 Gas 消耗 {{ $value }},超过阈值 500,000'
- alert: AgentSlowTransaction
expr: histogram_quantile(0.95, rate(msg_agent_transactions_duration_seconds_bucket[15m])) > 30
for: 5m
labels:
severity: warning
team: agent-sre
annotations:
summary: 'Agent {{ $labels.agent_id }} 交易确认延迟'
description: 'P95 交易确认时间 {{ $value }}s,超过 30s 阈值'
- alert: AgentRateLimitApproaching
expr: msg_gateway_agent_rate_limit_remaining < 100
for: 5m
labels:
severity: warning
team: agent-sre
annotations:
summary: 'Agent {{ $labels.agent_id }} 速率限制即将耗尽'
description: '剩余配额: {{ $value }},请评估是否需要扩容'
- alert: AgentTaskQueueGrowing
expr: msg_agent_task_queue_depth > 50
for: 10m
labels:
severity: warning
team: agent-sre
annotations:
summary: 'Agent {{ $labels.agent_id }} 任务队列积压'
description: '任务队列深度 {{ $value }},Agent 可能处理能力不足'
- alert: AgentHighLLMLatency
expr: histogram_quantile(0.95, rate(msg_agent_llm_latency_seconds_bucket[5m])) > 30
for: 5m
labels:
severity: warning
team: agent-ai
annotations:
summary: 'LLM 调用 P95 延迟 > 30s'
description: 'Agent {{ $labels.agent_id }} 调用模型 {{ $labels.model }} 的 P95 延迟达到 {{ $value }}s'
- alert: AgentLLMCostSurge
expr: rate(msg_agent_llm_cost_total[1h]) > 0.1
for: 15m
labels:
severity: warning
team: agent-ai
annotations:
summary: 'Agent {{ $labels.agent_id }} LLM 成本激增'
description: 'LLM 调用成本已达 ${{ $value }}/h,请检查调用模式'
- alert: AgentUnexpectedShutdown
expr: msg_agent_uptime_seconds < 60
for: 30s
labels:
severity: warning
team: agent-sre
annotations:
summary: 'Agent {{ $labels.agent_id }} 异常重启'
description: 'Agent {{ $labels.agent_name }} 运行时间不足 60 秒,疑似频繁重启'
- alert: AgentTaskHighFailureRate
expr: |
rate(msg_agent_tasks_completed{outcome="failed"}[30m])
/
rate(msg_agent_tasks_completed[30m])
> 0.15
for: 10m
labels:
severity: warning
team: agent-sre
annotations:
summary: 'Agent {{ $labels.agent_id }} 任务失败率 > 15%'
description: '最近 30 分钟内任务失败率 {{ $value | humanizePercentage }}'
- alert: AgentHighOpenConnections
expr: msg_agent_open_connections > 1000
for: 5m
labels:
severity: warning
team: agent-sre
annotations:
summary: 'Agent {{ $labels.agent_id }} 连接数过高'
description: '当前连接数 {{ $value }} > 1000,可能遭受 DoS 攻击或存在连接泄漏'
5.3 链健康告警
# alerts/chain-health.yml
groups:
- name: chain-health
interval: 15s
rules:
- alert: ChainNodeDown
expr: up{job="msg-chain-node"} == 0
for: 30s
labels:
severity: critical
team: chain-infra
annotations:
summary: 'MSG Chain 节点 {{ $labels.instance }} 离线'
description: '节点 {{ $labels.instance }} 已停止响应超过 30 秒'
- alert: ChainBlockProductionStalled
expr: time() - msg_chain_block_height * avg(msg_chain_block_time_seconds) > 30
for: 1m
labels:
severity: critical
team: chain-infra
annotations:
summary: 'MSG Chain 出块停止'
description: '最近 {{ $value }} 秒未产生新区块,网络可能已停止'
- alert: ChainTxPoolBacklog
expr: msg_chain_tx_pool_size > 5000
for: 5m
labels:
severity: warning
team: chain-infra
annotations:
summary: '交易池积压 > 5000'
description: '当前待处理交易 {{ $value }} 笔,可能需要提高 Gas 或扩容'
- alert: ChainHighGasPrice
expr: msg_chain_gas_price_wei > 50000000000
for: 5m
labels:
severity: warning
team: chain-infra
annotations:
summary: 'Gas 价格异常偏高'
description: '当前 Gas 价格 {{ $value }} wei > 50 Gwei'
- alert: ChainLowPeerCount
expr: msg_chain_peer_count < 3
for: 2m
labels:
severity: critical
team: chain-infra
annotations:
summary: '节点 {{ $labels.node_id }} 对等连接 < 3'
description: '当前对等连接数 {{ $value }},节点可能被网络孤立'
- alert: ChainValidatorMissedBlocks
expr: increase(cosmos_slashing_signed_blocks_window_missed[1h]) > 10
for: 5m
labels:
severity: warning
team: chain-infra
annotations:
summary: '验证人 {{ $labels.validator }} 漏块 > 10'
description: '验证人最近 1 小时漏块 {{ $value }},可能存在稳定性问题'
5.4 告警疲劳预防
# alerts/inhibition-rules.yml
groups:
- name: inhibition-rules
rules:
# 如果节点离线,抑制所有子告警
- alert: InhibitAgentDownAlerts
expr: up{job="msg-chain-node"} == 0
labels:
severity: critical
inhibition:
target_matchers:
- severity = warning
- severity = info
source_matchers:
- alertname = ChainNodeDown
# Agent 下线时抑制非关键业务告警
- alert: InhibitAgentNonCritical
expr: up{job="agent-instance"} == 0
labels:
severity: critical
inhibition:
target_matchers:
- alertname =~ "Agent(HighGas|Slow|RateLimit|TaskQueue).*"
source_matchers:
- alertname = AgentDown
# 告警静默窗口
- alert: SilenceAgentRestart
expr: changes(msg_agent_uptime_seconds[10m]) > 0
labels:
severity: info
annotations:
summary: 'Agent {{ $labels.agent_id }} 最近重启过,静默后续告警 10 分钟'
# alertmanager.yml
route:
receiver: 'default'
group_by: ['alertname', 'agent_id', 'severity']
group_wait: 30s
group_interval: 5m
repeat_interval: 4h
routes:
- receiver: 'pagerduty-critical'
matchers:
- severity = critical
repeat_interval: 30m
continue: true
- receiver: 'slack-warning'
matchers:
- severity = warning
repeat_interval: 2h
- receiver: 'telegram-alerts'
matchers:
- team = agent-sre
continue: true
- receiver: 'discord-security'
matchers:
- team = agent-security
receivers:
- name: 'pagerduty-critical'
pagerduty_configs:
- routing_key: '${PAGERDUTY_ROUTING_KEY}'
severity: critical
description: '{{ .GroupLabels.alertname }} - {{ .GroupLabels.agent_id }}'
- name: 'slack-warning'
slack_configs:
- api_url: '${SLACK_WEBHOOK_URL}'
channel: '#agent-alerts'
title: '{{ .GroupLabels.alertname }}'
text: '{{ .CommonAnnotations.description }}'
- name: 'telegram-alerts'
telegram_configs:
- bot_token: '${TELEGRAM_BOT_TOKEN}'
chat_id: -1001234567890
message: '🚨 {{ .GroupLabels.alertname }}\n{{ .CommonAnnotations.description }}'
- name: 'discord-security'
discord_configs:
- webhook_url: '${DISCORD_WEBHOOK_URL}'
title: '🔒 {{ .GroupLabels.alertname }}'
message: '{{ .CommonAnnotations.description }}'
6. Grafana 仪表盘
6.1 Agent Fleet 总览仪表盘
{
"dashboard": {
"title": "AI Agent Fleet 总览",
"version": 1,
"time": {
"from": "now-6h",
"to": "now"
},
"refresh": "30s",
"timezone": "utc",
"panels": [
{
"title": "Agent 存活状态",
"type": "stat",
"gridPos": {
"h": 4,
"w": 4,
"x": 0,
"y": 0
},
"targets": [
{
"expr": "count(up{job=\"agent-instance\"} == 1)",
"legendFormat": "在线"
},
{
"expr": "count(up{job=\"agent-instance\"} == 0) OR vector(0)",
"legendFormat": "离线"
}
],
"options": {
"colorMode": "background",
"graphMode": "none",
"reduceOptions": {
"calcs": [
"lastNotNull"
]
}
}
},
{
"title": "总交易数",
"type": "stat",
"gridPos": {
"h": 4,
"w": 4,
"x": 4,
"y": 0
},
"targets": [
{
"expr": "sum(increase(msg_agent_transactions_total[1h]))",
"legendFormat": "1h 交易数"
}
],
"options": {
"colorMode": "value",
"graphMode": "area"
}
},
{
"title": "支付成功率",
"type": "gauge",
"gridPos": {
"h": 4,
"w": 4,
"x": 8,
"y": 0
},
"targets": [
{
"expr": "avg(msg_agent_payment_success_rate) * 100",
"legendFormat": "成功率"
}
],
"options": {
"min": 0,
"max": 100,
"thresholds": {
"mode": "absolute",
"steps": [
{ "color": "red", "value": null },
{ "color": "orange", "value": 90 },
{ "color": "green", "value": 99 }
]
}
}
},
{
"title": "Agent 任务完成率",
"type": "stat",
"gridPos": {
"h": 4,
"w": 4,
"x": 12,
"y": 0
},
"targets": [
{
"expr": "sum(rate(msg_agent_tasks_completed{outcome=\"success\"}[30m])) / sum(rate(msg_agent_tasks_completed[30m])) * 100",
"legendFormat": "完成率"
}
],
"options": {
"colorMode": "value",
"graphMode": "none"
}
},
{
"title": "当前活跃 Agent",
"type": "stat",
"gridPos": {
"h": 4,
"w": 4,
"x": 16,
"y": 0
},
"targets": [
{
"expr": "sum(msg_agent_tasks_in_flight > 0)",
"legendFormat": "活跃数"
}
]
},
{
"title": "链最新区块高度",
"type": "stat",
"gridPos": {
"h": 4,
"w": 4,
"x": 20,
"y": 0
},
"targets": [
{
"expr": "msg_chain_block_height",
"legendFormat": "高度"
}
]
},
{
"title": "各 Agent 交易率 (1h)",
"type": "barchart",
"gridPos": {
"h": 8,
"w": 12,
"x": 0,
"y": 4
},
"targets": [
{
"expr": "sum by (agent_id) (increase(msg_agent_transactions_total[1h]))",
"legendFormat": "{{ agent_id }}"
}
],
"options": {
"orientation": "horizontal",
"sort": "desc"
}
},
{
"title": "A2A 消息流量",
"type": "timeseries",
"gridPos": {
"h": 8,
"w": 12,
"x": 12,
"y": 4
},
"targets": [
{
"expr": "sum by (protocol) (rate(msg_agent_a2a_messages_sent[5m]))",
"legendFormat": "{{ protocol }}"
}
],
"options": {
"legend": {
"displayMode": "table",
"placement": "bottom"
}
}
},
{
"title": "Gas 消耗 Top 10",
"type": "table",
"gridPos": {
"h": 8,
"w": 12,
"x": 0,
"y": 12
},
"targets": [
{
"expr": "topk(10, sum by (agent_id) (increase(msg_agent_gas_spent_sum[24h]) / increase(msg_agent_gas_spent_count[24h])))",
"legendFormat": "{{ agent_id }}"
}
],
"transformations": [
{
"id": "organize",
"options": {
"indexByName": {},
"renameByName": {
"agent_id": "Agent ID",
"Value": "平均 Gas"
}
}
}
]
},
{
"title": "告警事件时间线",
"type": "state-timeline",
"gridPos": {
"h": 8,
"w": 12,
"x": 12,
"y": 12
},
"targets": [
{
"expr": "ALERTS{severity=\"critical\"}",
"legendFormat": "{{ alertname }}"
}
],
"options": {
"merge": true
}
}
]
}
}
6.2 单个 Agent 详情仪表盘
{
"dashboard": {
"title": "Agent 详情 - {{ agent_id }}",
"templating": {
"list": [
{
"name": "agent_id",
"type": "query",
"query": "label_values(msg_agent_uptime_seconds, agent_id)"
}
]
},
"panels": [
{
"title": "Agent 基本信息",
"type": "stat",
"gridPos": { "h": 3, "w": 4, "x": 0, "y": 0 },
"targets": [
{ "expr": "msg_agent_uptime_seconds{agent_id=\"$agent_id\"}", "legendFormat": "运行时长(s)" },
{ "expr": "msg_agent_chain_block_height{agent_id=\"$agent_id\"}", "legendFormat": "已同步区块" }
]
},
{
"title": "Agent 余额",
"type": "gauge",
"gridPos": { "h": 3, "w": 4, "x": 4, "y": 0 },
"targets": [
{ "expr": "msg_agent_balance_amsg{agent_id=\"$agent_id\"}", "legendFormat": "余额(amsg)" }
]
},
{
"title": "交易延迟 P95",
"type": "timeseries",
"gridPos": { "h": 8, "w": 12, "x": 0, "y": 3 },
"targets": [
{
"expr": "histogram_quantile(0.95, rate(msg_agent_transactions_duration_seconds_bucket{agent_id=\"$agent_id\"}[5m]))",
"legendFormat": "P95"
},
{
"expr": "histogram_quantile(0.50, rate(msg_agent_transactions_duration_seconds_bucket{agent_id=\"$agent_id\"}[5m]))",
"legendFormat": "P50"
}
]
},
{
"title": "任务执行分布",
"type": "timeseries",
"gridPos": { "h": 8, "w": 12, "x": 12, "y": 3 },
"targets": [
{
"expr": "sum by (task_type) (rate(msg_agent_task_duration_seconds_sum{agent_id=\"$agent_id\"}[5m])) / sum by (task_type) (rate(msg_agent_task_duration_seconds_count{agent_id=\"$agent_id\"}[5m]))",
"legendFormat": "{{ task_type }}"
}
]
},
{
"title": "A2A 通信拓扑",
"type": "nodeGraph",
"gridPos": { "h": 10, "w": 24, "x": 0, "y": 11 },
"targets": [
{
"expr": "msg_agent_a2a_messages_sent{agent_id=\"$agent_id\"}",
"legendFormat": "__auto"
}
]
},
{
"title": "资源使用",
"type": "timeseries",
"gridPos": { "h": 6, "w": 12, "x": 0, "y": 21 },
"targets": [
{ "expr": "msg_agent_memory_usage_bytes{agent_id=\"$agent_id\"}", "legendFormat": "内存" },
{ "expr": "msg_agent_cpu_usage_ratio{agent_id=\"$agent_id\"}", "legendFormat": "CPU" },
{ "expr": "msg_agent_goroutine_count{agent_id=\"$agent_id\"}", "legendFormat": "Goroutines" }
]
},
{
"title": "LLM 调用分析",
"type": "timeseries",
"gridPos": { "h": 6, "w": 12, "x": 12, "y": 21 },
"targets": [
{
"expr": "sum by (model) (rate(msg_agent_llm_calls_total{agent_id=\"$agent_id\"}[5m]))",
"legendFormat": "{{ model }}"
}
]
}
]
}
}
6.3 Chain 健康仪表盘
{
"dashboard": {
"title": "MSG Chain 健康状态",
"panels": [
{
"title": "出块进度",
"type": "timeseries",
"gridPos": { "h": 8, "w": 12, "x": 0, "y": 0 },
"targets": [
{ "expr": "msg_chain_block_height", "legendFormat": "区块高度" },
{ "expr": "deriv(msg_chain_block_height[5m]) * 3600", "legendFormat": "出块速率(块/h)" }
]
},
{
"title": "出块间隔",
"type": "timeseries",
"gridPos": { "h": 8, "w": 12, "x": 12, "y": 0 },
"targets": [
{ "expr": "msg_chain_block_time_seconds", "legendFormat": "出块间隔(s)" }
],
"options": {
"thresholds": [
{ "value": 15, "color": "red" },
{ "value": 10, "color": "orange" }
]
}
},
{
"title": "交易池状态",
"type": "timeseries",
"gridPos": { "h": 8, "w": 8, "x": 0, "y": 8 },
"targets": [
{ "expr": "msg_chain_tx_pool_size", "legendFormat": "待处理" },
{ "expr": "rate(msg_chain_tx_daily_total[5m])", "legendFormat": "交易率(tx/s)" }
]
},
{
"title": "Gas 价格",
"type": "timeseries",
"gridPos": { "h": 8, "w": 8, "x": 8, "y": 8 },
"targets": [
{ "expr": "msg_chain_gas_price_wei / 1e9", "legendFormat": "Gas 价格(Gwei)" }
],
"options": {
"thresholds": [
{ "value": 50, "color": "red" },
{ "value": 20, "color": "orange" }
]
}
},
{
"title": "验证人节点",
"type": "table",
"gridPos": { "h": 8, "w": 8, "x": 16, "y": 8 },
"targets": [
{ "expr": "msg_chain_peer_count", "legendFormat": "{{ node_id }}" }
]
},
{
"title": "活跃 Agent 数量",
"type": "timeseries",
"gridPos": { "h": 6, "w": 8, "x": 0, "y": 16 },
"targets": [
{ "expr": "msg_chain_active_agents", "legendFormat": "活跃 Agent" }
]
},
{
"title": "质押总量",
"type": "timeseries",
"gridPos": { "h": 6, "w": 8, "x": 8, "y": 16 },
"targets": [
{ "expr": "sum(msg_chain_staked_amount)", "legendFormat": "总质押" }
]
},
{
"title": "区块 Gas 利用率",
"type": "timeseries",
"gridPos": { "h": 6, "w": 8, "x": 16, "y": 16 },
"targets": [
{ "expr": "msg_chain_avg_block_gas", "legendFormat": "平均 Gas/块" }
]
},
{
"title": "链事件日志流",
"type": "logs",
"gridPos": { "h": 10, "w": 24, "x": 0, "y": 22 },
"targets": [
{
"expr": "{job=\"msg-chain-node\"} |= \"\"",
"legendFormat": ""
}
],
"options": {
"showTime": true,
"wrapLogMessage": true,
"sortOrder": "Descending"
}
}
]
}
}
6.4 支付分析仪表盘
{
"dashboard": {
"title": "Agent 支付分析",
"panels": [
{
"title": "支付总额 (24h)",
"type": "stat",
"gridPos": { "h": 4, "w": 4, "x": 0, "y": 0 },
"targets": [
{ "expr": "sum(increase(msg_agent_payment_amount_total[24h]))", "legendFormat": "总金额(amsg)" }
]
},
{
"title": "支付成功率 (24h)",
"type": "stat",
"gridPos": { "h": 4, "w": 4, "x": 4, "y": 0 },
"targets": [
{ "expr": "avg(msg_agent_payment_success_rate) * 100", "legendFormat": "成功率(%)" }
],
"options": {
"thresholds": {
"mode": "absolute",
"steps": [
{ "color": "red", "value": null },
{ "color": "yellow", "value": 95 },
{ "color": "green", "value": 99 }
]
}
}
},
{
"title": "活跃支付会话",
"type": "stat",
"gridPos": { "h": 4, "w": 4, "x": 8, "y": 0 },
"targets": [
{ "expr": "sum(msg_agent_payment_sessions_total{status=\"pending\"})", "legendFormat": "进行中" }
]
},
{
"title": "平均结算延迟",
"type": "stat",
"gridPos": { "h": 4, "w": 4, "x": 12, "y": 0 },
"targets": [
{ "expr": "avg(msg_agent_payment_settlement_latency)", "legendFormat": "延迟(s)" }
]
},
{
"title": "支付方式分布",
"type": "piechart",
"gridPos": { "h": 8, "w": 8, "x": 0, "y": 4 },
"targets": [
{ "expr": "sum by (payment_method) (increase(msg_agent_payment_sessions_total[24h]))" }
]
},
{
"title": "支付金额时间序列",
"type": "timeseries",
"gridPos": { "h": 8, "w": 16, "x": 8, "y": 4 },
"targets": [
{ "expr": "sum by (agent_id) (rate(msg_agent_payment_amount_total[5m]))" }
]
},
{
"title": "支付失败分布",
"type": "timeseries",
"gridPos": { "h": 8, "w": 12, "x": 0, "y": 12 },
"targets": [
{ "expr": "rate(msg_agent_transactions_total{tx_type=\"payment\", status=\"failed\"}[5m])" }
]
},
{
"title": "结算延迟 P95",
"type": "timeseries",
"gridPos": { "h": 8, "w": 12, "x": 12, "y": 12 },
"targets": [
{ "expr": "histogram_quantile(0.95, rate(msg_agent_payment_settlement_latency_bucket[5m]))" }
]
},
{
"title": "Agent 支付排名",
"type": "table",
"gridPos": { "h": 8, "w": 24, "x": 0, "y": 20 },
"targets": [
{ "expr": "topk(20, sum by (agent_id) (increase(msg_agent_payment_amount_total[7d])))" }
],
"transformations": [
{
"id": "organize",
"options": {
"renameByName": {
"agent_id": "Agent ID",
"Value": "7d 支付总额(amsg)"
}
}
}
]
}
]
}
}
6.5 Loki 日志面板
{
"dashboard": {
"title": "Agent 日志分析",
"panels": [
{
"title": "实时日志流",
"type": "logs",
"gridPos": { "h": 12, "w": 24, "x": 0, "y": 0 },
"targets": [
{
"expr": "{agent_id=~\"$agent_id\"} |= \"\"",
"legendFormat": ""
}
],
"options": {
"showTime": true,
"wrapLogMessage": true,
"sortOrder": "Descending",
"maxLines": 500
}
},
{
"title": "日志级别分布",
"type": "timeseries",
"gridPos": { "h": 6, "w": 8, "x": 0, "y": 12 },
"targets": [
{
"expr": "sum by (level) (rate({agent_id=~\"$agent_id\"} | json | level != \"\" [$5m]))",
"legendFormat": "{{ level }}"
}
]
},
{
"title": "事件类型分布",
"type": "piechart",
"gridPos": { "h": 6, "w": 8, "x": 8, "y": 12 },
"targets": [
{
"expr": "count by (event_type) ({agent_id=~\"$agent_id\"} | json | event_type != \"\" [$24h])",
"legendFormat": "{{ event_type }}"
}
]
},
{
"title": "错误日志趋势",
"type": "timeseries",
"gridPos": { "h": 6, "w": 8, "x": 16, "y": 12 },
"targets": [
{
"expr": "rate({agent_id=~\"$agent_id\", level=\"error\"}[5m])",
"legendFormat": "错误率"
}
]
}
]
}
}
7. 告警通知渠道
7.1 PagerDuty 集成
# pagerduty_integration.py
import os
import json
import hashlib
import hmac
import requests
from typing import Dict, Optional
from datetime import datetime
class PagerDutyClient:
"""PagerDuty 告警集成"""
def __init__(self, routing_key: str = None):
self.routing_key = routing_key or os.environ.get('PAGERDUTY_ROUTING_KEY')
self.api_url = 'https://events.pagerduty.com/v2/enqueue'
self.session = requests.Session()
def trigger_alert(
self,
summary: str,
source: str,
severity: str = 'critical',
component: Optional[str] = None,
group: Optional[str] = None,
class_: Optional[str] = None,
custom_details: Optional[Dict] = None,
dedup_key: Optional[str] = None,
):
"""触发 PagerDuty 告警"""
payload = {
'routing_key': self.routing_key,
'event_action': 'trigger',
'payload': {
'summary': summary[:1024],
'source': source,
'severity': severity,
'timestamp': datetime.utcnow().isoformat() + 'Z',
'component': component or 'ai-agent',
'group': group or 'agent-fleet',
'class': class_ or 'observability',
'custom_details': custom_details or {},
},
}
if dedup_key:
payload['dedup_key'] = dedup_key
resp = self.session.post(self.api_url, json=payload, timeout=10)
resp.raise_for_status()
return resp.json()
def acknowledge_alert(self, dedup_key: str):
"""确认告警"""
payload = {
'routing_key': self.routing_key,
'event_action': 'acknowledge',
'dedup_key': dedup_key,
}
resp = self.session.post(self.api_url, json=payload, timeout=10)
return resp.json()
def resolve_alert(self, dedup_key: str):
"""解决告警"""
payload = {
'routing_key': self.routing_key,
'event_action': 'resolve',
'dedup_key': dedup_key,
}
resp = self.session.post(self.api_url, json=payload, timeout=10)
return resp.json()
def send_agent_down_alert(self, agent_id: str, agent_name: str, instance: str):
"""发送 Agent 下线告警"""
dedup_key = f'agent-down-{agent_id}'
return self.trigger_alert(
summary=f'Agent {agent_name} ({agent_id}) 已下线',
source=instance,
severity='critical',
component='ai-agent',
group=f'agent-{agent_id}',
class_='availability',
custom_details={
'agent_id': agent_id,
'agent_name': agent_name,
'instance': instance,
'environment': 'production',
'chain': 'msg-chain-1',
},
dedup_key=dedup_key,
)
def send_constitution_violation_alert(
self,
agent_id: str,
agent_name: str,
article: str,
severity: str,
details: Dict,
):
"""发送宪法违规告警"""
dedup_key = f'constitution-{agent_id}-{article}'
return self.trigger_alert(
summary=f'Agent {agent_name} 宪法违规 - 第{article}条',
source=agent_id,
severity='critical',
component='ai-agent',
group=f'agent-{agent_id}',
class_='security',
custom_details={
'agent_id': agent_id,
'agent_name': agent_name,
'article': article,
'severity': severity,
'details': details,
},
dedup_key=dedup_key,
)
7.2 Telegram Bot 通知
// telegram-bot.ts
import axios from 'axios';
interface TelegramConfig {
botToken: string;
chatId: string;
}
interface AlertPayload {
title: string;
description: string;
severity: 'critical' | 'warning' | 'info';
agentId?: string;
timestamp: number;
fields?: Record<string, string>;
}
export class TelegramNotifier {
private readonly apiBase: string;
private readonly chatId: string;
constructor(config: TelegramConfig) {
this.apiBase = `https://api.telegram.org/bot${config.botToken}`;
this.chatId = config.chatId;
}
async sendAlert(payload: AlertPayload): Promise<void> {
const emoji = payload.severity === 'critical'
? '🚨' : payload.severity === 'warning'
? '⚠️' : 'ℹ️';
const lines = [
`${emoji} *${payload.title}*`,
'',
payload.description,
'',
`Severity: \`${payload.severity}\``,
];
if (payload.agentId) {
lines.push(`Agent: \`${payload.agentId}\``);
}
if (payload.fields) {
for (const [key, value] of Object.entries(payload.fields)) {
lines.push(`${key}: \`${value}\``);
}
}
lines.push(`Time: ${new Date(payload.timestamp).toISOString()}`);
const text = lines.join('\n');
try {
await axios.post(`${this.apiBase}/sendMessage`, {
chat_id: this.chatId,
text,
parse_mode: 'Markdown',
disable_web_page_preview: true,
});
} catch (err) {
console.error('Failed to send Telegram alert:', err);
}
}
async sendPaymentFailure(params: {
agentId: string;
sessionId: string;
amount: number;
error: string;
}): Promise<void> {
await this.sendAlert({
title: '支付失败',
description: `Agent ${params.agentId} 支付失败`,
severity: 'critical',
agentId: params.agentId,
timestamp: Date.now(),
fields: {
Session: params.sessionId,
Amount: `${params.amount} amsg`,
Error: params.error,
},
});
}
async sendAgentDown(params: {
agentId: string;
agentName: string;
instance: string;
}): Promise<void> {
await this.sendAlert({
title: `Agent 下线: ${params.agentName}`,
description: `Agent ${params.agentName} (${params.agentId}) 在 ${params.instance} 上已停止响应`,
severity: 'critical',
agentId: params.agentId,
timestamp: Date.now(),
fields: {
Instance: params.instance,
},
});
}
}
7.3 Discord Webhook
// discord-webhook.ts
import axios from 'axios';
interface EmbedField {
name: string;
value: string;
inline?: boolean;
}
interface DiscordEmbed {
title: string;
description: string;
color: number;
fields?: EmbedField[];
timestamp?: string;
footer?: { text: string };
}
export class DiscordNotifier {
private readonly webhookUrl: string;
constructor(webhookUrl: string) {
this.webhookUrl = webhookUrl;
}
private getColor(severity: string): number {
switch (severity) {
case 'critical': return 0xFF0000;
case 'warning': return 0xFFA500;
case 'info': return 0x3498DB;
default: return 0x808080;
}
}
async sendAlert(params: {
title: string;
description: string;
severity: string;
fields?: EmbedField[];
}): Promise<void> {
const embed: DiscordEmbed = {
title: params.title,
description: params.description,
color: this.getColor(params.severity),
fields: params.fields,
timestamp: new Date().toISOString(),
footer: { text: 'MSG Chain Agent Monitoring' },
};
try {
await axios.post(this.webhookUrl, {
embeds: [embed],
username: 'Agent Monitor',
avatar_url: 'https://img.nicehash.com/logo.svg',
});
} catch (err) {
console.error('Failed to send Discord alert:', err);
}
}
async sendSecurityAlert(params: {
agentId: string;
event: string;
severity: string;
details: Record<string, string>;
}): Promise<void> {
const fields: EmbedField[] = [
{ name: 'Agent ID', value: params.agentId, inline: true },
{ name: '事件', value: params.event, inline: true },
{ name: '严重级别', value: params.severity, inline: true },
...Object.entries(params.details).map(([k, v]) => ({
name: k,
value: v,
inline: true,
})),
];
await this.sendAlert({
title: `🔒 Agent 安全事件`,
description: `Agent ${params.agentId} 触发安全事件: ${params.event}`,
severity: params.severity,
fields,
});
}
}
7.4 值班轮换配置
# oncall-schedule.yml
schedule:
team: agent-sre
rotation:
- name: primary
duration: 7d # 7 天轮换
members:
- name: alice
slack: "@alice"
pagerduty: "alice@company.com"
telegram: 123456789
- name: bob
slack: "@bob"
pagerduty: "bob@company.com"
telegram: 987654321
- name: secondary
duration: 7d
members:
- name: charlie
slack: "@charlie"
pagerduty: "charlie@company.com"
telegram: 456789123
- name: diana
slack: "@diana"
pagerduty: "diana@company.com"
telegram: 321654987
escalation:
- notify: primary
timeout: 5m
- notify: secondary
timeout: 10m
- notify: manager
timeout: 20m
holiday_calendar: "Asia/Shanghai"
override_enabled: true
8. 部署架构
8.1 Docker Compose 全栈部署
# docker-compose.yml
version: '3.8'
networks:
agent-monitoring:
driver: bridge
volumes:
prometheus_data:
grafana_data:
loki_data:
tempo_data:
services:
# ===== Agent 实例 =====
agent-gateway-01:
image: msgchain/agent-gateway:latest
ports:
- "8080:8080"
- "9464:9464"
environment:
AGENT_ID: "msg1agent000000000000000000000000000001"
AGENT_NAME: "trading-bot-01"
MSG_RPC_ENDPOINT: "http://validator-01:26657"
LOG_LEVEL: "INFO"
OTEL_EXPORTER_OTLP_ENDPOINT: "http://tempo:4317"
labels:
agent_id: "msg1agent000000000000000000000000000001"
agent_name: "trading-bot-01"
networks:
- agent-monitoring
restart: unless-stopped
healthcheck:
test: ["CMD", "curl", "-f", "http://localhost:8080/health"]
interval: 15s
timeout: 5s
retries: 3
agent-gateway-02:
image: msgchain/agent-gateway:latest
ports:
- "8081:8080"
- "9465:9464"
environment:
AGENT_ID: "msg1agent000000000000000000000000000002"
AGENT_NAME: "nft-curator-01"
MSG_RPC_ENDPOINT: "http://validator-02:26657"
LOG_LEVEL: "INFO"
OTEL_EXPORTER_OTLP_ENDPOINT: "http://tempo:4317"
labels:
agent_id: "msg1agent000000000000000000000000000002"
agent_name: "nft-curator-01"
networks:
- agent-monitoring
restart: unless-stopped
agent-gateway-03:
image: msgchain/agent-gateway:latest
ports:
- "8082:8080"
- "9466:9464"
environment:
AGENT_ID: "msg1agent000000000000000000000000000003"
AGENT_NAME: "defi-analyst-01"
MSG_RPC_ENDPOINT: "http://rpc-node-01:26657"
LOG_LEVEL: "INFO"
OTEL_EXPORTER_OTLP_ENDPOINT: "http://tempo:4317"
labels:
agent_id: "msg1agent000000000000000000000000000003"
agent_name: "defi-analyst-01"
networks:
- agent-monitoring
restart: unless-stopped
# ===== 监控基础设施 =====
prometheus:
image: prom/prometheus:v2.50.0
volumes:
- ./prometheus.yml:/etc/prometheus/prometheus.yml
- ./alerts:/etc/prometheus/alerts
- prometheus_data:/prometheus
command:
- '--config.file=/etc/prometheus/prometheus.yml'
- '--storage.tsdb.path=/prometheus'
- '--storage.tsdb.retention.time=30d'
- '--web.console.libraries=/etc/prometheus/console_libraries'
- '--web.console.templates=/etc/prometheus/consoles'
- '--web.enable-lifecycle'
ports:
- "9090:9090"
networks:
- agent-monitoring
restart: unless-stopped
loki:
image: grafana/loki:2.9.0
volumes:
- ./loki-config.yaml:/etc/loki/local-config.yaml
- loki_data:/data/loki
command: -config.file=/etc/loki/local-config.yaml
ports:
- "3100:3100"
- "9095:9095"
networks:
- agent-monitoring
restart: unless-stopped
promtail:
image: grafana/promtail:2.9.0
volumes:
- ./promtail-config.yaml:/etc/promtail/config.yml
- /var/log:/var/log
- /var/run/docker.sock:/var/run/docker.sock
command: -config.file=/etc/promtail/config.yml
networks:
- agent-monitoring
restart: unless-stopped
tempo:
image: grafana/tempo:2.3.0
volumes:
- ./tempo-config.yaml:/etc/tempo.yaml
- tempo_data:/data/tempo
command: -config.file=/etc/tempo.yaml
ports:
- "3200:3200"
- "4317:4317"
- "4318:4318"
- "9411:9411"
- "14268:14268"
networks:
- agent-monitoring
restart: unless-stopped
grafana:
image: grafana/grafana:10.3.0
volumes:
- grafana_data:/var/lib/grafana
- ./grafana-dashboards:/etc/grafana/provisioning/dashboards
- ./grafana-datasources:/etc/grafana/provisioning/datasources
environment:
GF_SECURITY_ADMIN_PASSWORD: "${GRAFANA_PASSWORD:-admin}"
GF_INSTALL_PLUGINS: "grafana-piechart-panel"
GF_AUTH_ANONYMOUS_ENABLED: "false"
ports:
- "3000:3000"
networks:
- agent-monitoring
restart: unless-stopped
# ===== 告警 =====
alertmanager:
image: prom/alertmanager:v0.26.0
volumes:
- ./alertmanager.yml:/etc/alertmanager/alertmanager.yml
command:
- '--config.file=/etc/alertmanager/alertmanager.yml'
- '--storage.path=/alertmanager'
ports:
- "9093:9093"
- "9094:9094"
networks:
- agent-monitoring
restart: unless-stopped
# ===== 可选: 通知服务 =====
telegram-bot:
image: msgchain/telegram-alert-bot:latest
environment:
TELEGRAM_BOT_TOKEN: "${TELEGRAM_BOT_TOKEN}"
TELEGRAM_CHAT_ID: "${TELEGRAM_CHAT_ID}"
PROMETHEUS_ALERTMANAGER_URL: "http://alertmanager:9093"
networks:
- agent-monitoring
restart: unless-stopped
8.2 Kubernetes 部署
# agent-monitoring-ns.yaml
apiVersion: v1
kind: Namespace
metadata:
name: agent-monitoring
# agent-deployment.yaml
apiVersion: apps/v1
kind: Deployment
metadata:
name: agent-gateway
namespace: agent-monitoring
labels:
app.kubernetes.io/name: agent-gateway
app.kubernetes.io/component: ai-agent
app.kubernetes.io/part-of: msg-chain
spec:
replicas: 3
selector:
matchLabels:
app.kubernetes.io/name: agent-gateway
template:
metadata:
labels:
app.kubernetes.io/name: agent-gateway
app.kubernetes.io/component: ai-agent
annotations:
prometheus.io/scrape: "true"
prometheus.io/port: "9464"
prometheus.io/path: "/metrics"
spec:
containers:
- name: agent
image: msgchain/agent-gateway:latest
ports:
- containerPort: 8080
name: http
- containerPort: 9464
name: metrics
env:
- name: AGENT_ID
valueFrom:
fieldRef:
fieldPath: metadata.labels['agent-id']
- name: MSG_RPC_ENDPOINT
value: "http://msg-chain-rpc:26657"
- name: OTEL_EXPORTER_OTLP_ENDPOINT
value: "http://tempo:4317"
- name: LOG_LEVEL
value: "INFO"
resources:
requests:
cpu: "500m"
memory: "512Mi"
limits:
cpu: "2"
memory: "2Gi"
livenessProbe:
httpGet:
path: /health
port: 8080
initialDelaySeconds: 30
periodSeconds: 15
readinessProbe:
httpGet:
path: /health
port: 8080
initialDelaySeconds: 5
periodSeconds: 10
# prometheus-operator.yaml
apiVersion: monitoring.coreos.com/v1
kind: ServiceMonitor
metadata:
name: agent-gateway-monitor
namespace: agent-monitoring
spec:
selector:
matchLabels:
app.kubernetes.io/name: agent-gateway
endpoints:
- port: metrics
interval: 15s
path: /metrics
namespaceSelector:
matchNames:
- agent-monitoring
8.3 架构总览图
┌─────────────────────────────────────────────────────────────────┐
│ Agent Fleet │
│ ┌──────────┐ ┌──────────┐ ┌──────────┐ │
│ │ Agent 01 │ │ Agent 02 │ │ Agent 03 │ ... │
│ │ :9464 │ │ :9465 │ │ :9466 │ │
│ └────┬─────┘ └────┬─────┘ └────┬─────┘ │
│ │ │ │ │
│ └───────────────┼───────────────┘ │
│ │ │
│ ┌────────────┴────────────┐ │
│ │ Prometheus Scraper │ :9090 │
│ │ (pull metrics every │ │
│ │ 15s) │ │
│ └────────────┬────────────┘ │
│ │ │
├───────────────────────┼─────────────────────────────────────────┤
│ │ │
│ ┌────────────┴────────────┐ │
│ │ Alertmanager │ :9093 │
│ └────────────┬────────────┘ │
│ │ │
│ ┌────────────┼────────────┐ │
│ ▼ ▼ ▼ │
│ ┌──────────┐ ┌──────────┐ ┌──────────┐ │
│ │PagerDuty │ │ Telegram │ │ Discord │ │
│ └──────────┘ └──────────┘ └──────────┘ │
│ │
├──────────────────────────────┬──────────────────────────────────┤
│ │ │
│ ┌───────────────────────────┴──────────────────────────────┐ │
│ │ Grafana :3000 │ │
│ │ ┌─────────────┐ ┌─────────────┐ ┌────────────────┐ │ │
│ │ │ Agent 面板 │ │ Chain 面板 │ │ 日志面板 │ │ │
│ │ └─────────────┘ └─────────────┘ └────────────────┘ │ │
│ │ │ │
│ │ ┌──────────────────────────────────────────────────┐ │ │
│ │ │ 数据源: Prometheus / Loki / Tempo │ │ │
│ │ └──────────────────────────────────────────────────┘ │ │
│ └──────────────────────────────────────────────────────┘ │
│ ▲ │
├──────────────────────────────┼───────────────────────────────┤
│ │ │
│ ┌───────────┐ ┌───────────┴────────────┐ ┌───────────┐ │
│ │ Loki │ │ Tempo │ │ Prometheus │ │
│ │ (日志) │ │ (链路追踪) │ │ (指标) │ │
│ │ :3100 │ │ :4317 OTLP │ │ :9090 │ │
│ │ │ │ :3200 API │ │ │ │
│ └─────┬─────┘ └──────────┬────────────┘ └────────────┘ │
│ │ │ │
│ ┌─────┴─────┐ ┌───────┴────────┐ │
│ │ Promtail │ │ OpenTelemetry │ │
│ │ (日志采集) │ │ Collector │ │
│ └───────────┘ └────────────────┘ │
│ │
│ ┌──────────────────────────────────────────────────────┐ │
│ │ MSG Chain 节点 (RPC) │ │
│ │ ┌──────────┐ ┌──────────┐ ┌──────────┐ │ │
│ │ │Validator │ │Validator │ │ RPC Node│ │ │
│ │ │ 01 │ │ 02 │ │ 01 │ │ │
│ │ └──────────┘ └──────────┘ └──────────┘ │ │
│ └──────────────────────────────────────────────────────┘ │
└─────────────────────────────────────────────────────────────────┘
8.4 性能与容量规划
| 组件 | 建议配置 | 存储 | 预估规模 |
|---|---|---|---|
| Prometheus | 4 CPU / 16GB RAM | 200GB SSD (30d) | 100 个 Agent |
| Loki | 4 CPU / 8GB RAM | 500GB SSD (30d) | 50GB/天日志 |
| Tempo | 4 CPU / 8GB RAM | 200GB SSD (14d) | 100 个 Trace/s |
| Grafana | 2 CPU / 4GB RAM | 10GB | 10 并发用户 |
| Alertmanager | 1 CPU / 2GB RAM | 5GB | 100 条告警规则 |
| 单个 Agent | 2 CPU / 2GB RAM | 本地日志 10GB | 500 tx/h |
8.5 启动与验证
# 1. 准备工作
mkdir -p /opt/agent-monitoring/{alerts,grafana-dashboards,grafana-datasources}
cd /opt/agent-monitoring
# 2. 配置 Grafana 数据源
cat > grafana-datasources/prometheus.yaml << 'EOF'
apiVersion: 1
datasources:
- name: Prometheus
type: prometheus
access: proxy
url: http://prometheus:9090
isDefault: true
- name: Loki
type: loki
access: proxy
url: http://loki:3100
- name: Tempo
type: tempo
access: proxy
url: http://tempo:3200
jsonData:
tracesToLogs:
datasourceUid: loki
EOF
# 3. 启动全栈
docker compose up -d
# 4. 验证各组件
curl -s http://localhost:9090/-/ready # Prometheus
curl -s http://localhost:3100/ready # Loki
curl -s http://localhost:3200/ready # Tempo
curl -s http://localhost:3000/api/health # Grafana
# 5. 验证 Agent 指标采集
curl -s http://localhost:9464/metrics | head -50
# 6. 导入仪表盘
curl -X POST http://admin:${GRAFANA_PASSWORD}@localhost:3000/api/dashboards/db \
-d @agent-fleet-dashboard.json \
-H "Content-Type: application/json"
# 7. 配置告警
curl -X POST http://localhost:9093/-/reload
9. 附录
A. 参考工具与链接
| 工具 | 用途 | 官方文档 |
|---|---|---|
| Prometheus | 指标采集与告警 | https://prometheus.io/docs/ |
| Grafana | 可视化仪表盘 | https://grafana.com/docs/ |
| Loki | 日志聚合 | https://grafana.com/docs/loki/ |
| Tempo | 链路追踪 | https://grafana.com/docs/tempo/ |
| OpenTelemetry | 遥测标准 | https://opentelemetry.io/docs/ |
| Alertmanager | 告警管理 | https://prometheus.io/docs/alerting/ |
| PagerDuty | 事件响应 | https://developer.pagerduty.com/ |
| structlog | Python 结构化日志 | https://www.structlog.org/ |
| prom-client | Node.js Prometheus | https://github.com/siimon/prom-client |
B. MSG Chain 相关地址前缀
| 用途 | 前缀 | 示例 |
|---|---|---|
| Agent 地址 | msg1agent |
msg1agent000000000000000000000000000001 |
| 用户地址 | msg1 |
msg1useraddress00000000000000000000000 |
| 合约地址 | msg1contract |
msg1contract00000000000000000000000001 |
| 验证人地址 | msg1valoper |
msg1valoper000000000000000000000000001 |
C. 常见问题排查
Q: Agent 指标不显示?
- 检查
prometheus.yml中scrape_configs的目标地址和端口 - 确认 Agent 的
/metrics端点可以正常 curl - 检查 Prometheus Target 页面
http://prometheus:9090/targets
Q: 日志没有写入 Loki?
- 检查 Promtail 是否挂载了 Docker Socket
- 验证 Loki 的 readiness 端点
/ready - 检查 Promtail 日志中的错误信息
Q: Trace 无法关联?
- 确认 OpenTelemetry 版本兼容
- 检查
OTEL_EXPORTER_OTLP_ENDPOINT配置 - 验证 Tempo 的 gRPC 端口
4317可访问
Q: 告警没有触发通知?
- 检查 Alertmanager 配置中的路由规则
- 验证 Webhook URL 和 API Key 正确
- 查看 Alertmanager 日志
/alertmanager/logs
D. 安全注意事项
- Metrics 端点应使用 HTTP Basic Auth 或 mTLS 保护
- 不要在日志中输出私钥、助记词或 API Key
- Alertmanager Webhook URL 应使用环境变量注入
- Grafana 应开启 SSO 或 OAuth2 认证
- 定期轮换监控系统的访问凭证
- 日志和 Trace 数据包含 Agent 行为细节,应控制访问权限
E. 许可
本文档遵循 MIT 许可协议,适用于 MSG Chain 生态的 AI Agent 开发者社区。
本文档由 MSG Chain 可观测性团队维护,如有问题请提交 Issue 或联系运维值班。
