dApp Docs/AI Agent 监控告警规则配置指南
Development reference. Not independently verified for production.

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 可能在余额耗尽后静默失效,或在违反章程后继续运行造成损失。本指南提供从指标暴露到告警通知的完整方案。

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