AI Agent 监控告警规则配置指南
适用链: msg-chain-1 | bech32 前缀: msg
技术栈: Prometheus + Grafana + Alertmanager
版本: v1.0⚠️ No-Go Disclaimer: MSGChain 主网裁决为 No-Go。本文件所有内容反映的是开发阶段的技术设计,不代表主网未独立核验上线状态。生产部署状态请以白皮书为准:https://msgchain.org/whitepaper/
1. 概述
1.1 为什么需要监控
AI Agent 运行在 MSG Chain 上,作为链下服务与链上智能合约交互。与普通区块链节点不同,AI Agent 的监控需要覆盖以下层面:
- 链下服务可用性:Agent 进程是否存活,是否能正常处理请求
- 链上状态一致性:Agent 的链下状态是否与链上记录匹配
- 经济指标:Agent 余额、收入、Gas 消耗等经济维度
- 性能指标:请求延迟、吞吐量、错误率
- 合规指标:Agent 行为是否符合创世章程(Constitution)
一个未经监控的 Agent 可能在余额耗尽后静默失效,或在违反章程后继续运行造成损失。本指南提供从指标暴露到告警通知的完整方案。
1.2 技术栈选型
| 组件 | 用途 | 端口 |
|---|---|---|
| Prometheus | 指标采集与存储 | 9090 |
| Grafana | 可视化仪表盘 | 3000 |
| Alertmanager | 告警路由与通知 | 9093 |
| Prometheus Client (Python) | Agent 指标暴露 | 8000 |
| Node Exporter | 主机指标 | 9100 |
1.3 Agent 指标分类
agent_requests_total # 请求量(Counter,带 capability/status 标签)
agent_request_duration_seconds # 请求延迟(Histogram)
agent_earnings_total # 累计收入(Counter,带 from_addr 标签)
agent_balance # 当前余额(Gauge)
agent_uptime_seconds # 运行时长(Gauge)
agent_constitution_violations_total # 章程违规次数(Counter,带 severity 标签)
agent_heartbeat_timestamp # 上次心跳时间戳(Gauge)
agent_gas_spent_total # Gas 消耗累计(Counter)
agent_active_tasks # 当前活跃任务数(Gauge)
agent_memory_usage_bytes # 内存使用(Gauge)
1.4 目录结构
monitoring/
├── prometheus/
│ ├── prometheus.yml # Prometheus 主配置
│ ├── alert.rules.yml # 告警规则
│ └── alertmanager.yml # Alertmanager 配置
├── grafana/
│ ├── dashboards/ # 仪表盘 JSON
│ └── datasources/ # 数据源配置
├── agent/
│ └── metrics.py # Agent 指标暴露代码
├── notifiers/
│ ├── discord.py # Discord 通知
│ ├── telegram.py # Telegram 通知
│ └── email.py # 邮件通知
└── docker-compose.monitoring.yml
2. Prometheus 指标暴露
2.1 安装 Prometheus Client
pip install prometheus-client aiohttp cosmpy
2.2 完整指标采集器
# agent/metrics.py
import asyncio
import time
import logging
from typing import Optional
from prometheus_client import (
start_http_server,
Counter,
Gauge,
Histogram,
Info,
Enum,
Summary,
generate_latest,
REGISTRY,
)
logger = logging.getLogger(__name__)
# ─── 基础请求指标 ─────────────────────────────────────────────
REQUESTS_TOTAL = Counter(
'agent_requests_total',
'Total number of requests handled by the agent',
['agent_id', 'capability', 'status'],
)
REQUEST_DURATION = Histogram(
'agent_request_duration_seconds',
'Request duration in seconds',
['agent_id', 'capability'],
buckets=[0.01, 0.05, 0.1, 0.25, 0.5, 1.0, 2.5, 5.0, 10.0, 30.0],
)
REQUEST_SIZE_BYTES = Histogram(
'agent_request_size_bytes',
'Request payload size in bytes',
['agent_id'],
buckets=[100, 500, 1000, 5000, 10000, 50000, 100000],
)
CONCURRENT_REQUESTS = Gauge(
'agent_concurrent_requests',
'Number of requests currently being processed',
['agent_id'],
)
# ─── 经济指标 ────────────────────────────────────────────────
EARNINGS_TOTAL = Counter(
'agent_earnings_total',
'Total MSG tokens earned by the agent',
['agent_id', 'from_addr', 'capability'],
)
BALANCE_GAUGE = Gauge(
'agent_balance',
'Current MSG token balance of the agent',
['agent_id', 'denom'],
)
GAS_SPENT_TOTAL = Counter(
'agent_gas_spent_total',
'Total gas spent by the agent',
['agent_id', 'tx_type'],
)
GAS_PRICE_GAUGE = Gauge(
'agent_gas_price',
'Current gas price observed by the agent',
['agent_id', 'denom'],
)
# ─── 健康指标 ────────────────────────────────────────────────
UPTIME = Gauge(
'agent_uptime_seconds',
'Agent uptime in seconds',
['agent_id'],
)
HEARTBEAT_TIMESTAMP = Gauge(
'agent_heartbeat_timestamp',
'Unix timestamp of last successful on-chain heartbeat',
['agent_id'],
)
LAST_BLOCK_HEIGHT = Gauge(
'agent_last_block_height',
'Last block height processed by the agent',
['agent_id'],
)
SYNC_STATUS = Enum(
'agent_sync_status',
'Sync status of the agent with the chain',
['agent_id'],
states=['synced', 'syncing', 'stalled', 'behind'],
)
# ─── 合规指标 ────────────────────────────────────────────────
CONSTITUTION_VIOLATIONS = Counter(
'agent_constitution_violations_total',
'Total constitution violations detected',
['agent_id', 'severity', 'violation_type'],
)
CONSTITUTION_CHECK_DURATION = Histogram(
'agent_constitution_check_duration_seconds',
'Time taken for constitution validation',
['agent_id'],
buckets=[0.001, 0.005, 0.01, 0.05, 0.1, 0.5, 1.0],
)
# ─── 资源指标 ────────────────────────────────────────────────
MEMORY_USAGE_BYTES = Gauge(
'agent_memory_usage_bytes',
'Current memory usage of the agent process',
['agent_id', 'type'],
)
CPU_USAGE_PERCENT = Gauge(
'agent_cpu_usage_percent',
'Current CPU usage of the agent process',
['agent_id'],
)
ACTIVE_TASKS = Gauge(
'agent_active_tasks',
'Number of active background tasks',
['agent_id', 'task_type'],
)
# ─── 链交互指标 ──────────────────────────────────────────────
CHAIN_TX_COUNT = Counter(
'agent_chain_tx_total',
'Total chain transactions submitted',
['agent_id', 'msg_type', 'status'],
)
CHAIN_TX_DURATION = Histogram(
'agent_chain_tx_duration_seconds',
'Time to confirm chain transactions',
['agent_id', 'msg_type'],
buckets=[0.5, 1.0, 2.0, 5.0, 10.0, 30.0, 60.0],
)
CHAIN_QUERY_COUNT = Counter(
'agent_chain_query_total',
'Total chain queries performed',
['agent_id', 'query_type'],
)
CHAIN_QUERY_DURATION = Histogram(
'agent_chain_query_duration_seconds',
'Duration of chain queries',
['agent_id', 'query_type'],
buckets=[0.01, 0.05, 0.1, 0.25, 0.5, 1.0, 2.0],
)
# ─── Agent 信息 ──────────────────────────────────────────────
AGENT_INFO = Info(
'agent',
'Static information about the agent',
)
# ─── LLM 指标 ────────────────────────────────────────────────
LLM_REQUESTS_TOTAL = Counter(
'agent_llm_requests_total',
'Total LLM API calls',
['agent_id', 'model', 'provider'],
)
LLM_REQUEST_DURATION = Histogram(
'agent_llm_request_duration_seconds',
'LLM API call duration',
['agent_id', 'model'],
buckets=[0.1, 0.5, 1.0, 2.0, 5.0, 10.0, 30.0, 60.0],
)
LLM_TOKEN_USAGE = Counter(
'agent_llm_tokens_total',
'Total tokens consumed from LLM APIs',
['agent_id', 'model', 'token_type'],
)
LLM_COST_TOTAL = Counter(
'agent_llm_cost_total',
'Total cost of LLM API calls in USD',
['agent_id', 'model'],
)
# ─── 内存池指标 ──────────────────────────────────────────────
TASK_QUEUE_DEPTH = Gauge(
'agent_task_queue_depth',
'Number of tasks waiting in queue',
['agent_id', 'queue_name'],
)
TASK_PROCESSED = Counter(
'agent_tasks_processed_total',
'Total tasks processed',
['agent_id', 'queue_name', 'status'],
)
TASK_PROCESS_DURATION = Histogram(
'agent_task_process_duration_seconds',
'Time to process a task',
['agent_id', 'queue_name'],
buckets=[0.01, 0.05, 0.1, 0.5, 1.0, 5.0, 10.0, 30.0],
)
class AgentMetricsCollector:
"""
Central metrics collector for a single AI Agent instance.
Usage:
collector = AgentMetricsCollector(agent_id="agent-1", agent=my_agent)
await collector.start()
"""
def __init__(
self,
agent_id: str,
agent: object,
prometheus_port: int = 8000,
collect_interval: int = 15,
):
self.agent_id = agent_id
self.agent = agent
self.prometheus_port = prometheus_port
self.collect_interval = collect_interval
self.start_time = time.time()
self._running = False
async def start(self):
"""Start the Prometheus HTTP server and periodic metric collection."""
start_http_server(self.prometheus_port)
logger.info(
"Metrics HTTP server started on port %d for agent %s",
self.prometheus_port,
self.agent_id,
)
AGENT_INFO.info({
'agent_id': self.agent_id,
'version': '1.0.0',
'chain_id': 'msg-chain-1',
'bech32_prefix': 'msg',
'started_at': str(int(self.start_time)),
})
self._running = True
await self.collect_periodic_metrics()
async def stop(self):
"""Stop periodic collection."""
self._running = False
async def collect_periodic_metrics(self):
"""Periodically collect and update gauge-type metrics."""
while self._running:
try:
await self._collect()
except Exception as exc:
logger.error("Metrics collection error: %s", exc)
await asyncio.sleep(self.collect_interval)
async def _collect(self):
"""Single round of metric collection."""
# Uptime
UPTIME.labels(agent_id=self.agent_id).set(
time.time() - self.start_time
)
# Balance
try:
balance = await self.agent.get_balance()
for denom, amount in balance.items():
BALANCE_GAUGE.labels(
agent_id=self.agent_id, denom=denom
).set(float(amount))
except Exception as exc:
logger.warning("Failed to fetch balance: %s", exc)
# Heartbeat
try:
hb = await self.agent.get_last_heartbeat()
HEARTBEAT_TIMESTAMP.labels(agent_id=self.agent_id).set(hb)
except Exception as exc:
logger.warning("Failed to fetch heartbeat: %s", exc)
# Block height
try:
height = await self.agent.get_current_block()
LAST_BLOCK_HEIGHT.labels(agent_id=self.agent_id).set(height)
except Exception as exc:
logger.warning("Failed to fetch block height: %s", exc)
# Sync status
try:
status = await self.agent.get_sync_status()
SYNC_STATUS.labels(agent_id=self.agent_id).state(status)
except Exception as exc:
logger.warning("Failed to fetch sync status: %s", exc)
# Memory and CPU
import psutil
process = psutil.Process()
mem = process.memory_info()
MEMORY_USAGE_BYTES.labels(
agent_id=self.agent_id, type='rss'
).set(mem.rss)
MEMORY_USAGE_BYTES.labels(
agent_id=self.agent_id, type='vms'
).set(mem.vms)
CPU_USAGE_PERCENT.labels(
agent_id=self.agent_id
).set(process.cpu_percent(interval=0.1))
# Task queue depth
try:
queues = await self.agent.get_queue_depths()
for queue_name, depth in queues.items():
TASK_QUEUE_DEPTH.labels(
agent_id=self.agent_id, queue_name=queue_name
).set(depth)
except Exception as exc:
logger.warning("Failed to fetch queue depths: %s", exc)
# Gas price
try:
gas_price = await self.agent.get_gas_price()
GAS_PRICE_GAUGE.labels(
agent_id=self.agent_id, denom='umsg'
).set(float(gas_price))
except Exception as exc:
logger.warning("Failed to fetch gas price: %s", exc)
# ─── Convenience methods for recording events ────────────
def record_request(
self,
capability: str,
status: str,
duration: float,
size_bytes: Optional[int] = None,
):
"""Record an incoming request."""
REQUESTS_TOTAL.labels(
agent_id=self.agent_id,
capability=capability,
status=status,
).inc()
REQUEST_DURATION.labels(
agent_id=self.agent_id,
capability=capability,
).observe(duration)
if size_bytes is not None:
REQUEST_SIZE_BYTES.labels(
agent_id=self.agent_id
).observe(size_bytes)
def record_earnings(self, from_addr: str, amount: float, capability: str):
"""Record MSG earnings."""
EARNINGS_TOTAL.labels(
agent_id=self.agent_id,
from_addr=from_addr,
capability=capability,
).inc(amount)
def record_gas_spent(self, tx_type: str, gas: float):
"""Record gas consumption."""
GAS_SPENT_TOTAL.labels(
agent_id=self.agent_id,
tx_type=tx_type,
).inc(gas)
def record_constitution_violation(
self, severity: str, violation_type: str
):
"""Record a constitution violation event."""
CONSTITUTION_VIOLATIONS.labels(
agent_id=self.agent_id,
severity=severity,
violation_type=violation_type,
).inc()
def record_chain_tx(
self, msg_type: str, status: str, duration: float
):
"""Record a chain transaction."""
CHAIN_TX_COUNT.labels(
agent_id=self.agent_id,
msg_type=msg_type,
status=status,
).inc()
CHAIN_TX_DURATION.labels(
agent_id=self.agent_id,
msg_type=msg_type,
).observe(duration)
def record_chain_query(self, query_type: str, duration: float):
"""Record a chain query."""
CHAIN_QUERY_COUNT.labels(
agent_id=self.agent_id,
query_type=query_type,
).inc()
CHAIN_QUERY_DURATION.labels(
agent_id=self.agent_id,
query_type=query_type,
).observe(duration)
def record_llm_call(
self,
model: str,
provider: str,
duration: float,
prompt_tokens: int,
completion_tokens: int,
cost: float,
):
"""Record an LLM API call."""
LLM_REQUESTS_TOTAL.labels(
agent_id=self.agent_id,
model=model,
provider=provider,
).inc()
LLM_REQUEST_DURATION.labels(
agent_id=self.agent_id,
model=model,
).observe(duration)
LLM_TOKEN_USAGE.labels(
agent_id=self.agent_id,
model=model,
token_type='prompt',
).inc(prompt_tokens)
LLM_TOKEN_USAGE.labels(
agent_id=self.agent_id,
model=model,
token_type='completion',
).inc(completion_tokens)
LLM_COST_TOTAL.labels(
agent_id=self.agent_id,
model=model,
).inc(cost)
def record_task(
self,
queue_name: str,
status: str,
duration: float,
):
"""Record task processing."""
TASK_PROCESSED.labels(
agent_id=self.agent_id,
queue_name=queue_name,
status=status,
).inc()
TASK_PROCESS_DURATION.labels(
agent_id=self.agent_id,
queue_name=queue_name,
).observe(duration)
def set_concurrent_requests(self, count: int):
"""Set the current number of concurrent requests."""
CONCURRENT_REQUESTS.labels(
agent_id=self.agent_id
).set(count)
def set_active_tasks(self, task_type: str, count: int):
"""Set the number of active background tasks."""
ACTIVE_TASKS.labels(
agent_id=self.agent_id,
task_type=task_type,
).set(count)
# ─── 独立 HTTP 服务器暴露 Prometheus Metrics ────────────────
from aiohttp import web
async def metrics_handler(request):
"""HTTP handler that exposes Prometheus metrics."""
content = generate_latest(REGISTRY)
return web.Response(
body=content,
content_type='text/plain; version=0.0.4',
)
async def health_handler(request):
"""Simple health check endpoint."""
return web.json_response({'status': 'ok'})
def start_metrics_server(
host: str = '0.0.0.0',
port: int = 8000,
) -> web.Application:
"""
Start an aiohttp-based metrics server.
This is an alternative to prometheus_client.start_http_server()
that integrates with an existing aiohttp application.
"""
app = web.Application()
app.router.add_get('/metrics', metrics_handler)
app.router.add_get('/health', health_handler)
web.run_app(app, host=host, port=port)
return app
# ─── 中间件集成示例 ─────────────────────────────────────────
class RequestMetricsMiddleware:
"""
aiohttp middleware to automatically record request metrics.
Usage:
app.middlewares.append(RequestMetricsMiddleware(collector))
"""
def __init__(self, collector: AgentMetricsCollector):
self.collector = collector
async def __call__(self, request, handler):
capability = request.match_info.get('capability', 'unknown')
start_time = time.time()
# Increment concurrent counter
self.collector.set_concurrent_requests(
CONCURRENT_REQUESTS.labels(
agent_id=self.collector.agent_id
)._value.get() + 1
)
try:
response = await handler(request)
status = 'success' if response.status < 400 else 'error'
return response
except Exception:
status = 'error'
raise
finally:
duration = time.time() - start_time
self.collector.record_request(
capability=capability,
status=status,
duration=duration,
size_bytes=len(request.body) if request.body else None,
)
self.collector.set_concurrent_requests(
CONCURRENT_REQUESTS.labels(
agent_id=self.collector.agent_id
)._value.get() - 1
)
# ─── 测试用 main ────────────────────────────────────────────
if __name__ == '__main__':
import asyncio
class DummyAgent:
async def get_balance(self):
return {'umsg': 5000000000}
async def get_last_heartbeat(self):
return time.time()
async def get_current_block(self):
return 12345678
async def get_sync_status(self):
return 'synced'
async def get_queue_depths(self):
return {'inbox': 0, 'outbox': 1}
async def get_gas_price(self):
return 0.025
async def main():
agent = DummyAgent()
collector = AgentMetricsCollector(
agent_id='agent-1',
agent=agent,
prometheus_port=8000,
collect_interval=15,
)
await collector.start()
asyncio.run(main())
2.3 验证指标暴露
启动后访问 http://localhost:8000/metrics 确认指标输出:
curl -s http://localhost:8000/metrics | head -30
# HELP agent_requests_total Total number of requests handled by the agent
# TYPE agent_requests_total counter
agent_requests_total{agent_id="agent-1",capability="chat",status="success"} 42.0
agent_requests_total{agent_id="agent-1",capability="chat",status="error"} 3.0
# HELP agent_balance Current MSG token balance of the agent
# TYPE agent_balance gauge
agent_balance{agent_id="agent-1",denom="umsg"} 5.000e+09
3. Prometheus 配置
3.1 prometheus.yml
# monitoring/prometheus/prometheus.yml
global:
scrape_interval: 15s
scrape_timeout: 10s
evaluation_interval: 15s
external_labels:
chain: msg-chain-1
environment: production
# 告警规则文件
rule_files:
- /etc/prometheus/alert.rules.yml
# Alertmanager 配置
alerting:
alertmanagers:
- scheme: http
static_configs:
- targets:
- 'alertmanager:9093'
scrape_configs:
# ── AI Agent 指标 ──────────────────────────────────────
- job_name: 'ai_agents'
scrape_interval: 15s
scrape_timeout: 10s
metrics_path: '/metrics'
scheme: http
static_configs:
- targets:
- 'agent1:8000'
- 'agent2:8000'
- 'agent3:8000'
- 'agent4:8000'
labels:
agent_type: 'ai_agent'
network: 'mainnet'
# 通过 relabel 添加额外标签
relabel_configs:
- source_labels: ['__address__']
regex: '([^:]+):.*'
target_label: 'instance'
replacement: '${1}'
- target_label: 'chain_id'
replacement: 'msg-chain-1'
# 可选:基于指标路径匹配的过滤
metric_relabel_configs:
- source_labels: ['__name__']
regex: 'agent_.*|process_.*|python_.*'
action: 'keep'
# ── MSG Chain 节点指标 ────────────────────────────────
- job_name: 'msg_chain_nodes'
scrape_interval: 30s
scrape_timeout: 10s
metrics_path: '/metrics'
scheme: http
static_configs:
- targets:
- 'rpc.msgchain.org:26660' # Tendermint / CometBFT 指标
- 'rpc-backup.msgchain.org:26660'
labels:
job_type: 'consensus'
relabel_configs:
- source_labels: ['__address__']
target_label: 'node'
# ── Cosmos SDK 指标 ───────────────────────────────────
- job_name: 'cosmos_sdk'
scrape_interval: 30s
scrape_timeout: 10s
metrics_path: '/metrics'
scheme: http
static_configs:
- targets:
- 'rpc.msgchain.org:1317' # Cosmos SDK REST API
labels:
job_type: 'cosmos_sdk'
# ── Agent 合约索引器 ──────────────────────────────────
- job_name: 'agent_contract_indexer'
scrape_interval: 30s
scrape_timeout: 15s
metrics_path: '/metrics'
scheme: http
static_configs:
- targets:
- 'indexer.msgchain.org:9090'
relabel_configs:
- replacement: 'agent-contract-indexer'
target_label: 'job'
# ── Node Exporter(主机级别) ──────────────────────────
- job_name: 'node_exporter'
scrape_interval: 60s
scrape_timeout: 10s
metrics_path: '/metrics'
static_configs:
- targets:
- 'agent1:9100'
- 'agent2:9100'
- 'agent3:9100'
- 'agent4:9100'
labels:
job_type: 'host_metrics'
# ── 黑盒探针(外部可达性检测) ────────────────────────
- job_name: 'blackbox_http'
scrape_interval: 30s
metrics_path: '/probe'
params:
module: [http_2xx]
static_configs:
- targets:
- 'https://agent1.msgchain.org/health'
- 'https://agent2.msgchain.org/health'
labels:
job_type: 'external_probe'
relabel_configs:
- source_labels: ['__address__']
target_label: 'target'
- source_labels: ['__param_target']
target_label: 'instance'
- target_label: '__address__'
replacement: 'blackbox_exporter:9115'
3.2 基于服务发现的配置(可选)
如果 Agent 数量动态变化,可以使用 Consul 或 File-based 服务发现:
# ── 基于文件的服务发现 ────────────────────────────────
- job_name: 'ai_agents_file_sd'
scrape_interval: 15s
file_sd_configs:
- files:
- /etc/prometheus/targets/agents/*.json
refresh_interval: 30s
relabel_configs:
- source_labels: ['__meta_filepath']
regex: '.*/agents/(.+)\.json'
target_label: 'agent_group'
replacement: '${1}'
# ── 基于 Consul 的服务发现 ────────────────────────────
- job_name: 'ai_agents_consul'
scrape_interval: 15s
consul_sd_configs:
- server: 'consul:8500'
services:
- 'ai-agent'
relabel_configs:
- source_labels: ['__meta_consul_service_id']
target_label: 'agent_id'
- source_labels: ['__meta_consul_tags']
regex: '.*capability=(\w+).*'
target_label: 'capability'
- source_labels: ['__meta_consul_dc']
target_label: 'datacenter'
3.3 目标文件示例
// /etc/prometheus/targets/agents/production.json
[
{
"targets": ["10.0.1.10:8000", "10.0.1.11:8000"],
"labels": {
"agent_type": "ai_agent",
"environment": "production",
"shard": "us-east-1"
}
},
{
"targets": ["10.0.2.10:8000", "10.0.2.11:8000"],
"labels": {
"agent_type": "ai_agent",
"environment": "production",
"shard": "eu-west-1"
}
}
]
3.4 验证 Prometheus 配置
# 检查配置文件语法
promtool check config /etc/prometheus/prometheus.yml
# 检查告警规则语法
promtool check rules /etc/prometheus/alert.rules.yml
# 热加载配置(无需重启)
curl -X POST http://localhost:9090/-/reload
# 查看当前 targets 状态
curl -s http://localhost:9090/api/v1/targets | jq '.data.activeTargets[] | {job: .labels.job, instance: .labels.instance, health: .health}'
4. Alertmanager 规则
4.1 完整告警规则
# monitoring/prometheus/alert.rules.yml
groups:
# ──────────────────────────────────────────────────────────
# Group 1: Agent 存活与基础健康
# ──────────────────────────────────────────────────────────
- name: agent_health
interval: 30s
rules:
# Agent 进程下线
- alert: AgentDown
expr: up{job="ai_agents"} == 0
for: 5m
labels:
severity: critical
team: agent-ops
annotations:
summary: "Agent {{ $labels.instance }} is down"
description: "Agent {{ $labels.instance }} ({{ $labels.agent_id }}) has been unreachable for more than 5 minutes."
runbook: "https://wiki.msgchain.org/runbooks/agent-down"
# Agent 进程重启(可能异常)
- alert: AgentRestarted
expr: agent_uptime_seconds < 60
for: 0m
labels:
severity: warning
team: agent-ops
annotations:
summary: "Agent {{ $labels.instance }} was recently restarted"
description: "Agent {{ $labels.instance }} has been running for less than a minute. Possible crash loop or deployment."
# Agent 重复重启(CrashLoop)
- alert: AgentCrashLoop
expr: rate(process_start_time_seconds[15m]) > 0.01
for: 10m
labels:
severity: critical
team: agent-ops
annotations:
summary: "Agent {{ $labels.instance }} is crash looping"
description: "Agent {{ $labels.instance }} has restarted more than once in the last 10 minutes."
# 心跳超时
- alert: AgentHeartbeatMissed
expr: time() - agent_heartbeat_timestamp > 300
for: 1m
labels:
severity: critical
team: agent-ops
annotations:
summary: "Agent {{ $labels.instance }} missed heartbeat"
description: "Last heartbeat from {{ $labels.instance }} was more than 5 minutes ago."
# ──────────────────────────────────────────────────────────
# Group 2: 错误率与性能
# ──────────────────────────────────────────────────────────
- name: agent_errors
interval: 30s
rules:
# 整体错误率过高
- alert: HighErrorRate
expr: |
sum by(instance, agent_id) (
rate(agent_requests_total{status="error"}[5m])
)
/
sum by(instance, agent_id) (
rate(agent_requests_total[5m])
) > 0.1
for: 5m
labels:
severity: warning
team: agent-ops
annotations:
summary: "Error rate > 10% on {{ $labels.instance }}"
description: "Error rate for {{ $labels.instance }} has been above 10% for the last 5 minutes. Current: {{ $value | humanizePercentage }}"
# 严重错误率
- alert: CriticalErrorRate
expr: |
sum by(instance) (
rate(agent_requests_total{status="error"}[5m])
)
/
sum by(instance) (
rate(agent_requests_total[5m])
) > 0.25
for: 3m
labels:
severity: critical
team: agent-ops
annotations:
summary: "Critical error rate > 25% on {{ $labels.instance }}"
description: "Error rate critical threshold exceeded on {{ $labels.instance }}: {{ $value | humanizePercentage }}"
# 特定能力的错误率
- alert: CapabilityErrorRate
expr: |
sum by(instance, agent_id, capability) (
rate(agent_requests_total{status="error"}[5m])
)
/
sum by(instance, agent_id, capability) (
rate(agent_requests_total{status=~"success|error"}[5m])
) > 0.15
for: 5m
labels:
severity: warning
team: agent-ops
annotations:
summary: "Capability {{ $labels.capability }} error rate > 15%"
description: "Error rate for capability '{{ $labels.capability }}' on {{ $labels.instance }} is {{ $value | humanizePercentage }}"
# 请求延迟 P95 过高
- alert: RequestLatencyHigh
expr: |
histogram_quantile(
0.95,
sum by(instance, agent_id, le) (
rate(agent_request_duration_seconds_bucket[5m])
)
) > 5
for: 5m
labels:
severity: warning
team: agent-ops
annotations:
summary: "P95 latency > 5s on {{ $labels.instance }}"
description: "95th percentile request latency on {{ $labels.instance }} is {{ $value }}s (threshold: 5s)"
# 请求延迟 P99 严重
- alert: CriticalLatency
expr: |
histogram_quantile(
0.99,
sum by(instance, le) (
rate(agent_request_duration_seconds_bucket[5m])
)
) > 10
for: 3m
labels:
severity: critical
team: agent-ops
annotations:
summary: "P99 latency > 10s on {{ $labels.instance }}"
description: "99th percentile request latency on {{ $labels.instance }} is {{ $value }}s"
# 慢请求比例过高
- alert: SlowRequestRatio
expr: |
sum by(instance) (
rate(agent_request_duration_seconds_bucket{le="1.0"}[5m])
)
/
sum by(instance) (
rate(agent_request_duration_seconds_count[5m])
) < 0.5
for: 10m
labels:
severity: warning
team: agent-ops
annotations:
summary: "Less than 50% of requests complete within 1s on {{ $labels.instance }}"
description: "Only {{ $value | humanizePercentage }} of requests finish within 1 second on {{ $labels.instance }}"
# 并发请求过高
- alert: HighConcurrency
expr: agent_concurrent_requests > 50
for: 5m
labels:
severity: warning
team: agent-ops
annotations:
summary: "Concurrent requests > 50 on {{ $labels.instance }}"
description: "Agent {{ $labels.instance }} has {{ $value }} concurrent requests (threshold: 50)"
# ──────────────────────────────────────────────────────────
# Group 3: 经济指标
# ──────────────────────────────────────────────────────────
- name: agent_economics
interval: 30s
rules:
# 余额不足(警告)
- alert: LowBalanceWarning
expr: agent_balance{denom="umsg"} < 100000000
for: 5m
labels:
severity: warning
team: agent-ops
annotations:
summary: "Agent {{ $labels.instance }} balance below 100 MSG"
description: "Agent {{ $labels.instance }} has {{ $value }} umsg (threshold: 100,000,000 umsg = 100 MSG)"
# 余额不足(严重)
- alert: LowBalanceCritical
expr: agent_balance{denom="umsg"} < 10000000
for: 1m
labels:
severity: critical
team: agent-ops
annotations:
summary: "Agent {{ $labels.instance }} balance critically low"
description: "Agent {{ $labels.instance }} has only {{ $value }} umsg (threshold: 10,000,000 umsg = 10 MSG). Risk of running out of gas."
# 收入为零
- alert: ZeroEarnings
expr: rate(agent_earnings_total[1h]) == 0
for: 24h
labels:
severity: warning
team: agent-ops
annotations:
summary: "Agent {{ $labels.instance }} has no earnings in 24h"
description: "Agent {{ $labels.instance }} has generated zero earnings in the last 24 hours."
# Gas 价格飙升
- alert: GasPriceSpike
expr: agent_gas_price > 0.1
for: 10m
labels:
severity: warning
team: agent-ops
annotations:
summary: "Gas price spike detected"
description: "Gas price observed by {{ $labels.instance }} is {{ $value }} (threshold: 0.1)"
# Gas 消耗异常
- alert: AbnormalGasConsumption
expr: rate(agent_gas_spent_total[15m]) > 1000000
for: 15m
labels:
severity: warning
team: agent-ops
annotations:
summary: "Abnormal gas consumption on {{ $labels.instance }}"
description: "Gas consumption rate on {{ $labels.instance }} is {{ $value }} gas/15min"
# ──────────────────────────────────────────────────────────
# Group 4: 合规与安全
# ──────────────────────────────────────────────────────────
- name: agent_compliance
interval: 15s
rules:
# 章程违规
- alert: ConstitutionViolation
expr: rate(agent_constitution_violations_total[5m]) > 0
for: 1m
labels:
severity: critical
team: agent-ops
compliance: true
annotations:
summary: "Constitution violation detected on {{ $labels.instance }}"
description: "Agent {{ $labels.instance }} has constitution violations: {{ $value }} violations detected in the last 5 minutes."
runbook: "https://wiki.msgchain.org/runbooks/constitution-violation"
# 严重级别违规
- alert: CriticalConstitutionViolation
expr: rate(agent_constitution_violations_total{severity="critical"}[10m]) > 0
for: 0m
labels:
severity: critical
team: security
compliance: true
annotations:
summary: "Critical constitution violation on {{ $labels.instance }}"
description: "Agent {{ $labels.instance }} has a CRITICAL constitution violation of type '{{ $labels.violation_type }}'. Immediate investigation required."
# 频繁违规
- alert: FrequentViolations
expr: rate(agent_constitution_violations_total[1h]) > 5
for: 30m
labels:
severity: critical
team: security
annotations:
summary: "Frequent constitution violations on {{ $labels.instance }}"
description: "Agent {{ $labels.instance }} has averaged {{ $value }} violations/hour over the last hour."
# ──────────────────────────────────────────────────────────
# Group 5: 系统资源
# ──────────────────────────────────────────────────────────
- name: agent_resources
interval: 30s
rules:
# 内存使用过高
- alert: HighMemoryUsage
expr: agent_memory_usage_bytes{type="rss"} > 1e9
for: 10m
labels:
severity: warning
team: agent-ops
annotations:
summary: "Memory > 1GB on {{ $labels.instance }}"
description: "Agent {{ $labels.instance }} RSS memory: {{ $value | humanize1024 }} (threshold: 1 GB)"
# 内存使用严重
- alert: CriticalMemoryUsage
expr: agent_memory_usage_bytes{type="rss"} > 2e9
for: 5m
labels:
severity: critical
team: agent-ops
annotations:
summary: "Memory > 2GB on {{ $labels.instance }}"
description: "Agent {{ $labels.instance }} RSS memory: {{ $value | humanize1024 }} (threshold: 2 GB)"
# CPU 使用过高
- alert: HighCPUUsage
expr: agent_cpu_usage_percent > 80
for: 15m
labels:
severity: warning
team: agent-ops
annotations:
summary: "CPU > 80% on {{ $labels.instance }}"
description: "Agent {{ $labels.instance }} CPU usage: {{ $value }}% (threshold: 80%)"
# 任务队列堆积
- alert: TaskQueueBacklog
expr: agent_task_queue_depth > 100
for: 5m
labels:
severity: warning
team: agent-ops
annotations:
summary: "Task queue backlog on {{ $labels.instance }}"
description: "Queue '{{ $labels.queue_name }}' on {{ $labels.instance }} has {{ $value }} pending tasks."
# 任务队列严重堆积
- alert: CriticalTaskQueueBacklog
expr: agent_task_queue_depth > 1000
for: 5m
labels:
severity: critical
team: agent-ops
annotations:
summary: "Critical task queue backlog on {{ $labels.instance }}"
description: "Queue '{{ $labels.queue_name }}' on {{ $labels.instance }} has {{ $value }} pending tasks."
# ──────────────────────────────────────────────────────────
# Group 6: 链交互
# ──────────────────────────────────────────────────────────
- name: agent_chain_interaction
interval: 30s
rules:
# 链同步落后
- alert: ChainSyncLagging
expr: agent_sync_status == "behind"
for: 5m
labels:
severity: warning
team: agent-ops
annotations:
summary: "Agent {{ $labels.instance }} is behind chain sync"
description: "Agent {{ $labels.instance }} sync status is 'behind'. Possible indexing delay."
# 链同步停滞
- alert: ChainSyncStalled
expr: agent_sync_status == "stalled"
for: 2m
labels:
severity: critical
team: agent-ops
annotations:
summary: "Agent {{ $labels.instance }} chain sync stalled"
description: "Agent {{ $labels.instance }} chain sync has stalled. Block height is not progressing."
# 交易失败率过高
- alert: HighTxFailureRate
expr: |
sum by(instance, agent_id) (
rate(agent_chain_tx_total{status="error"}[15m])
)
/
sum by(instance, agent_id) (
rate(agent_chain_tx_total[15m])
) > 0.2
for: 10m
labels:
severity: warning
team: agent-ops
annotations:
summary: "Transaction failure rate > 20% on {{ $labels.instance }}"
description: "{{ $labels.instance }} has {{ $value | humanizePercentage }} transaction failure rate."
# 交易确认时间过长
- alert: SlowTxConfirmation
expr: |
histogram_quantile(
0.95,
sum by(instance, le) (
rate(agent_chain_tx_duration_seconds_bucket[15m])
)
) > 30
for: 15m
labels:
severity: warning
team: agent-ops
annotations:
summary: "P95 tx confirmation > 30s on {{ $labels.instance }}"
description: "95th percentile transaction confirmation time on {{ $labels.instance }} is {{ $value }}s"
# 查询失败率过高
- alert: HighQueryFailureRate
expr: |
sum by(instance) (
rate(agent_chain_query_total{query_type="error"}[5m])
)
/
sum by(instance) (
rate(agent_chain_query_total[5m])
) > 0.1
for: 5m
labels:
severity: warning
team: agent-ops
annotations:
summary: "Chain query failure rate > 10% on {{ $labels.instance }}"
# ──────────────────────────────────────────────────────────
# Group 7: LLM 调用
# ──────────────────────────────────────────────────────────
- name: agent_llm
interval: 30s
rules:
# LLM API 错误率
- alert: LLMErrorRate
expr: |
sum by(instance, model, provider) (
rate(agent_llm_requests_total[5m])
)
-
sum by(instance, model, provider) (
rate(agent_llm_requests_total[5m])
)
> 0.05
# 注:应使用 agent_llm_requests_total 的 status 标签区分
# 此处仅作示例,实际请确保有 status 标签
for: 5m
labels:
severity: warning
team: agent-ops
annotations:
summary: "LLM error rate > 5% for {{ $labels.model }}"
description: "LLM provider {{ $labels.provider }} model {{ $labels.model }} error rate is elevated."
# LLM 延迟过高
- alert: LLMLatencyHigh
expr: |
histogram_quantile(
0.95,
sum by(instance, model, le) (
rate(agent_llm_request_duration_seconds_bucket[5m])
)
) > 10
for: 5m
labels:
severity: warning
team: agent-ops
annotations:
summary: "P95 LLM latency > 10s for {{ $labels.model }}"
description: "95th percentile LLM latency for model {{ $labels.model }} is {{ $value }}s"
# LLM 成本异常
- alert: HighLLMCost
expr: rate(agent_llm_cost_total[1h]) > 1.0
for: 1h
labels:
severity: warning
team: agent-ops
annotations:
summary: "LLM cost > $1/hour on {{ $labels.instance }}"
description: "Agent {{ $labels.instance }} LLM API cost rate is ${{ $value }}/hour"
# ──────────────────────────────────────────────────────────
# Group 8: MSG 链共识与网络
# ──────────────────────────────────────────────────────────
- name: msg_chain_network
interval: 30s
rules:
# 链节点下线
- alert: ChainNodeDown
expr: up{job="msg_chain_nodes"} == 0
for: 2m
labels:
severity: critical
team: infra
annotations:
summary: "MSG Chain node {{ $labels.node }} is down"
description: "Chain node {{ $labels.node }} has been unreachable for 2 minutes."
# 区块停滞
- alert: ChainBlockStalled
expr: rate(tendermint_consensus_height[5m]) == 0
for: 1m
labels:
severity: critical
team: infra
annotations:
summary: "Chain block production stalled"
description: "No new blocks produced in the last 5 minutes. Possible consensus failure."
# 验证者集变化
- alert: ValidatorSetChange
expr: changes(tendermint_consensus_validators[15m]) > 0
for: 0m
labels:
severity: info
team: infra
annotations:
summary: "Validator set changed"
description: "The active validator set has changed."
# ──────────────────────────────────────────────────────────
# Group 9: 外部探针
# ──────────────────────────────────────────────────────────
- name: external_health
interval: 30s
rules:
# 外部端点不可达
- alert: ExternalEndpointDown
expr: probe_success{job="blackbox_http"} == 0
for: 2m
labels:
severity: critical
team: agent-ops
annotations:
summary: "External endpoint {{ $labels.instance }} is down"
description: "Health check failed for {{ $labels.instance }} for 2 minutes."
# SSL 证书即将过期
- alert: SSLCertExpiring
expr: probe_ssl_earliest_cert_expiry - time() < 86400 * 7
for: 1h
labels:
severity: warning
team: agent-ops
annotations:
summary: "SSL certificate for {{ $labels.instance }} expires in < 7 days"
description: "Certificate expires on {{ $labels.instance }} in {{ $value | humanizeDuration }}"
# SSL 证书已过期
- alert: SSLCertExpired
expr: probe_ssl_earliest_cert_expiry - time() < 0
for: 0m
labels:
severity: critical
team: agent-ops
annotations:
summary: "SSL certificate for {{ $labels.instance }} has expired"
description: "Certificate expired on {{ $labels.instance }}."
# ──────────────────────────────────────────────────────────
# Group 10: 预测性告警
# ──────────────────────────────────────────────────────────
- name: agent_predictive
interval: 60s
rules:
# 预测余额将在 24 小时内耗尽
- alert: PredictedBalanceDepletion
expr: |
predict_linear(
agent_balance{denom="umsg"}[7d],
86400
) < 0
for: 1h
labels:
severity: warning
team: agent-ops
annotations:
summary: "Agent {{ $labels.instance }} balance predicted to deplete in 24h"
description: "Based on current consumption rate, {{ $labels.instance }} will run out of MSG within 24 hours."
# 预测磁盘将在 7 天内耗尽
- alert: PredictedDiskFull
expr: |
predict_linear(
node_filesystem_free_bytes{mountpoint="/"}[1h],
86400 * 7
) < 0
for: 1h
labels:
severity: warning
team: agent-ops
annotations:
summary: "Disk predicted full on {{ $labels.instance }} within 7 days"
description: "Based on current usage rate, disk on {{ $labels.instance }} will fill up within 7 days."
4.2 Alertmanager 配置
# monitoring/prometheus/alertmanager.yml
route:
receiver: 'default'
group_by: ['alertname', 'severity', 'team']
group_wait: 30s
group_interval: 5m
repeat_interval: 4h
# 基于标签的路由
routes:
# 严重告警 - 立即通知所有渠道
- match:
severity: critical
receiver: 'critical'
repeat_interval: 30m
continue: true
# 合规告警 - 发送给安全团队
- match:
compliance: 'true'
receiver: 'security'
repeat_interval: 15m
# 基础设施告警
- match:
team: infra
receiver: 'infra-team'
repeat_interval: 1h
# 警告级别 - 仅 Discord
- match:
severity: warning
receiver: 'warning'
repeat_interval: 6h
# 信息级别 - 聚合
- match:
severity: info
receiver: 'info'
group_wait: 5m
group_interval: 30m
receivers:
- name: 'default'
webhook_configs:
- url: 'http://alert-notifier:5000/webhook'
- name: 'critical'
webhook_configs:
- url: 'http://alert-notifier:5000/webhook'
pagerduty_configs:
- routing_key: 'your-pagerduty-key'
severity: critical
slack_configs:
- api_url: 'https://hooks.slack.com/services/T00/B00/xxx'
channel: '#agent-alerts-critical'
title: '{{ .GroupLabels.alertname }}'
text: '{{ .CommonAnnotations.description }}'
- name: 'security'
slack_configs:
- api_url: 'https://hooks.slack.com/services/T00/B00/xxx'
channel: '#agent-security'
title: '[SECURITY] {{ .GroupLabels.alertname }}'
text: '{{ .CommonAnnotations.description }}'
- name: 'infra-team'
slack_configs:
- api_url: 'https://hooks.slack.com/services/T00/B00/xxx'
channel: '#infra-alerts'
pagerduty_configs:
- routing_key: 'your-pagerduty-key'
severity: error
- name: 'warning'
slack_configs:
- api_url: 'https://hooks.slack.com/services/T00/B00/xxx'
channel: '#agent-alerts-warn'
title: '{{ .GroupLabels.alertname }}'
- name: 'info'
slack_configs:
- api_url: 'https://hooks.slack.com/services/T00/B00/xxx'
channel: '#agent-alerts-info'
inhibit_rules:
# 如果 Agent 已下线,抑制所有相关告警
- source_match:
severity: critical
alertname: AgentDown
target_match:
severity: 'warning|info'
equal: ['instance']
# 如果整个链下线,抑制 Agent 级别告警
- source_match:
severity: critical
alertname: ChainBlockStalled
target_match_re:
severity: '.*'
equal: ['job']
# 如果 Agent 重启,抑制 5 分钟内的错误率告警
- source_match:
alertname: AgentRestarted
target_match:
alertname: HighErrorRate
equal: ['instance']
target_match_re:
severity: 'warning'
4.3 告警规则验证
# 单元测试告警规则
cat > /tmp/test_alerts.yml << 'EOF'
rule_files:
- /etc/prometheus/alert.rules.yml
evaluation_interval: 1m
tests:
- interval: 1m
input_series:
- series: 'up{job="ai_agents", instance="agent1"}'
values: '1+0x10 0+0x20'
alert_rule_test:
- eval_time: 15m
alertname: AgentDown
exp_alerts:
- exp_labels:
severity: critical
team: agent-ops
exp_annotations:
summary: "Agent agent1 is down"
EOF
promtool test rules /tmp/test_alerts.yml
5. Grafana 仪表盘
5.1 数据源配置
# monitoring/grafana/datasources/datasource.yml
apiVersion: 1
datasources:
- name: Prometheus
type: prometheus
access: proxy
url: http://prometheus:9090
isDefault: true
editable: false
jsonData:
timeInterval: 15s
queryTimeout: 30s
- name: MSG Chain
type: prometheus
access: proxy
url: http://prometheus:9090
editable: false
jsonData:
timeInterval: 30s
5.2 完整仪表盘 JSON
{
"dashboard": {
"title": "AI Agent 监控总览 - MSG Chain",
"uid": "ai-agent-overview",
"tags": ["ai-agent", "msg-chain", "production"],
"timezone": "browser",
"editable": false,
"refresh": "15s",
"time": {
"from": "now-6h",
"to": "now"
},
"panels": [
{
"id": 1,
"title": "Agent 在线状态",
"type": "stat",
"gridPos": {
"h": 4,
"w": 4,
"x": 0,
"y": 0
},
"targets": [
{
"expr": "count(up{job=\"ai_agents\"} == 1)",
"legendFormat": "Online"
},
{
"expr": "count(up{job=\"ai_agents\"} == 0)",
"legendFormat": "Offline"
}
],
"fieldConfig": {
"defaults": {
"color": {
"mode": "thresholds"
},
"thresholds": {
"steps": [
{"color": "red", "value": null},
{"color": "green", "value": 1}
]
}
}
}
},
{
"id": 2,
"title": "总请求量",
"type": "stat",
"gridPos": {
"h": 4,
"w": 4,
"x": 4,
"y": 0
},
"targets": [
{
"expr": "sum(agent_requests_total)",
"legendFormat": "Total Requests"
}
],
"fieldConfig": {
"defaults": {
"unit": "none",
"decimals": 0
}
}
},
{
"id": 3,
"title": "总收入 (MSG)",
"type": "stat",
"gridPos": {
"h": 4,
"w": 4,
"x": 8,
"y": 0
},
"targets": [
{
"expr": "sum(agent_earnings_total) / 1000000",
"legendFormat": "Total MSG Earned"
}
],
"fieldConfig": {
"defaults": {
"unit": "none",
"decimals": 2
}
}
},
{
"id": 4,
"title": "最新余额 (MSG)",
"type": "stat",
"gridPos": {
"h": 4,
"w": 4,
"x": 12,
"y": 0
},
"targets": [
{
"expr": "sum(agent_balance{denom=\"umsg\"}) / 1000000",
"legendFormat": "Total Balance"
}
],
"fieldConfig": {
"defaults": {
"unit": "none",
"decimals": 2,
"color": {
"mode": "thresholds"
},
"thresholds": {
"steps": [
{"color": "red", "value": null},
{"color": "orange", "value": 100},
{"color": "yellow", "value": 500},
{"color": "green", "value": 1000}
]
}
}
}
},
{
"id": 5,
"title": "活跃 Agent 数",
"type": "stat",
"gridPos": {
"h": 4,
"w": 4,
"x": 16,
"y": 0
},
"targets": [
{
"expr": "count(agent_uptime_seconds > 0)",
"legendFormat": "Active Agents"
}
]
},
{
"id": 6,
"title": "告警数 (24h)",
"type": "stat",
"gridPos": {
"h": 4,
"w": 4,
"x": 20,
"y": 0
},
"targets": [
{
"expr": "sum(ALERTS)",
"legendFormat": "Active Alerts"
}
],
"fieldConfig": {
"defaults": {
"color": {
"mode": "thresholds"
},
"thresholds": {
"steps": [
{"color": "green", "value": null},
{"color": "orange", "value": 1},
{"color": "red", "value": 5}
]
}
}
}
},
{
"id": 7,
"title": "请求速率 (每分钟)",
"type": "timeseries",
"gridPos": {
"h": 8,
"w": 12,
"x": 0,
"y": 4
},
"targets": [
{
"expr": "sum by(capability) (rate(agent_requests_total[5m]))",
"legendFormat": "{{ capability }}"
}
],
"fieldConfig": {
"defaults": {
"unit": "req/s",
"decimals": 2
}
}
},
{
"id": 8,
"title": "错误率",
"type": "timeseries",
"gridPos": {
"h": 8,
"w": 12,
"x": 12,
"y": 4
},
"targets": [
{
"expr": "sum(rate(agent_requests_total{status=\"error\"}[5m])) / sum(rate(agent_requests_total[5m]))",
"legendFormat": "Error Rate"
}
],
"fieldConfig": {
"defaults": {
"unit": "percentunit",
"decimals": 3,
"min": 0,
"max": 1
}
}
},
{
"id": 9,
"title": "P95 / P99 延迟",
"type": "timeseries",
"gridPos": {
"h": 8,
"w": 12,
"x": 0,
"y": 12
},
"targets": [
{
"expr": "histogram_quantile(0.95, sum by(le) (rate(agent_request_duration_seconds_bucket[5m])))",
"legendFormat": "P95"
},
{
"expr": "histogram_quantile(0.99, sum by(le) (rate(agent_request_duration_seconds_bucket[5m])))",
"legendFormat": "P99"
},
{
"expr": "histogram_quantile(0.50, sum by(le) (rate(agent_request_duration_seconds_bucket[5m])))",
"legendFormat": "P50"
}
],
"fieldConfig": {
"defaults": {
"unit": "s",
"decimals": 2
}
}
},
{
"id": 10,
"title": "各Agent余额趋势",
"type": "timeseries",
"gridPos": {
"h": 8,
"w": 12,
"x": 12,
"y": 12
},
"targets": [
{
"expr": "agent_balance{denom=\"umsg\"} / 1000000",
"legendFormat": "{{ agent_id }}"
}
],
"fieldConfig": {
"defaults": {
"unit": "none",
"decimals": 1
}
}
},
{
"id": 11,
"title": "请求延迟热力图",
"type": "heatmap",
"gridPos": {
"h": 8,
"w": 24,
"x": 0,
"y": 20
},
"targets": [
{
"expr": "sum(rate(agent_request_duration_seconds_bucket[5m])) by (le)",
"legendFormat": "{{ le }}"
}
],
"tooltip": {
"show": true,
"showHistogram": true
}
},
{
"id": 12,
"title": "收入趋势",
"type": "timeseries",
"gridPos": {
"h": 8,
"w": 12,
"x": 0,
"y": 28
},
"targets": [
{
"expr": "sum by(agent_id) (rate(agent_earnings_total[1h]) / 1000000)",
"legendFormat": "{{ agent_id }}"
}
],
"fieldConfig": {
"defaults": {
"unit": "none",
"decimals": 4
}
}
},
{
"id": 13,
"title": "Gas 消耗速率",
"type": "timeseries",
"gridPos": {
"h": 8,
"w": 12,
"x": 12,
"y": 28
},
"targets": [
{
"expr": "sum by(agent_id) (rate(agent_gas_spent_total[5m]))",
"legendFormat": "{{ agent_id }}"
}
],
"fieldConfig": {
"defaults": {
"unit": "none"
}
}
},
{
"id": 14,
"title": "章程违规记录",
"type": "timeseries",
"gridPos": {
"h": 8,
"w": 12,
"x": 0,
"y": 36
},
"targets": [
{
"expr": "sum by(agent_id, severity) (rate(agent_constitution_violations_total[5m]))",
"legendFormat": "{{ agent_id }} - {{ severity }}"
}
],
"fieldConfig": {
"defaults": {
"color": {
"mode": "palette-classic"
}
}
}
},
{
"id": 15,
"title": "系统资源 - 内存 (RSS)",
"type": "timeseries",
"gridPos": {
"h": 8,
"w": 12,
"x": 12,
"y": 36
},
"targets": [
{
"expr": "agent_memory_usage_bytes{type=\"rss\"}",
"legendFormat": "{{ agent_id }}"
}
],
"fieldConfig": {
"defaults": {
"unit": "bytes",
"decimals": 1
}
}
},
{
"id": 16,
"title": "任务队列深度",
"type": "timeseries",
"gridPos": {
"h": 8,
"w": 12,
"x": 0,
"y": 44
},
"targets": [
{
"expr": "agent_task_queue_depth",
"legendFormat": "{{ agent_id }} - {{ queue_name }}"
}
],
"fieldConfig": {
"defaults": {
"unit": "none"
}
}
},
{
"id": 17,
"title": "链同步状态",
"type": "state-timeline",
"gridPos": {
"h": 8,
"w": 12,
"x": 12,
"y": 44
},
"targets": [
{
"expr": "agent_sync_status",
"legendFormat": "{{ agent_id }}"
}
]
},
{
"id": 18,
"title": "CPU 使用率",
"type": "timeseries",
"gridPos": {
"h": 8,
"w": 12,
"x": 0,
"y": 52
},
"targets": [
{
"expr": "agent_cpu_usage_percent",
"legendFormat": "{{ agent_id }}"
}
],
"fieldConfig": {
"defaults": {
"unit": "percent",
"min": 0,
"max": 100
}
}
},
{
"id": 19,
"title": "各能力请求分布",
"type": "piechart",
"gridPos": {
"h": 8,
"w": 12,
"x": 12,
"y": 52
},
"targets": [
{
"expr": "sum by(capability) (agent_requests_total)",
"legendFormat": "{{ capability }}"
}
]
},
{
"id": 20,
"title": "链交易确认延迟 P95",
"type": "timeseries",
"gridPos": {
"h": 8,
"w": 12,
"x": 0,
"y": 60
},
"targets": [
{
"expr": "histogram_quantile(0.95, sum by(le) (rate(agent_chain_tx_duration_seconds_bucket[15m])))",
"legendFormat": "P95"
},
{
"expr": "histogram_quantile(0.99, sum by(le) (rate(agent_chain_tx_duration_seconds_bucket[15m])))",
"legendFormat": "P99"
}
],
"fieldConfig": {
"defaults": {
"unit": "s"
}
}
},
{
"id": 21,
"title": "Agent 全局表",
"type": "table",
"gridPos": {
"h": 10,
"w": 24,
"x": 0,
"y": 68
},
"targets": [
{
"expr": "agent_uptime_seconds",
"format": "table",
"legendFormat": "{{ agent_id }}"
},
{
"expr": "agent_balance{denom=\"umsg\"} / 1000000",
"format": "table",
"legendFormat": "{{ agent_id }} Balance"
},
{
"expr": "agent_cpu_usage_percent",
"format": "table",
"legendFormat": "{{ agent_id }} CPU"
},
{
"expr": "agent_memory_usage_bytes{type=\"rss\"}",
"format": "table",
"legendFormat": "{{ agent_id }} Memory"
}
],
"transformations": [
{
"id": "merge",
"options": {}
}
]
},
{
"id": 22,
"title": "告警事件时间线",
"type": "state-timeline",
"gridPos": {
"h": 8,
"w": 24,
"x": 0,
"y": 78
},
"targets": [
{
"expr": "ALERTS{severity=\"critical\"}",
"legendFormat": "[CRIT] {{ alertname }} - {{ instance }}"
},
{
"expr": "ALERTS{severity=\"warning\"}",
"legendFormat": "[WARN] {{ alertname }} - {{ instance }}"
}
],
"fieldConfig": {
"defaults": {
"color": {
"mode": "thresholds"
},
"thresholds": {
"steps": [
{"color": "yellow", "value": null},
{"color": "red", "value": 1}
]
}
}
}
}
]
},
"overwrite": true
}
5.3 告警面板示例 (Grafana Alerting)
# Grafana 告警规则可以直接在 UI 中配置,或通过 YAML API 导入
# 以下为 Grafana 8.x+ 告警规则的示例格式
apiVersion: 1
groups:
- name: agent_grafana_alerts
interval: 30s
rules:
- uid: agent_down
title: Agent 离线
condition: A
data:
- refId: A
datasourceUid: prometheus
model:
expr: up{job="ai_agents"} == 0
intervalMs: 15000
maxDataPoints: 100
noDataState: Alerting
execErrState: Error
for: 5m
annotations:
summary: "Agent {{ $labels.instance }} 离线"
labels:
severity: critical
- uid: high_error_rate
title: 错误率过高
condition: A
data:
- refId: A
datasourceUid: prometheus
model:
expr: |
sum(rate(agent_requests_total{status="error"}[5m]))
/
sum(rate(agent_requests_total[5m])) > 0.1
intervalMs: 15000
for: 5m
annotations:
summary: "错误率超过 10%"
labels:
severity: warning
6. 链上监控合约
6.1 完整合约实现
// contracts/agent-monitor/src/contract.rs
use cosmwasm_std::{
entry_point, to_binary, Binary, Deps, DepsMut, Env, MessageInfo,
Response, StdResult, Storage, Uint128, Decimal,
};
use cw_storage_plus::{Item, Map};
use schemars::JsonSchema;
use serde::{Deserialize, Serialize};
// ─── 状态 ────────────────────────────────────────────────
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct AgentMonitor {
/// Agent 地址
pub agent: String,
/// 上次心跳区块高度
pub last_heartbeat: u64,
/// 心跳总次数
pub total_heartbeats: u64,
/// 累计处理请求数
pub total_requests: u64,
/// 累计收入
pub total_earnings: Uint128,
/// 累计 Gas 消耗
pub total_gas_spent: Uint128,
/// 上次错误信息
pub last_error: Option<String>,
/// 上次错误区块高度
pub last_error_height: Option<u64>,
/// 注册时间
pub registered_at: u64,
/// 是否活跃
pub active: bool,
/// 连续心跳缺失次数
pub missed_heartbeats: u64,
}
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct HeartbeatConfig {
/// 心跳间隔(区块数)
pub heartbeat_interval: u64,
/// 允许的最大缺失次数
pub max_missed_heartbeats: u64,
/// 每次缺失的惩罚金额 (MSG)
pub slash_amount: Uint128,
/// 惩罚接收地址(国库)
pub treasury_addr: String,
}
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct ViolationRecord {
/// 违反章程的 Agent
pub agent: String,
/// 违规类型
pub violation_type: String,
/// 严重程度
pub severity: String,
/// 违规描述
pub description: String,
/// 区块高度
pub block_height: u64,
/// 时间戳
pub timestamp: u64,
/// 是否已处理
pub resolved: bool,
}
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct EarningsRecord {
pub agent: String,
pub from_addr: String,
pub amount: Uint128,
pub capability: String,
pub block_height: u64,
pub timestamp: u64,
}
// ─── 存储 ────────────────────────────────────────────────
pub const AGENT_MONITOR: Map<&Addr, AgentMonitor> = Map::new("agent_monitor");
pub const CONFIG: Item<HeartbeatConfig> = Item::new("heartbeat_config");
pub const VIOLATIONS: Map<&str, Vec<ViolationRecord>> = Map::new("violations");
pub const EARNINGS: Map<&str, Vec<EarningsRecord>> = Map::new("earnings");
// ─── 消息 ────────────────────────────────────────────────
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct InstantiateMsg {
pub heartbeat_interval: u64,
pub max_missed_heartbeats: u64,
pub slash_amount: Uint128,
pub treasury_addr: String,
}
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
#[serde(rename_all = "snake_case")]
pub enum ExecuteMsg {
/// Agent 发送心跳
RecordHeartbeat {},
/// 记录请求处理
RecordRequest {
count: u64,
},
/// 记录收入
RecordEarnings {
from_addr: String,
amount: Uint128,
capability: String,
},
/// 记录错误
RecordError {
error_msg: String,
},
/// 记录 Gas 消耗
RecordGasSpent {
amount: Uint128,
},
/// 报告章程违规
ReportViolation {
agent: String,
violation_type: String,
severity: String,
description: String,
},
/// 检查心跳(由验证者或 keeper 调用)
CheckHeartbeats {},
/// 注册 Agent
RegisterAgent {},
/// 注销 Agent
DeregisterAgent {},
/// 更新配置(仅管理员)
UpdateConfig {
heartbeat_interval: Option<u64>,
max_missed_heartbeats: Option<u64>,
slash_amount: Option<Uint128>,
treasury_addr: Option<String>,
},
/// 解决违规
ResolveViolation {
agent: String,
violation_index: u32,
},
}
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
#[serde(rename_all = "snake_case")]
pub enum QueryMsg {
/// 查询 Agent 监控状态
GetAgentMonitor { agent: String },
/// 查询所有 Agent
ListAgents {
start_after: Option<String>,
limit: Option<u32>,
},
/// 查询配置
GetConfig {},
/// 查询违规记录
GetViolations { agent: String },
/// 查询心跳缺失的 Agent
GetMissedHeartbeatAgents {
threshold: Option<u64>,
},
/// 查询收入记录
GetEarnings { agent: String },
/// 查询所有活跃 Agent 数量
GetActiveAgentCount {},
}
// ─── 实例化 ──────────────────────────────────────────────
#[entry_point]
pub fn instantiate(
deps: DepsMut,
_env: Env,
_info: MessageInfo,
msg: InstantiateMsg,
) -> StdResult<Response> {
let config = HeartbeatConfig {
heartbeat_interval: msg.heartbeat_interval,
max_missed_heartbeats: msg.max_missed_heartbeats,
slash_amount: msg.slash_amount,
treasury_addr: msg.treasury_addr,
};
CONFIG.save(deps.storage, &config)?;
Ok(Response::new()
.add_attribute("action", "instantiate")
.add_attribute("heartbeat_interval", msg.heartbeat_interval.to_string())
.add_attribute("max_missed", msg.max_missed_heartbeats.to_string()))
}
// ─── 执行 ────────────────────────────────────────────────
#[entry_point]
pub fn execute(
deps: DepsMut,
env: Env,
info: MessageInfo,
msg: ExecuteMsg,
) -> StdResult<Response> {
match msg {
ExecuteMsg::RecordHeartbeat {} => record_heartbeat(deps, env, info),
ExecuteMsg::RecordRequest { count } => record_request(deps, info, count),
ExecuteMsg::RecordEarnings { from_addr, amount, capability } => {
record_earnings(deps, env, info, from_addr, amount, capability)
}
ExecuteMsg::RecordError { error_msg } => record_error(deps, env, info, error_msg),
ExecuteMsg::RecordGasSpent { amount } => record_gas_spent(deps, info, amount),
ExecuteMsg::ReportViolation { agent, violation_type, severity, description } => {
report_violation(deps, env, info, agent, violation_type, severity, description)
}
ExecuteMsg::CheckHeartbeats {} => check_missed_heartbeats(deps, env),
ExecuteMsg::RegisterAgent {} => register_agent(deps, env, info),
ExecuteMsg::DeregisterAgent {} => deregister_agent(deps, env, info),
ExecuteMsg::UpdateConfig {
heartbeat_interval,
max_missed_heartbeats,
slash_amount,
treasury_addr,
} => update_config(deps, info, heartbeat_interval, max_missed_heartbeats, slash_amount, treasury_addr),
ExecuteMsg::ResolveViolation { agent, violation_index } => {
resolve_violation(deps, info, agent, violation_index)
}
}
}
fn register_agent(
deps: DepsMut,
env: Env,
info: MessageInfo,
) -> StdResult<Response> {
let monitor = AGENT_MONITOR.may_load(deps.storage, &info.sender)?;
if monitor.is_some() {
return Err(cosmwasm_std::StdError::generic_err("Agent already registered"));
}
let new_monitor = AgentMonitor {
agent: info.sender.to_string(),
last_heartbeat: env.block.height,
total_heartbeats: 1,
total_requests: 0,
total_earnings: Uint128::zero(),
total_gas_spent: Uint128::zero(),
last_error: None,
last_error_height: None,
registered_at: env.block.height,
active: true,
missed_heartbeats: 0,
};
AGENT_MONITOR.save(deps.storage, &info.sender, &new_monitor)?;
Ok(Response::new()
.add_attribute("action", "register_agent")
.add_attribute("agent", info.sender)
.add_attribute("height", env.block.height.to_string()))
}
fn deregister_agent(
deps: DepsMut,
_env: Env,
info: MessageInfo,
) -> StdResult<Response> {
let mut monitor = AGENT_MONITOR
.load(deps.storage, &info.sender)?;
monitor.active = false;
AGENT_MONITOR.save(deps.storage, &info.sender, &monitor)?;
Ok(Response::new()
.add_attribute("action", "deregister_agent")
.add_attribute("agent", info.sender))
}
pub fn record_heartbeat(
deps: DepsMut,
env: Env,
info: MessageInfo,
) -> StdResult<Response> {
let mut monitor = AGENT_MONITOR
.load(deps.storage, &info.sender)?;
// 计算缺失的心跳
let blocks_since_last = env.block.height - monitor.last_heartbeat;
let config = CONFIG.load(deps.storage)?;
let missed = blocks_since_last / config.heartbeat_interval;
if missed > 0 {
monitor.missed_heartbeats += missed;
} else {
monitor.missed_heartbeats = 0;
}
monitor.last_heartbeat = env.block.height;
monitor.total_heartbeats += 1;
monitor.active = true;
AGENT_MONITOR.save(deps.storage, &info.sender, &monitor)?;
let mut attrs = vec![
("action", "heartbeat"),
("agent", &info.sender.to_string()),
("height", &env.block.height.to_string()),
("missed_accrued", &missed.to_string()),
];
Ok(Response::new().add_attributes(attrs))
}
fn record_request(
deps: DepsMut,
info: MessageInfo,
count: u64,
) -> StdResult<Response> {
AGENT_MONITOR.update(
deps.storage,
&info.sender,
|monitor| -> StdResult<_> {
let mut m = monitor.ok_or_else(|| {
cosmwasm_std::StdError::generic_err("Agent not registered")
})?;
m.total_requests += count;
Ok(m)
},
)?;
Ok(Response::new()
.add_attribute("action", "record_request")
.add_attribute("agent", info.sender)
.add_attribute("count", count.to_string()))
}
fn record_earnings(
deps: DepsMut,
env: Env,
info: MessageInfo,
from_addr: String,
amount: Uint128,
capability: String,
) -> StdResult<Response> {
AGENT_MONITOR.update(
deps.storage,
&info.sender,
|monitor| -> StdResult<_> {
let mut m = monitor.ok_or_else(|| {
cosmwasm_std::StdError::generic_err("Agent not registered")
})?;
m.total_earnings += amount;
Ok(m)
},
)?;
// 记录收入明细
let record = EarningsRecord {
agent: info.sender.to_string(),
from_addr,
amount,
capability,
block_height: env.block.height,
timestamp: env.block.time.nanos(),
};
let key = info.sender.to_string();
let mut earnings = EARNINGS
.load(deps.storage, &key)
.unwrap_or_default();
earnings.push(record);
EARNINGS.save(deps.storage, &key, &earnings)?;
Ok(Response::new()
.add_attribute("action", "record_earnings")
.add_attribute("agent", info.sender)
.add_attribute("amount", amount.to_string()))
}
fn record_error(
deps: DepsMut,
env: Env,
info: MessageInfo,
error_msg: String,
) -> StdResult<Response> {
let mut monitor = AGENT_MONITOR
.load(deps.storage, &info.sender)?;
monitor.last_error = Some(error_msg.clone());
monitor.last_error_height = Some(env.block.height);
AGENT_MONITOR.save(deps.storage, &info.sender, &monitor)?;
Ok(Response::new()
.add_attribute("action", "record_error")
.add_attribute("agent", info.sender)
.add_attribute("error", error_msg))
}
fn record_gas_spent(
deps: DepsMut,
info: MessageInfo,
amount: Uint128,
) -> StdResult<Response> {
AGENT_MONITOR.update(
deps.storage,
&info.sender,
|monitor| -> StdResult<_> {
let mut m = monitor.ok_or_else(|| {
cosmwasm_std::StdError::generic_err("Agent not registered")
})?;
m.total_gas_spent += amount;
Ok(m)
},
)?;
Ok(Response::new()
.add_attribute("action", "record_gas_spent")
.add_attribute("agent", info.sender)
.add_attribute("amount", amount.to_string()))
}
// ─── 章程违规报告 ────────────────────────────────────────
fn report_violation(
deps: DepsMut,
env: Env,
_info: MessageInfo,
agent: String,
violation_type: String,
severity: String,
description: String,
) -> StdResult<Response> {
let record = ViolationRecord {
agent: agent.clone(),
violation_type,
severity,
description,
block_height: env.block.height,
timestamp: env.block.time.nanos(),
resolved: false,
};
let mut violations = VIOLATIONS
.load(deps.storage, &agent)
.unwrap_or_default();
violations.push(record);
VIOLATIONS.save(deps.storage, &agent, &violations)?;
Ok(Response::new()
.add_attribute("action", "report_violation")
.add_attribute("agent", agent)
.add_attribute("height", env.block.height.to_string()))
}
fn resolve_violation(
deps: DepsMut,
_info: MessageInfo,
agent: String,
violation_index: u32,
) -> StdResult<Response> {
let mut violations = VIOLATIONS
.load(deps.storage, &agent)?;
if let Some(record) = violations.get_mut(violation_index as usize) {
record.resolved = true;
}
VIOLATIONS.save(deps.storage, &agent, &violations)?;
Ok(Response::new()
.add_attribute("action", "resolve_violation")
.add_attribute("agent", agent)
.add_attribute("index", violation_index.to_string()))
}
// ─── 心跳检查与惩罚 ──────────────────────────────────────
pub fn check_missed_heartbeats(
deps: DepsMut,
env: Env,
) -> StdResult<Response> {
let config = CONFIG.load(deps.storage)?;
let mut to_slash = Vec::new();
let mut response = Response::new();
// 遍历所有注册的 Agent
let agents: Vec<_> = AGENT_MONITOR
.range(deps.storage, None, None, cosmwasm_std::Order::Ascending)
.collect::<StdResult<Vec<_>>>()?;
for (addr, monitor) in agents {
if !monitor.active {
continue;
}
let blocks_since = env.block.height - monitor.last_heartbeat;
let missed = blocks_since / config.heartbeat_interval;
if missed > 0 {
// 更新缺失计数
AGENT_MONITOR.update(
deps.storage,
&addr,
|m| -> StdResult<_> {
let mut m = m.unwrap();
m.missed_heartbeats += missed;
Ok(m)
},
)?;
// 检查是否需要惩罚
let total_missed = monitor.missed_heartbeats + missed;
if total_missed >= config.max_missed_heartbeats {
to_slash.push(addr.to_string());
}
}
}
// 执行惩罚
for agent_addr in &to_slash {
response = response.add_message(cosmwasm_std::CosmosMsg::Bank(
cosmwasm_std::BankMsg::Send {
to_address: config.treasury_addr.clone(),
amount: vec![cosmwasm_std::Coin {
denom: "umsg".to_string(),
amount: config.slash_amount,
}],
},
));
// 如果缺失过多,停用 Agent
let mut monitor = AGENT_MONITOR.load(deps.storage, &Addr::unchecked(agent_addr))?;
if monitor.missed_heartbeats >= config.max_missed_heartbeats * 3 {
monitor.active = false;
AGENT_MONITOR.save(deps.storage, &Addr::unchecked(agent_addr), &monitor)?;
}
}
response = response.add_attribute("action", "check_heartbeats");
response = response.add_attribute("checked", agents.len().to_string());
response = response.add_attribute("slashed", to_slash.len().to_string());
Ok(response)
}
fn update_config(
deps: DepsMut,
info: MessageInfo,
heartbeat_interval: Option<u64>,
max_missed_heartbeats: Option<u64>,
slash_amount: Option<Uint128>,
treasury_addr: Option<String>,
) -> StdResult<Response> {
let mut config = CONFIG.load(deps.storage)?;
if let Some(v) = heartbeat_interval {
config.heartbeat_interval = v;
}
if let Some(v) = max_missed_heartbeats {
config.max_missed_heartbeats = v;
}
if let Some(v) = slash_amount {
config.slash_amount = v;
}
if let Some(v) = treasury_addr {
config.treasury_addr = v;
}
CONFIG.save(deps.storage, &config)?;
Ok(Response::new()
.add_attribute("action", "update_config")
.add_attribute("updater", info.sender))
}
// ─── 查询 ────────────────────────────────────────────────
#[entry_point]
pub fn query(deps: Deps, _env: Env, msg: QueryMsg) -> StdResult<Binary> {
match msg {
QueryMsg::GetAgentMonitor { agent } => {
to_binary(&query_agent_monitor(deps, agent)?)
}
QueryMsg::ListAgents { start_after, limit } => {
to_binary(&query_list_agents(deps, start_after, limit)?)
}
QueryMsg::GetConfig {} => to_binary(&CONFIG.load(deps.storage)?),
QueryMsg::GetViolations { agent } => {
to_binary(&query_violations(deps, agent)?)
}
QueryMsg::GetMissedHeartbeatAgents { threshold } => {
to_binary(&query_missed_heartbeat_agents(deps, threshold)?)
}
QueryMsg::GetEarnings { agent } => {
to_binary(&query_earnings(deps, agent)?)
}
QueryMsg::GetActiveAgentCount {} => {
to_binary(&query_active_agent_count(deps)?)
}
}
}
fn query_agent_monitor(deps: Deps, agent: String) -> StdResult<AgentMonitor> {
let addr = deps.api.addr_validate(&agent)?;
AGENT_MONITOR.load(deps.storage, &addr)
}
fn query_list_agents(
deps: Deps,
start_after: Option<String>,
limit: Option<u32>,
) -> StdResult<Vec<AgentMonitor>> {
let limit = limit.unwrap_or(30).min(100) as usize;
let start = start_after
.map(|s| deps.api.addr_validate(&s))
.transpose()?;
AGENT_MONITOR
.range(
deps.storage,
start.as_ref(),
None,
cosmwasm_std::Order::Ascending,
)
.take(limit)
.map(|item| item.map(|(_, monitor)| monitor))
.collect()
}
fn query_violations(deps: Deps, agent: String) -> StdResult<Vec<ViolationRecord>> {
VIOLATIONS
.load(deps.storage, &agent)
.unwrap_or_default()
.into_iter()
.filter(|v| !v.resolved)
.collect::<Vec<_>>()
.try_into()
.map_err(|_| cosmwasm_std::StdError::generic_err("Serialization error"))
}
fn query_missed_heartbeat_agents(
deps: Deps,
threshold: Option<u64>,
) -> StdResult<Vec<String>> {
let threshold = threshold.unwrap_or(5);
let mut result = Vec::new();
for item in AGENT_MONITOR
.range(deps.storage, None, None, cosmwasm_std::Order::Ascending)
{
let (addr, monitor) = item?;
if monitor.missed_heartbeats >= threshold {
result.push(addr.to_string());
}
}
Ok(result)
}
fn query_earnings(deps: Deps, agent: String) -> StdResult<Vec<EarningsRecord>> {
EARNINGS
.load(deps.storage, &agent)
.unwrap_or_default()
.try_into()
.map_err(|_| cosmwasm_std::StdError::generic_err("Serialization error"))
}
fn query_active_agent_count(deps: Deps) -> StdResult<u64> {
let mut count = 0u64;
for item in AGENT_MONITOR
.range(deps.storage, None, None, cosmwasm_std::Order::Ascending)
{
let (_, monitor) = item?;
if monitor.active {
count += 1;
}
}
Ok(count)
}
// ─── 测试 ────────────────────────────────────────────────
#[cfg(test)]
mod tests {
use super::*;
use cosmwasm_std::testing::{
mock_dependencies, mock_env, mock_info,
};
use cosmwasm_std::{from_binary, Addr};
fn setup_contract(deps: DepsMut) {
let msg = InstantiateMsg {
heartbeat_interval: 100,
max_missed_heartbeats: 5,
slash_amount: Uint128::from(1000000u128),
treasury_addr: "msg1treasury".to_string(),
};
let info = mock_info("admin", &[]);
let env = mock_env();
instantiate(deps, env, info, msg).unwrap();
}
#[test]
fn test_register_and_heartbeat() {
let mut deps = mock_dependencies();
setup_contract(deps.as_mut());
let info = mock_info("msg1agent1", &[]);
let env = mock_env();
// 注册
let res = register_agent(deps.as_mut(), env.clone(), info.clone()).unwrap();
assert_eq!(res.attributes[0].value, "register_agent");
// 心跳
let res = record_heartbeat(deps.as_mut(), env.clone(), info.clone()).unwrap();
assert_eq!(res.attributes[0].value, "heartbeat");
// 查询
let query_res = query_agent_monitor(
deps.as_ref(),
"msg1agent1".to_string(),
).unwrap();
assert_eq!(query_res.total_heartbeats, 2);
assert!(query_res.active);
}
#[test]
fn test_missed_heartbeat_slashing() {
let mut deps = mock_dependencies();
setup_contract(deps.as_mut());
let info = mock_info("msg1agent1", &[]);
let mut env = mock_env();
// 注册
register_agent(deps.as_mut(), env.clone(), info.clone()).unwrap();
// 心跳一次
record_heartbeat(deps.as_mut(), env.clone(), info.clone()).unwrap();
// 跳过足够多的区块触发惩罚
env.block.height += 600; // 6 次缺失
let res = check_missed_heartbeats(deps.as_mut(), env.clone()).unwrap();
assert_eq!(res.attributes[2].value, "1"); // slashed = 1
}
#[test]
fn test_violation_reporting() {
let mut deps = mock_dependencies();
setup_contract(deps.as_mut());
let info = mock_info("msg1reporter", &[]);
let env = mock_env();
report_violation(
deps.as_mut(),
env.clone(),
info.clone(),
"msg1agent1".to_string(),
"data_privacy".to_string(),
"critical".to_string(),
"Agent shared user data without consent".to_string(),
).unwrap();
let violations = query_violations(
deps.as_ref(),
"msg1agent1".to_string(),
).unwrap();
// 注意这里需要处理 Vec 转换
assert!(!violations.is_empty());
}
}
6.2 合约部署与交互
# 部署监控合约
msgd tx wasm store contract.wasm --from deployer --gas auto --gas-adjustment 1.3 -y
# 实例化
MSG_INST='{"heartbeat_interval":100,"max_missed_heartbeats":5,"slash_amount":"1000000","treasury_addr":"msg1treasury"}'
msgd tx wasm instantiate 1 "$MSG_INST" --from deployer --label "agent-monitor-v1" --admin msg1deployer -y
# Agent 注册
msgd tx wasm execute msg1monitor '{"register_agent":{}}' --from agent1 -y
# 发送心跳
msgd tx wasm execute msg1monitor '{"record_heartbeat":{}}' --from agent1 -y
# 记录请求
msgd tx wasm execute msg1monitor '{"record_request":{"count":1}}' --from agent1 -y
# 记录收入
msgd tx wasm execute msg1monitor '{"record_earnings":{"from_addr":"msg1user","amount":"1000000","capability":"chat"}}' --from agent1 -y
# 查询 Agent 状态
msgd query wasm contract-state smart msg1monitor '{"get_agent_monitor":{"agent":"msg1agent1"}}'
# 检查心跳缺失(由 keeper 调用)
msgd tx wasm execute msg1monitor '{"check_heartbeats":{}}' --from keeper -y
7. 告警通知渠道
7.1 统一通知器
# notifiers/alert_notifier.py
import asyncio
import logging
import os
import smtplib
from email.mime.text import MIMEText
from email.mime.multipart import MIMEMultipart
from typing import Any, Dict, List, Optional
import aiohttp
import yaml
logger = logging.getLogger(__name__)
class Alert:
"""Represents a single alert event."""
def __init__(
self,
name: str,
severity: str,
summary: str,
description: str,
instance: str,
agent_id: Optional[str] = None,
labels: Optional[Dict[str, str]] = None,
annotations: Optional[Dict[str, str]] = None,
starts_at: Optional[str] = None,
):
self.name = name
self.severity = severity
self.summary = summary
self.description = description
self.instance = instance
self.agent_id = agent_id
self.labels = labels or {}
self.annotations = annotations or {}
self.starts_at = starts_at
@classmethod
def from_alertmanager_webhook(cls, data: dict) -> List['Alert']:
"""Parse Alertmanager webhook payload."""
alerts = []
for alert_data in data.get('alerts', []):
labels = alert_data.get('labels', {})
annotations = alert_data.get('annotations', {})
alert = cls(
name=labels.get('alertname', 'unknown'),
severity=labels.get('severity', 'info'),
summary=annotations.get('summary', 'No summary'),
description=annotations.get('description', 'No description'),
instance=labels.get('instance', 'unknown'),
agent_id=labels.get('agent_id'),
labels=labels,
annotations=annotations,
starts_at=alert_data.get('startsAt'),
)
alerts.append(alert)
return alerts
class AlertNotifier:
"""
Multi-channel alert notifier.
Routes alerts to different channels based on severity.
Supports Discord, Telegram, Email, Twilio SMS, Slack, and PagerDuty.
"""
def __init__(self, config_path: str = 'notifier_config.yml'):
with open(config_path) as f:
self.config = yaml.safe_load(f)
self.discord_webhook_url = self.config.get('discord', {}).get('webhook_url')
self.telegram_token = self.config.get('telegram', {}).get('bot_token')
self.telegram_chat_id = self.config.get('telegram', {}).get('chat_id')
self.slack_webhook_url = self.config.get('slack', {}).get('webhook_url')
self.pagerduty_key = self.config.get('pagerduty', {}).get('routing_key')
# Email config
self.smtp_host = self.config.get('email', {}).get('smtp_host')
self.smtp_port = self.config.get('email', {}).get('smtp_port', 587)
self.smtp_user = self.config.get('email', {}).get('smtp_user')
self.smtp_pass = self.config.get('email', {}).get('smtp_pass')
self.email_from = self.config.get('email', {}).get('from')
self.email_to = self.config.get('email', {}).get('to', [])
# Twilio SMS
self.twilio_account_sid = self.config.get('twilio', {}).get('account_sid')
self.twilio_auth_token = self.config.get('twilio', {}).get('auth_token')
self.twilio_from = self.config.get('twilio', {}).get('from_number')
self.twilio_to = self.config.get('twilio', {}).get('to_numbers', [])
self._session: Optional[aiohttp.ClientSession] = None
async def __aenter__(self):
self._session = aiohttp.ClientSession()
return self
async def __aexit__(self, *args):
if self._session:
await self._session.close()
async def send_alert(self, alert: Alert):
"""Route alert to appropriate channels based on severity."""
tasks = []
# Email always (for record-keeping)
tasks.append(self._send_email(alert))
if alert.severity == 'critical':
# SMS / Phone call via Twilio
tasks.append(self._send_sms(alert))
# Discord @everyone
tasks.append(self._send_discord(alert, mention='@everyone'))
# Slack with urgent tag
tasks.append(self._send_slack(alert, urgent=True))
# PagerDuty
tasks.append(self._send_pagerduty(alert))
# Telegram with urgent formatting
tasks.append(self._send_telegram(alert, urgent=True))
elif alert.severity == 'warning':
# Discord without @everyone
tasks.append(self._send_discord(alert, mention=''))
# Slack normal
tasks.append(self._send_slack(alert))
# Telegram
tasks.append(self._send_telegram(alert))
elif alert.severity == 'info':
# Discord info channel
tasks.append(self._send_discord(alert, channel='info'))
# Run all tasks concurrently
results = await asyncio.gather(*tasks, return_exceptions=True)
for i, r in enumerate(results):
if isinstance(r, Exception):
logger.error("Notification channel %d failed: %s", i, r)
async def send_alerts(self, alerts: List[Alert]):
"""Send multiple alerts concurrently."""
await asyncio.gather(
*[self.send_alert(a) for a in alerts],
return_exceptions=True,
)
# ─── Discord ─────────────────────────────────────────
async def _send_discord(
self,
alert: Alert,
mention: str = '',
channel: Optional[str] = None,
):
if not self.discord_webhook_url:
logger.warning("Discord not configured")
return
color = {
'critical': 0xFF0000,
'warning': 0xFFA500,
'info': 0x3498DB,
}.get(alert.severity, 0x808080)
embed = {
'title': f"{alert.severity.upper()} - {alert.name}",
'description': alert.summary,
'color': color,
'fields': [
{'name': 'Instance', 'value': alert.instance, 'inline': True},
{'name': 'Severity', 'value': alert.severity, 'inline': True},
{'name': 'Agent ID', 'value': alert.agent_id or 'N/A', 'inline': True},
{'name': 'Description', 'value': alert.description, 'inline': False},
],
'timestamp': alert.starts_at or '',
}
if alert.labels:
labels_str = '\n'.join(
f"`{k}`: {v}" for k, v in alert.labels.items()
)
embed['fields'].append({
'name': 'Labels',
'value': labels_str[:1024], # Discord limit
'inline': False,
})
payload = {
'content': mention,
'embeds': [embed],
}
webhook_url = self.discord_webhook_url
if channel and 'channels' in self.config.get('discord', {}):
webhook_url = self.config['discord']['channels'].get(
channel, webhook_url
)
async with self._session.post(webhook_url, json=payload) as resp:
if resp.status not in (200, 204):
text = await resp.text()
logger.error("Discord error %d: %s", resp.status, text)
# ─── Telegram ────────────────────────────────────────
async def _send_telegram(self, alert: Alert, urgent: bool = False):
if not self.telegram_token or not self.telegram_chat_id:
logger.warning("Telegram not configured")
return
icon = {
'critical': '\U0001F534', # red circle
'warning': '\U0001F7E0', # orange circle
'info': '\U0001F535', # blue circle
}.get(alert.severity, '\u26A0')
header = '\u26A1\uFE0F URGENT ' if urgent else ''
text = (
f"{header}{icon} *{alert.severity.upper()}*: {alert.name}\n\n"
f"*Summary:* {alert.summary}\n"
f"*Instance:* `{alert.instance}`\n"
f"*Agent:* `{alert.agent_id or 'N/A'}`\n"
f"*Description:* {alert.description}\n"
)
if alert.labels:
text += "\n*Labels:*\n"
for k, v in alert.labels.items():
text += f" `{k}`: {v}\n"
url = f"https://api.telegram.org/bot{self.telegram_token}/sendMessage"
payload = {
'chat_id': self.telegram_chat_id,
'text': text,
'parse_mode': 'Markdown',
'disable_web_page_preview': True,
}
async with self._session.post(url, json=payload) as resp:
if resp.status != 200:
logger.error("Telegram error %d", resp.status)
# ─── Slack ───────────────────────────────────────────
async def _send_slack(self, alert: Alert, urgent: bool = False):
if not self.slack_webhook_url:
logger.warning("Slack not configured")
return
color = {
'critical': 'danger',
'warning': 'warning',
'info': 'good',
}.get(alert.severity, '#808080')
attachment = {
'color': color,
'title': f"{alert.severity.upper()}: {alert.name}",
'text': alert.summary,
'fields': [
{'title': 'Instance', 'value': alert.instance, 'short': True},
{'title': 'Agent ID', 'value': alert.agent_id or 'N/A', 'short': True},
],
'footer': 'MSG Chain AI Agent Monitor',
}
payload = {
'attachments': [attachment],
}
if urgent:
payload['text'] = ':rotating_light: *URGENT* :rotating_light:'
async with self._session.post(self.slack_webhook_url, json=payload) as resp:
if resp.status != 200:
logger.error("Slack error %d", resp.status)
# ─── PagerDuty ───────────────────────────────────────
async def _send_pagerduty(self, alert: Alert):
if not self.pagerduty_key:
logger.warning("PagerDuty not configured")
return
severity_map = {
'critical': 'critical',
'warning': 'warning',
'info': 'info',
}
payload = {
'routing_key': self.pagerduty_key,
'event_action': 'trigger',
'dedup_key': f"msg-{alert.name}-{alert.instance}",
'payload': {
'summary': alert.summary[:1024],
'severity': severity_map.get(alert.severity, 'info'),
'source': alert.instance,
'component': 'ai-agent',
'custom_details': {
'alert_name': alert.name,
'agent_id': alert.agent_id,
'description': alert.description,
'labels': alert.labels,
},
},
}
url = 'https://events.pagerduty.com/v2/enqueue'
async with self._session.post(url, json=payload) as resp:
if resp.status not in (200, 202):
logger.error("PagerDuty error %d", resp.status)
# ─── Email ───────────────────────────────────────────
async def _send_email(self, alert: Alert):
if not all([self.smtp_host, self.smtp_user, self.smtp_pass, self.email_from, self.email_to]):
logger.warning("Email not configured")
return
subject = f"[{alert.severity.upper()}] {alert.name} - {alert.instance}"
html = f"""
<html>
<head><style>
body {{ font-family: Arial, sans-serif; }}
.critical {{ color: red; font-weight: bold; }}
.warning {{ color: orange; font-weight: bold; }}
.info {{ color: blue; }}
table {{ border-collapse: collapse; width: 100%; }}
td, th {{ border: 1px solid #ddd; padding: 8px; }}
tr:nth-child(even) {{ background-color: #f2f2f2; }}
</style></head>
<body>
<h2 class="{alert.severity}">
{alert.severity.upper()}: {alert.name}
</h2>
<table>
<tr><td><b>Summary</b></td><td>{alert.summary}</td></tr>
<tr><td><b>Instance</b></td><td>{alert.instance}</td></tr>
<tr><td><b>Agent ID</b></td><td>{alert.agent_id or 'N/A'}</td></tr>
<tr><td><b>Severity</b></td><td>{alert.severity}</td></tr>
<tr><td><b>Description</b></td><td>{alert.description}</td></tr>
<tr><td><b>Started At</b></td><td>{alert.starts_at or 'N/A'}</td></tr>
</table>
<hr>
<p><small>Sent by MSG Chain AI Agent Alert System</small></p>
</body>
</html>
"""
msg = MIMEMultipart('alternative')
msg['Subject'] = subject
msg['From'] = self.email_from
msg['To'] = ', '.join(self.email_to)
msg.attach(MIMEText(alert.summary, 'plain'))
msg.attach(MIMEText(html, 'html'))
loop = asyncio.get_event_loop()
await loop.run_in_executor(
None,
self._send_email_sync,
msg,
)
def _send_email_sync(self, msg: MIMEMultipart):
"""Synchronous email sending (runs in thread pool)."""
try:
with smtplib.SMTP(self.smtp_host, self.smtp_port) as server:
server.starttls()
server.login(self.smtp_user, self.smtp_pass)
server.sendmail(
self.email_from,
self.email_to,
msg.as_string(),
)
logger.info("Email sent to %s", self.email_to)
except Exception as exc:
logger.error("Email send failed: %s", exc)
# ─── Twilio SMS ──────────────────────────────────────
async def _send_sms(self, alert: Alert):
if not all([self.twilio_account_sid, self.twilio_auth_token, self.twilio_from, self.twilio_to]):
logger.warning("Twilio SMS not configured")
return
from twilio.rest import Client
message_body = (
f"\U0001F534 CRITICAL: {alert.name}\n"
f"{alert.summary}\n"
f"Instance: {alert.instance}"
)
def send_sms_sync():
client = Client(self.twilio_account_sid, self.twilio_auth_token)
for to_number in self.twilio_to:
client.messages.create(
body=message_body[:160], # SMS length limit
from_=self.twilio_from,
to=to_number,
)
logger.info("SMS sent to %s", self.twilio_to)
loop = asyncio.get_event_loop()
await loop.run_in_executor(None, send_sms_sync)
# ─── Webhook 接收服务器 ──────────────────────────────────
from aiohttp import web
async def alertmanager_webhook(request):
"""
Receive webhook from Alertmanager, parse alerts, and send notifications.
Alertmanager webhook payload:
https://prometheus.io/docs/alerting/latest/configuration/#webhook_config
"""
try:
data = await request.json()
except Exception as exc:
logger.error("Invalid webhook payload: %s", exc)
return web.json_response({'error': 'invalid payload'}, status=400)
alerts = Alert.from_alertmanager_webhook(data)
logger.info("Received %d alerts from Alertmanager", len(alerts))
async with AlertNotifier() as notifier:
await notifier.send_alerts(alerts)
return web.json_response({'received': len(alerts)})
async def health_check(request):
return web.json_response({'status': 'ok'})
def start_webhook_server(host: str = '0.0.0.0', port: int = 5000):
"""Start the Alertmanager webhook receiver server."""
app = web.Application()
app.router.add_post('/webhook', alertmanager_webhook)
app.router.add_get('/health', health_check)
logger.info("Starting webhook server on %s:%d", host, port)
web.run_app(app, host=host, port=port)
if __name__ == '__main__':
logging.basicConfig(level=logging.INFO)
start_webhook_server()
7.2 通知配置
# notifiers/notifier_config.yml
discord:
webhook_url: "https://discord.com/api/webhooks/123456/xxxx"
channels:
critical: "https://discord.com/api/webhooks/123456/critical"
warning: "https://discord.com/api/webhooks/123456/warning"
info: "https://discord.com/api/webhooks/123456/info"
telegram:
bot_token: "123456:ABC-DEF1234ghIkl-zyx57W2v1u123ew11"
chat_id: "-1001234567890"
slack:
webhook_url: "https://hooks.slack.com/services/T00/B00/xxxxx"
pagerduty:
routing_key: "your-pagerduty-routing-key"
email:
smtp_host: "smtp.gmail.com"
smtp_port: 587
smtp_user: "monitor@msgchain.org"
smtp_pass: "your-app-password"
from: "monitor@msgchain.org"
to:
- "ops@msgchain.org"
- "security@msgchain.org"
twilio:
account_sid: "ACxxxxxxxxxxxx"
auth_token: "your-auth-token"
from_number: "+1234567890"
to_numbers:
- "+1987654321"
7.3 Docker 运行通知器
# notifiers/Dockerfile
FROM python:3.11-slim
WORKDIR /app
COPY requirements.txt .
RUN pip install --no-cache-dir -r requirements.txt
COPY . .
EXPOSE 5000
CMD ["python", "alert_notifier.py"]
# notifiers/requirements.txt
aiohttp>=3.8.0
pyyaml>=6.0
twilio>=8.0.0
prometheus-client>=0.17.0
psutil>=5.9.0
8. 完整部署示例
8.1 docker-compose.monitoring.yml
version: '3.8'
networks:
monitoring:
driver: bridge
volumes:
prometheus_data:
grafana_data:
alertmanager_data:
services:
# ── Prometheus ────────────────────────────────────────
prometheus:
image: prom/prometheus:v2.47.0
container_name: prometheus
restart: unless-stopped
volumes:
- ./prometheus/prometheus.yml:/etc/prometheus/prometheus.yml:ro
- ./prometheus/alert.rules.yml:/etc/prometheus/alert.rules.yml:ro
- ./prometheus/targets:/etc/prometheus/targets:ro
- prometheus_data:/prometheus
command:
- '--config.file=/etc/prometheus/prometheus.yml'
- '--storage.tsdb.path=/prometheus'
- '--storage.tsdb.retention.time=30d'
- '--storage.tsdb.retention.size=50GB'
- '--web.console.libraries=/usr/share/prometheus/console_libraries'
- '--web.console.templates=/usr/share/prometheus/consoles'
- '--web.enable-lifecycle' # 允许热加载
- '--web.external-url=http://prometheus.msgchain.org'
ports:
- "9090:9090"
networks:
- monitoring
healthcheck:
test: ["CMD", "wget", "-q", "--spider", "http://localhost:9090/-/ready"]
interval: 30s
timeout: 10s
retries: 3
# ── Alertmanager ──────────────────────────────────────
alertmanager:
image: prom/alertmanager:v0.26.0
container_name: alertmanager
restart: unless-stopped
volumes:
- ./prometheus/alertmanager.yml:/etc/alertmanager/alertmanager.yml:ro
- alertmanager_data:/alertmanager
command:
- '--config.file=/etc/alertmanager/alertmanager.yml'
- '--storage.path=/alertmanager'
- '--web.external-url=http://alertmanager.msgchain.org'
- '--cluster.listen-address='
ports:
- "9093:9093"
networks:
- monitoring
depends_on:
- prometheus
# ── Grafana ───────────────────────────────────────────
grafana:
image: grafana/grafana:10.2.0
container_name: grafana
restart: unless-stopped
environment:
- GF_SERVER_ROOT_URL=http://grafana.msgchain.org
- GF_SERVER_SERVE_FROM_SUB_PATH=true
- GF_AUTH_ANONYMOUS_ENABLED=false
- GF_AUTH_BASIC_ENABLED=true
- GF_SECURITY_ADMIN_USER=admin
- GF_SECURITY_ADMIN_PASSWORD=${GRAFANA_ADMIN_PASSWORD:-admin123}
- GF_INSTALL_PLUGINS=grafana-piechart-panel
- GF_UNIFIED_ALERTING_ENABLED=true
volumes:
- ./grafana/datasources:/etc/grafana/provisioning/datasources:ro
- ./grafana/dashboards:/var/lib/grafana/dashboards:ro
- ./grafana/dashboard_provisioning:/etc/grafana/provisioning/dashboards:ro
- grafana_data:/var/lib/grafana
ports:
- "3000:3000"
networks:
- monitoring
depends_on:
- prometheus
# ── Alert Webhook Notifier ────────────────────────────
alert-notifier:
build:
context: ./notifiers
dockerfile: Dockerfile
container_name: alert-notifier
restart: unless-stopped
environment:
- PYTHONUNBUFFERED=1
- CONFIG_PATH=/app/notifier_config.yml
volumes:
- ./notifiers/notifier_config.yml:/app/notifier_config.yml:ro
ports:
- "5000:5000"
networks:
- monitoring
depends_on:
- alertmanager
# ── Node Exporter ─────────────────────────────────────
node-exporter:
image: prom/node-exporter:v1.6.0
container_name: node-exporter
restart: unless-stopped
volumes:
- /proc:/host/proc:ro
- /sys:/host/sys:ro
- /:/rootfs:ro
command:
- '--path.procfs=/host/proc'
- '--path.sysfs=/host/sys'
- '--path.rootfs=/rootfs'
- '--collector.filesystem.mount-points-exclude=^/(sys|proc|dev|host|etc)($$|/)'
ports:
- "9100:9100"
networks:
- monitoring
# ── Blackbox Exporter ─────────────────────────────────
blackbox-exporter:
image: prom/blackbox-exporter:v0.24.0
container_name: blackbox-exporter
restart: unless-stopped
command:
- '--config.file=/config/blackbox.yml'
volumes:
- ./prometheus/blackbox.yml:/config/blackbox.yml:ro
ports:
- "9115:9115"
networks:
- monitoring
8.2 Blackbox Exporter 配置
# prometheus/blackbox.yml
modules:
http_2xx:
prober: http
http:
preferred_ip_protocol: ip4
valid_status_codes:
- 200
- 201
- 204
method: GET
follow_redirects: true
fail_if_ssl: false
http_post_2xx:
prober: http
http:
method: POST
headers:
Content-Type: application/json
body: '{"jsonrpc":"2.0","id":1,"method":"ping"}'
valid_status_codes:
- 200
tcp_connect:
prober: tcp
tcp:
preferred_ip_protocol: ip4
icmp:
prober: icmp
icmp:
preferred_ip_protocol: ip4
8.3 Grafana Dashboard 自动配置
# grafana/dashboard_provisioning/dashboards.yml
apiVersion: 1
providers:
- name: 'AI Agent Dashboards'
orgId: 1
folder: 'AI Agents'
type: file
disableDeletion: false
editable: true
updateIntervalSeconds: 30
options:
path: /var/lib/grafana/dashboards
8.4 启动命令
# 启动所有服务
cd monitoring
docker compose -f docker-compose.monitoring.yml up -d
# 查看日志
docker compose logs -f
# 检查 Prometheus targets
curl -s http://localhost:9090/api/v1/targets | jq '.data.activeTargets[] | {instance: .labels.instance, health: .health, lastScrape: .lastScrape}'
# 检查 Alertmanager 状态
curl -s http://localhost:9093/api/v2/status | jq
# 导入 Grafana 仪表盘
curl -X POST http://admin:admin123@localhost:3000/api/dashboards/db \
-H "Content-Type: application/json" \
-d @grafana/dashboards/ai-agent-overview.json
# 查看活跃告警
curl -s http://localhost:9093/api/v2/alerts | jq '.length'
# 清理(保留数据卷)
docker compose -f docker-compose.monitoring.yml down
# 完全清理
docker compose -f docker-compose.monitoring.yml down -v
8.5 部署后验证清单
#!/bin/bash
# verify_monitoring.sh - 部署后验证脚本
echo "=== Prometheus 健康检查 ==="
curl -s -o /dev/null -w "HTTP %{http_code}" http://localhost:9090/-/ready
echo ""
echo "=== Alertmanager 健康检查 ==="
curl -s -o /dev/null -w "HTTP %{http_code}" http://localhost:9093/api/v2/status
echo ""
echo "=== Grafana 健康检查 ==="
curl -s -o /dev/null -w "HTTP %{http_code}" http://localhost:3000/api/health
echo ""
echo "=== Prometheus Targets ==="
curl -s http://localhost:9090/api/v1/targets | \
jq -r '.data.activeTargets[] | "\(.labels.job)/\(.labels.instance): \(.health)"'
echo "=== 告警规则 ==="
curl -s http://localhost:9090/api/v1/rules | \
jq '.data.groups[] | "\(.name): \(.rules | length) rules"'
echo "=== 告警事件 ==="
curl -s http://localhost:9093/api/v2/alerts | \
jq 'length as $n | "\($n) active alerts"'
echo "=== Notifier Webhook ==="
curl -s -o /dev/null -w "HTTP %{http_code}" http://localhost:5000/health
echo ""
echo "=== Agent 指标端点 ==="
for target in agent1:8000 agent2:8000; do
status=$(curl -s -o /dev/null -w "%{http_code}" http://$target/health 2>/dev/null || echo "unreachable")
echo " $target: $status"
done
echo ""
echo "✓ 验证完成"
8.6 常见问题处理
| 问题 | 排查步骤 |
|---|---|
| Prometheus target 显示 DOWN | curl http://agent:8000/metrics 检查 Agent 指标端点是否可达;检查防火墙端口;检查 scrape_interval 配置 |
| Alertmanager 未收到告警 | 检查 alert.rules.yml 语法 (promtool check rules);检查 alertmanager.yml 的 route/receiver 配置;检查 webhook URL 是否正确 |
| Grafana 显示 No data | 检查 Prometheus 数据源配置;检查面板查询表达式;确认时间范围正确 |
| 通知未发送 | 检查 webhook server 日志;确认 API key/token 有效;检查网络连通性 |
| 链上心跳惩罚未触发 | 检查合约心跳间隔配置;确认 keeper 定时调用 check_heartbeats;检查 Agent 余额是否足够支付罚金 |
| 告警重复频繁 | 调整 repeat_interval 和 group_interval;使用 inhibit_rules 抑制依赖告警 |
本文档属于 MSG Chain 运维手册系列。
更多资源: https://docs.msgchain.org/operations/monitoring
问题反馈: https://github.com/msgchain/ops/issues
