dApp Docs/AI Agent 数据可视化与报告指南
Development reference. Not independently verified for production.

MSG Chain AI Agent 数据可视化与报告指南

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

适用于 MSG Chain 主网 (msg) 的 AI Agent 数据可视化、实时仪表盘与自动化报告系统。
链标识: msg-chain-1 | bech32 前缀: msg


目录

  1. 概述
  2. 链上数据获取
  3. Matplotlib 可视化
  4. Plotly 交互式图表
  5. Streamlit 实时仪表盘
  6. 自动化报告生成
  7. 报告分发
  8. 完整示例:监控与报告 Agent
  9. 附录:常见问题

1. 概述

1.1 为什么 Agent 需要数据可视化

AI Agent 在 MSG Chain 上运行时产生大量链上数据:请求日志、收益记录、声誉评分变化、执行错误等。原始数据难以洞察趋势,可视化将数字转化为直观图形,帮助开发者、持有者和利益相关者快速理解 Agent 状态。

一个具备可视化能力的 Agent 可以:

1.2 工具栈选择

工具 用途 优势
Matplotlib 静态图表生成 成熟稳定,适合 PDF/PNG 导出
Plotly 交互式图表 鼠标悬停、缩放、点击交互
Grafana 生产仪表盘 实时数据源、告警规则、团队共享
Streamlit 快速应用搭建 纯 Python、组件丰富、部署简单

1.3 报告类型

1.4 数据流架构

MSG Chain (链上数据)
    │
    ▼
Indexer / RPC Node (数据采集)
    │
    ├─► DataCollector (聚合、清洗)
    │       │
    │       ├─► Matplotlib Visualizer (静态图表)
    │       ├─► Plotly Dashboard (交互图表)
    │       └─► Streamlit App (实时看板)
    │
    ├─► ReportGenerator (HTML → IPFS)
    │       │
    │       └─► ReportDistributor (Email/Discord/Twitter/A2A)
    │
    └─► Grafana DataSource (Prometheus/InfluxDB)

2. 链上数据获取

2.1 MSG Chain 数据源

Agent 可以通过以下方式获取链上数据:

  1. RPC 节点:直接查询链状态 (msgd query ...)
  2. Indexer API:专用索引服务,提供历史聚合数据
  3. 事件日志:订阅链上事件,实时接收更新
  4. REST API:LCD (Light Client Daemon) 接口

2.2 Indexer 客户端封装

import time
import json
import asyncio
from typing import Optional, List, Dict, Any
from datetime import datetime, timedelta
import aiohttp
import logging

logger = logging.getLogger(__name__)


class IndexerClient:
    """MSG Chain Indexer HTTP 客户端"""

    BASE_URL = "https://indexer.msgchain.org/api/v1"

    def __init__(self, base_url: str = None):
        self.base_url = base_url or self.BASE_URL
        self.session: Optional[aiohttp.ClientSession] = None

    async def __aenter__(self):
        self.session = aiohttp.ClientSession(
            headers={"Content-Type": "application/json"}
        )
        return self

    async def __aexit__(self, *args):
        if self.session:
            await self.session.close()

    async def _get(self, path: str, params: dict = None) -> dict:
        if not self.session:
            self.session = aiohttp.ClientSession()
        url = f"{self.base_url}{path}"
        async with self.session.get(url, params=params) as resp:
            resp.raise_for_status()
            return await resp.json()

    async def count_agents(self) -> int:
        data = await self._get("/agents/count")
        return data["count"]

    async def total_volume_24h(self) -> float:
        data = await self._get("/stats/volume24h")
        return float(data["volume"])

    async def top_earners(self, limit: int = 10) -> List[dict]:
        return await self._get("/agents/top", {"limit": limit})

    async def capability_stats(self) -> Dict[str, int]:
        data = await self._get("/stats/capabilities")
        return {item["capability"]: item["count"] for item in data["items"]}

    async def query_agent_daily(
        self, agent_id: str, timestamp: int
    ) -> "AgentDailyStats":
        data = await self._get(
            f"/agents/{agent_id}/daily",
            {"timestamp": timestamp},
        )
        return AgentDailyStats(
            requests=data["requests"],
            earnings=data["earnings"],
            errors=data.get("errors", 0),
            timestamp=timestamp,
        )

    async def query_reputation_history(
        self, agent_id: str, limit: int = 100
    ) -> List[dict]:
        return await self._get(
            f"/agents/{agent_id}/reputation",
            {"limit": limit},
        )

    async def query_agent_balance(self, agent_id: str) -> float:
        data = await self._get(f"/agents/{agent_id}/balance")
        return float(data["balance"])

    async def query_agent_latency(self, agent_id: str) -> Dict[str, float]:
        data = await self._get(f"/agents/{agent_id}/latency")
        return {
            "avg": float(data.get("avg", 0)),
            "p50": float(data.get("p50", 0)),
            "p95": float(data.get("p95", 0)),
            "p99": float(data.get("p99", 0)),
        }


class AgentDailyStats:
    """Agent 单日统计数据"""

    def __init__(
        self,
        requests: int = 0,
        earnings: str = "0",
        errors: int = 0,
        timestamp: int = 0,
    ):
        self.requests = requests
        self.earnings = earnings
        self.errors = errors
        self.timestamp = timestamp

    @property
    def error_rate(self) -> float:
        if self.requests == 0:
            return 0.0
        return self.errors / self.requests

2.3 数据收集器

class DataCollector:
    """采集链上数据,供可视化与报告模块使用"""

    def __init__(self, indexer: Optional[IndexerClient] = None):
        self.indexer = indexer or IndexerClient()
        self._cache: Dict[str, Any] = {}

    async def collect_agent_stats(
        self, agent_id: str, days: int = 30
    ) -> Dict[str, list]:
        """采集 Agent 指定天数的历史统计数据"""
        stats: Dict[str, list] = {
            "daily_requests": [],
            "daily_earnings": [],
            "daily_errors": [],
            "reputation_history": [],
            "balance_history": [],
            "latency_history": [],
            "timestamps": [],
        }

        now = int(time.time())

        for day_offset in range(days):
            timestamp = now - day_offset * 86400
            try:
                daily = await self.indexer.query_agent_daily(agent_id, timestamp)
                stats["daily_requests"].insert(0, daily.requests)
                stats["daily_earnings"].insert(0, float(daily.earnings))
                stats["daily_errors"].insert(0, daily.errors)
                stats["timestamps"].insert(0, timestamp)
            except Exception as e:
                logger.warning(f"Failed to fetch day {day_offset}: {e}")
                stats["daily_requests"].insert(0, 0)
                stats["daily_earnings"].insert(0, 0.0)
                stats["daily_errors"].insert(0, 0)
                stats["timestamps"].insert(0, timestamp)

        reputation_data = await self.indexer.query_reputation_history(agent_id)
        stats["reputation_history"] = [
            {
                "score": item.get("score", 0),
                "change": item.get("change", 0),
                "reason": item.get("reason", ""),
                "timestamp": item.get("timestamp", 0),
            }
            for item in reputation_data
        ]

        try:
            balance = await self.indexer.query_agent_balance(agent_id)
            stats["balance_history"] = self._simulate_balance_history(
                balance, stats["daily_earnings"]
            )
        except Exception:
            stats["balance_history"] = [0.0] * days

        try:
            latency = await self.indexer.query_agent_latency(agent_id)
            stats["latency_history"] = self._simulate_latency_history(
                latency, days
            )
        except Exception:
            stats["latency_history"] = [0.0] * days

        return stats

    def _simulate_balance_history(
        self, current_balance: float, earnings: List[float]
    ) -> List[float]:
        """根据当前余额和每日收益反推余额历史"""
        balance = current_balance
        history = []
        for e in reversed(earnings):
            history.append(balance)
            balance -= e
        history.reverse()
        return history

    def _simulate_latency_history(
        self, latency_stats: Dict[str, float], days: int
    ) -> List[float]:
        """根据当前延迟统计数据模拟日延迟历史"""
        base = latency_stats.get("avg", 1.0)
        import random
        random.seed(42)
        history = []
        for _ in range(days):
            jitter = base * 0.2 * (random.random() - 0.5)
            history.append(max(0.05, base + jitter))
        return history

    async def collect_market_data(self) -> dict:
        """采集生态系统全局数据"""
        return {
            "total_agents": await self.indexer.count_agents(),
            "total_volume": await self.indexer.total_volume_24h(),
            "top_agents": await self.indexer.top_earners(10),
            "capability_distribution": await self.indexer.capability_stats(),
        }

    async def collect_agent_detail(self, agent_id: str) -> dict:
        """采集 Agent 详细信息(一次性快照)"""
        balance = await self.indexer.query_agent_balance(agent_id)
        latency = await self.indexer.query_agent_latency(agent_id)
        reputation = await self.indexer.query_reputation_history(agent_id)
        latest_rep = reputation[-1]["score"] if reputation else 0
        return {
            "agent_id": agent_id,
            "balance": balance,
            "latency": latency,
            "latest_reputation": latest_rep,
            "collected_at": int(time.time()),
        }

2.4 链上事件监听器

from websockets import connect
import json


class OnChainEventListener:
    """监听 MSG Chain 实时事件,为仪表盘提供流式数据"""

    def __init__(self, ws_url: str = "wss://rpc.msgchain.org/websocket"):
        self.ws_url = ws_url
        self.listeners = {}
        self._running = False

    def subscribe(self, event_type: str, callback):
        if event_type not in self.listeners:
            self.listeners[event_type] = []
        self.listeners[event_type].append(callback)

    async def start(self):
        self._running = True
        async with connect(self.ws_url) as ws:
            subscribe_msg = json.dumps({
                "jsonrpc": "2.0",
                "method": "subscribe",
                "params": ["tm.event='Tx'"],
                "id": 1,
            })
            await ws.send(subscribe_msg)
            while self._running:
                message = await ws.recv()
                data = json.loads(message)
                event_type = data.get("result", {}).get("query", "unknown")
                for cb in self.listeners.get(event_type, []):
                    await cb(data)

    def stop(self):
        self._running = False


class DataPipeline:
    """数据管道:采集 -> 清洗 -> 存储"""

    def __init__(self):
        self.collector = DataCollector()
        self.raw_buffer = []
        self.processed_buffer = []

    async def fetch_and_process(self, agent_id: str, days: int = 30) -> dict:
        raw = await self.collector.collect_agent_stats(agent_id, days)
        return self.clean(raw)

    def clean(self, raw: dict) -> dict:
        cleaned = {}
        for key, values in raw.items():
            if key == "reputation_history":
                cleaned[key] = values
                continue
            numeric = [v if isinstance(v, (int, float)) else 0 for v in values]
            cleaned[key] = self._fill_missing(numeric)
        return cleaned

    def _fill_missing(self, data):
        last_valid = 0.0
        result = []
        for v in data:
            if v == 0 and last_valid > 0:
                result.append(last_valid)
            else:
                result.append(v)
                if v > 0:
                    last_valid = v
        return result

    def detect_anomalies(self, data, threshold: float = 3.0):
        import numpy as np
        arr = np.array(data, dtype=float)
        mean = np.nanmean(arr)
        std = np.nanstd(arr)
        if std == 0:
            return []
        z_scores = np.abs((arr - mean) / std)
        return np.where(z_scores > threshold)[0].tolist()

2.5 MSG Chain 特定查询工具

class MsgChainQuerier:
    """封装 msgd CLI / REST 查询"""

    MSG_DECIMALS = 6

    def __init__(self, rest_url: str = "https://rest.msgchain.org"):
        self.rest_url = rest_url

    def convert_to_msg(self, amount: str) -> float:
        return int(amount) / (10 ** self.MSG_DECIMALS)

    async def query_agent_contract(self, contract_addr: str) -> dict:
        import aiohttp
        async with aiohttp.ClientSession() as session:
            url = f"{self.rest_url}/cosmwasm/wasm/v1/contract/{contract_addr}"
            async with session.get(url) as resp:
                return await resp.json()

    async def query_agent_txs(self, agent_id: str, limit: int = 100) -> list:
        import aiohttp
        async with aiohttp.ClientSession() as session:
            url = (f"{self.rest_url}/cosmos/tx/v1beta1/txs"
                   f"?events=agent_id%3D{agent_id}&pagination.limit={limit}")
            async with session.get(url) as resp:
                data = await resp.json()
                return data.get("tx_responses", [])

    def parse_tx_for_metrics(self, tx: dict) -> dict:
        gas_wanted = int(tx.get("gas_wanted", 0))
        gas_used = int(tx.get("gas_used", 0))
        height = int(tx.get("height", 0))
        timestamp = tx.get("timestamp", "")
        metrics = {
            "height": height,
            "gas_wanted": gas_wanted,
            "gas_used": gas_used,
            "gas_efficiency": gas_used / max(gas_wanted, 1),
            "timestamp": timestamp,
        }
        for log in tx.get("logs", []):
            for event in log.get("events", []):
                if event["type"] == "agent_action":
                    attrs = {a["key"]: a["value"] for a in event.get("attributes", [])}
                    metrics["action_type"] = attrs.get("action", "unknown")
                    metrics["tokens_used"] = int(attrs.get("tokens", 0))
                    metrics["response_time"] = float(attrs.get("response_ms", 0))
        return metrics

3. Matplotlib 可视化

3.1 基础配置

import matplotlib
matplotlib.use("Agg")

import matplotlib.pyplot as plt
import matplotlib.dates as mdates
import matplotlib.ticker as mticker
import numpy as np
import pandas as pd
from typing import List, Optional, Dict, Tuple
from datetime import datetime, timedelta
import os

plt.rcParams.update({
    "font.family": "DejaVu Sans",
    "axes.unicode_minus": False,
    "figure.dpi": 150,
    "savefig.dpi": 150,
    "savefig.bbox": "tight",
    "figure.facecolor": "#FAFAFA",
    "axes.facecolor": "#FFFFFF",
    "axes.edgecolor": "#E0E0E0",
    "axes.grid": True,
    "grid.alpha": 0.3,
    "grid.linestyle": "--",
})

3.2 Agent 可视化器

class AgentVisualizer:
    """将 Agent 数据转换为 Matplotlib 图表"""

    COLOR_PRIMARY = "#2196F3"
    COLOR_SUCCESS = "#4CAF50"
    COLOR_DANGER = "#F44336"
    COLOR_WARNING = "#FF9800"
    COLOR_PURPLE = "#9C27B0"
    COLOR_TEAL = "#009688"

    def __init__(self, output_dir: str = "./charts"):
        self.output_dir = output_dir
        os.makedirs(output_dir, exist_ok=True)

    def _prepare_dates(self, n: int, freq: str = "D") -> pd.DatetimeIndex:
        return pd.date_range(end=pd.Timestamp.now(), periods=n, freq=freq)

    def _save_figure(self, fig: plt.Figure, name: str) -> str:
        path = os.path.join(self.output_dir, name)
        fig.savefig(path)
        plt.close(fig)
        return path

3.3 收益趋势图

    def plot_earnings_trend(self, stats: dict, save_path: str = "earnings_trend.png") -> str:
        """日收益柱状图 + 累计收益折线图"""
        earnings = stats.get("daily_earnings", [])
        if not earnings:
            return ""

        dates = self._prepare_dates(len(earnings))
        cumulative = pd.Series(earnings).cumsum()

        fig, (ax1, ax2) = plt.subplots(2, 1, figsize=(14, 9))

        colors = [self.COLOR_SUCCESS if v > 0 else self.COLOR_DANGER for v in earnings]
        ax1.bar(dates, earnings, color=colors, alpha=0.75, edgecolor="white", width=0.7)
        ax1.set_title("Daily Earnings (MSG)", fontsize=15, fontweight="bold", pad=12)
        ax1.set_ylabel("MSG", fontsize=12)
        ax1.xaxis.set_major_formatter(mdates.DateFormatter("%m-%d"))
        ax1.tick_params(axis="x", rotation=45)
        ax1.axhline(y=np.mean(earnings), color="gray", linestyle=":", alpha=0.6,
                    label=f"Avg: {np.mean(earnings):.1f} MSG")
        ax1.legend(fontsize=10)

        ax2.plot(dates, cumulative, color=self.COLOR_PRIMARY, linewidth=2.5,
                 marker="o", markersize=4)
        ax2.fill_between(dates, cumulative, alpha=0.15, color=self.COLOR_PRIMARY)
        ax2.set_title("Cumulative Earnings", fontsize=15, fontweight="bold", pad=12)
        ax2.set_ylabel("MSG", fontsize=12)
        ax2.xaxis.set_major_formatter(mdates.DateFormatter("%m-%d"))
        ax2.tick_params(axis="x", rotation=45)

        total = cumulative.iloc[-1] if len(cumulative) > 0 else 0
        ax2.annotate(f"Total: {total:.2f} MSG",
                     xy=(dates[-1], cumulative.iloc[-1]),
                     xytext=(dates[-1] - timedelta(days=len(dates) * 0.3),
                             cumulative.iloc[-1] * 0.7),
                     arrowprops=dict(arrowstyle="->", color="black", alpha=0.5),
                     fontsize=12, fontweight="bold")

        fig.suptitle("Agent Earnings Overview", fontsize=16, fontweight="bold", y=1.01)
        plt.tight_layout()
        return self._save_figure(fig, save_path)

3.4 请求量柱状图

    def plot_request_volume(self, stats: dict, save_path: str = "request_volume.png") -> str:
        """请求量直方图和滚动均值"""
        requests = stats.get("daily_requests", [])
        if not requests:
            return ""

        dates = self._prepare_dates(len(requests))
        series = pd.Series(requests, index=dates)
        rolling_avg = series.rolling(window=7, min_periods=1).mean()

        fig, ax = plt.subplots(figsize=(14, 6))

        ax.bar(dates, requests, color=self.COLOR_PRIMARY, alpha=0.5,
               label="Daily Requests", width=0.7)
        ax.plot(dates, rolling_avg, color=self.COLOR_DANGER, linewidth=2.5,
                label="7-Day Moving Average")

        peak = series.max()
        peak_date = series.idxmax()
        ax.annotate(f"Peak: {int(peak)}",
                    xy=(peak_date, peak),
                    xytext=(peak_date, peak * 1.15),
                    arrowprops=dict(arrowstyle="->", color=self.COLOR_DANGER, alpha=0.6),
                    fontsize=11, color=self.COLOR_DANGER, fontweight="bold")

        ax.set_title("Daily Request Volume", fontsize=15, fontweight="bold")
        ax.set_ylabel("Requests", fontsize=12)
        ax.set_xlabel("Date", fontsize=12)
        ax.xaxis.set_major_formatter(mdates.DateFormatter("%m-%d"))
        ax.tick_params(axis="x", rotation=45)
        ax.legend(fontsize=11, loc="upper left")
        ax.yaxis.set_major_formatter(mticker.FuncFormatter(lambda x, _: f"{int(x):,}"))

        plt.tight_layout()
        return self._save_figure(fig, save_path)

3.5 错误率时间序列

    def plot_error_rate(self, stats: dict, save_path: str = "error_rate.png") -> str:
        """错误率时间序列及告警阈值标注"""
        requests = stats.get("daily_requests", [])
        errors = stats.get("daily_errors", [])
        if not requests or not errors:
            return ""

        error_rate = [e / r * 100 if r > 0 else 0 for e, r in zip(errors, requests)]
        dates = self._prepare_dates(len(error_rate))

        fig, ax = plt.subplots(figsize=(14, 6))

        ax.fill_between(dates, error_rate, alpha=0.2, color=self.COLOR_DANGER)
        ax.plot(dates, error_rate, color=self.COLOR_DANGER, linewidth=2,
                marker="s", markersize=4, label="Error Rate")

        ax.axhline(y=5, color=self.COLOR_WARNING, linestyle="--", linewidth=1.5,
                   label="Warning (5%)")
        ax.axhline(y=10, color=self.COLOR_DANGER, linestyle="--", linewidth=1.5,
                   label="Critical (10%)")

        for i, rate in enumerate(error_rate):
            if rate > 10:
                ax.plot(dates[i], rate, "v", color=self.COLOR_DANGER, markersize=10, zorder=5)

        ax.set_title("Daily Error Rate", fontsize=15, fontweight="bold")
        ax.set_ylabel("Error Rate (%)", fontsize=12)
        ax.set_xlabel("Date", fontsize=12)
        ax.xaxis.set_major_formatter(mdates.DateFormatter("%m-%d"))
        ax.tick_params(axis="x", rotation=45)
        ax.legend(fontsize=11, loc="upper right")
        ax.set_ylim(bottom=0)

        plt.tight_layout()
        return self._save_figure(fig, save_path)

3.6 声誉变化图

    def plot_reputation_history(self, reputation_data: list, save_path: str = "reputation.png") -> str:
        """声誉变化折线图"""
        if not reputation_data:
            return ""

        scores = [item["score"] for item in reputation_data]
        changes = [item.get("change", 0) for item in reputation_data]
        reasons = [item.get("reason", "") for item in reputation_data]
        timestamps = [datetime.fromtimestamp(item.get("timestamp", 0)) for item in reputation_data]

        fig, (ax1, ax2) = plt.subplots(2, 1, figsize=(14, 8), gridspec_kw={"height_ratios": [3, 1]})

        ax1.plot(timestamps, scores, color=self.COLOR_PURPLE, linewidth=2.5, marker="o", markersize=6)
        ax1.fill_between(timestamps, 0, scores, alpha=0.1, color=self.COLOR_PURPLE)
        ax1.set_title("Agent Reputation Score Over Time", fontsize=15, fontweight="bold")
        ax1.set_ylabel("Reputation Score", fontsize=12)
        ax1.set_ylim(0, 100)
        ax1.axhline(y=70, color=self.COLOR_WARNING, linestyle=":", alpha=0.5, label="Good Threshold")
        ax1.axhline(y=40, color=self.COLOR_DANGER, linestyle=":", alpha=0.5, label="Risk Threshold")
        ax1.legend(fontsize=10)

        bar_colors = [self.COLOR_SUCCESS if c >= 0 else self.COLOR_DANGER for c in changes]
        ax2.bar(timestamps, changes, color=bar_colors, alpha=0.7, edgecolor="white", width=0.6)
        ax2.set_title("Score Changes", fontsize=12)
        ax2.set_ylabel("Change", fontsize=11)
        ax2.axhline(y=0, color="black", linewidth=0.5)
        ax2.xaxis.set_major_formatter(mdates.DateFormatter("%m-%d"))
        ax2.tick_params(axis="x", rotation=45)

        plt.tight_layout()
        return self._save_figure(fig, save_path)

3.7 能力雷达图

    def plot_performance_radar(self, capabilities: Dict[str, float], save_path: str = "capability_radar.png") -> str:
        """能力表现雷达图"""
        if not capabilities:
            return ""

        categories = list(capabilities.keys())
        values = list(capabilities.values())
        num_vars = len(categories)

        angles = np.linspace(0, 2 * np.pi, num_vars, endpoint=False).tolist()
        values += values[:1]
        angles += angles[:1]

        fig, ax = plt.subplots(figsize=(9, 9), subplot_kw={"polar": True})
        ax.set_theta_offset(np.pi / 2)
        ax.set_theta_direction(-1)

        ax.plot(angles, values, "o-", linewidth=2, color=self.COLOR_TEAL, markersize=8)
        ax.fill(angles, values, alpha=0.25, color=self.COLOR_TEAL)

        ax.set_thetagrids(np.degrees(angles[:-1]), categories, fontsize=11)
        ax.set_ylim(0, 100)
        ax.set_yticks([20, 40, 60, 80, 100])
        ax.set_yticklabels(["20", "40", "60", "80", "100"], fontsize=9, color="gray")
        ax.set_title("Agent Capability Performance", fontsize=15, fontweight="bold", pad=20)

        plt.tight_layout()
        return self._save_figure(fig, save_path)

3.8 延迟分布图

    def plot_latency_distribution(self, latency_data: list, save_path: str = "latency_distribution.png") -> str:
        """响应延迟分布直方图 + 百分位线"""
        if not latency_data:
            return ""

        fig, (ax1, ax2) = plt.subplots(1, 2, figsize=(14, 5))

        ax1.hist(latency_data, bins=30, color=self.COLOR_PRIMARY, alpha=0.7, edgecolor="white")
        ax1.axvline(np.median(latency_data), color=self.COLOR_SUCCESS, linestyle="--", linewidth=2,
                    label=f"Median: {np.median(latency_data):.1f}s")
        ax1.axvline(np.percentile(latency_data, 95), color=self.COLOR_DANGER, linestyle="--", linewidth=2,
                    label=f"P95: {np.percentile(latency_data, 95):.1f}s")
        ax1.set_title("Latency Distribution", fontsize=14, fontweight="bold")
        ax1.set_xlabel("Response Time (s)", fontsize=11)
        ax1.set_ylabel("Frequency", fontsize=11)
        ax1.legend(fontsize=10)

        percentiles = [50, 75, 90, 95, 99]
        pvalues = [np.percentile(latency_data, p) for p in percentiles]
        bars = ax2.bar([str(p) for p in percentiles], pvalues,
                       color=[self.COLOR_PRIMARY, self.COLOR_SUCCESS, self.COLOR_WARNING,
                              self.COLOR_DANGER, self.COLOR_PURPLE],
                       alpha=0.8, edgecolor="white")
        for bar, val in zip(bars, pvalues):
            ax2.text(bar.get_x() + bar.get_width() / 2, bar.get_height(),
                     f"{val:.1f}s", ha="center", va="bottom", fontsize=11, fontweight="bold")
        ax2.set_title("Latency Percentiles", fontsize=14, fontweight="bold")
        ax2.set_xlabel("Percentile", fontsize=11)
        ax2.set_ylabel("Response Time (s)", fontsize=11)

        plt.tight_layout()
        return self._save_figure(fig, save_path)

3.9 Gas 消耗分析图

    def plot_gas_analysis(self, tx_metrics: pd.DataFrame, save_path: str = "gas_analysis.png") -> str:
        """Gas 消耗分析"""
        if tx_metrics.empty:
            return ""

        fig, ((ax1, ax2), (ax3, ax4)) = plt.subplots(2, 2, figsize=(16, 10))

        ax1.scatter(tx_metrics["gas_wanted"], tx_metrics["gas_used"], alpha=0.5, s=20, c=self.COLOR_PRIMARY)
        max_gas = max(tx_metrics["gas_wanted"].max(), tx_metrics["gas_used"].max())
        ax1.plot([0, max_gas], [0, max_gas], "r--", alpha=0.5, label="y=x (ideal)")
        ax1.set_title("Gas Wanted vs Used", fontsize=13, fontweight="bold")
        ax1.set_xlabel("Gas Wanted")
        ax1.set_ylabel("Gas Used")
        ax1.legend(fontsize=10)
        ax1.set_aspect("equal")

        efficiency = tx_metrics["gas_efficiency"]
        ax3.hist(efficiency, bins=30, color=self.COLOR_SUCCESS, alpha=0.7, edgecolor="white")
        ax3.axvline(x=efficiency.mean(), color=self.COLOR_DANGER, linestyle="--", linewidth=2,
                    label=f"Mean: {efficiency.mean():.2f}")
        ax3.set_title("Gas Efficiency Distribution", fontsize=13, fontweight="bold")
        ax3.set_xlabel("Efficiency")
        ax3.set_ylabel("Count")
        ax3.legend(fontsize=10)

        gas_cost = tx_metrics["gas_used"] * 0.001
        ax4.plot(gas_cost.cumsum(), color=self.COLOR_PURPLE, linewidth=2)
        ax4.set_title("Cumulative Gas Cost", fontsize=13, fontweight="bold")
        ax4.set_ylabel("Cost (MSG)")

        plt.tight_layout()
        return self._save_figure(fig, save_path)

3.10 综合性能面板

    def plot_performance_panel(self, stats: dict, save_path: str = "performance_panel.png") -> str:
        """综合性能面板:4 合一视图"""
        requests = stats.get("daily_requests", [])
        earnings = stats.get("daily_earnings", [])
        errors = stats.get("daily_errors", [])
        latency = stats.get("latency_history", [])

        if not requests:
            return ""

        dates = self._prepare_dates(len(requests))
        error_rate = [e / r * 100 if r > 0 else 0 for e, r in zip(errors, requests)]

        fig, ((ax1, ax2), (ax3, ax4)) = plt.subplots(2, 2, figsize=(16, 10))

        ax1.bar(dates, requests, color=self.COLOR_PRIMARY, alpha=0.6, width=0.7)
        ax1.set_title("Daily Requests", fontsize=13, fontweight="bold")
        ax1.set_ylabel("Count")
        ax1.xaxis.set_major_formatter(mdates.DateFormatter("%m-%d"))
        ax1.tick_params(axis="x", rotation=45)

        cumulative = pd.Series(earnings).cumsum()
        ax2.fill_between(dates, cumulative, alpha=0.3, color=self.COLOR_SUCCESS)
        ax2.plot(dates, cumulative, color=self.COLOR_SUCCESS, linewidth=2)
        ax2.set_title("Cumulative Earnings (MSG)", fontsize=13, fontweight="bold")
        ax2.set_ylabel("MSG")
        ax2.xaxis.set_major_formatter(mdates.DateFormatter("%m-%d"))
        ax2.tick_params(axis="x", rotation=45)

        ax3.plot(dates, error_rate, color=self.COLOR_DANGER, linewidth=2)
        ax3.fill_between(dates, error_rate, alpha=0.15, color=self.COLOR_DANGER)
        ax3.axhline(y=5, color=self.COLOR_WARNING, linestyle="--", alpha=0.5, label="5% Warning")
        ax3.set_title("Error Rate", fontsize=13, fontweight="bold")
        ax3.set_ylabel("%")
        ax3.xaxis.set_major_formatter(mdates.DateFormatter("%m-%d"))
        ax3.tick_params(axis="x", rotation=45)
        ax3.legend(fontsize=9)

        if latency:
            ax4.plot(dates, latency, color=self.COLOR_PURPLE, linewidth=2)
            ax4.fill_between(dates, latency, alpha=0.15, color=self.COLOR_PURPLE)
            ax4.set_title("Avg Response Time", fontsize=13, fontweight="bold")
            ax4.set_ylabel("Seconds")
        else:
            ax4.text(0.5, 0.5, "No Data", ha="center", va="center", transform=ax4.transAxes,
                     fontsize=14, color="gray")
        ax4.xaxis.set_major_formatter(mdates.DateFormatter("%m-%d"))
        ax4.tick_params(axis="x", rotation=45)

        fig.suptitle("Agent Performance Panel", fontsize=16, fontweight="bold", y=1.01)
        plt.tight_layout()
        return self._save_figure(fig, save_path)

3.11 对比分析图

    def plot_agent_comparison(self, agents_data: Dict[str, dict], save_path: str = "agent_comparison.png") -> str:
        """多 Agent 对比条形图"""
        if not agents_data:
            return ""

        names = list(agents_data.keys())
        metrics_data = {
            "7d Requests": [],
            "7d Earnings": [],
            "Reputation": [],
            "Success Rate": [],
        }

        for name in names:
            data = agents_data[name]
            metrics_data["7d Requests"].append(sum(data.get("daily_requests", [])[-7:]))
            metrics_data["7d Earnings"].append(sum(data.get("daily_earnings", [])[-7:]))
            metrics_data["Reputation"].append(data.get("latest_reputation", 0))
            reqs = data.get("daily_requests", [])
            errs = data.get("daily_errors", [])
            total_req = sum(reqs[-7:])
            total_err = sum(errs[-7:])
            success_rate = (total_req - total_err) / total_req * 100 if total_req > 0 else 100
            metrics_data["Success Rate"].append(success_rate)

        fig, axes = plt.subplots(2, 2, figsize=(16, 10))
        colors_list = [self.COLOR_PRIMARY, self.COLOR_SUCCESS, self.COLOR_PURPLE, self.COLOR_TEAL]

        for ax, (metric, values), color in zip(axes.flatten(), metrics_data.items(), colors_list):
            bars = ax.barh(names, values, color=color, alpha=0.75, edgecolor="white")
            for bar, val in zip(bars, values):
                ax.text(bar.get_width() + max(values) * 0.01,
                        bar.get_y() + bar.get_height() / 2,
                        f"{val:.1f}", va="center", fontsize=10)
            ax.set_title(metric, fontsize=13, fontweight="bold")
            ax.tick_params(labelsize=10)

        fig.suptitle("Agent Comparison (7-Day)", fontsize=16, fontweight="bold", y=1.01)
        plt.tight_layout()
        return self._save_figure(fig, save_path)

4. Plotly 交互式图表

4.1 基础配置

import plotly.graph_objects as go
import plotly.express as px
from plotly.subplots import make_subplots
from plotly.io import to_html, to_image
import numpy as np
import pandas as pd
from typing import Optional, List, Dict, Any


class InteractiveDashboard:
    """基于 Plotly 的交互式仪表盘"""

    THEME = {
        "bg_color": "#FAFAFA",
        "font_color": "#333333",
        "primary": "#2196F3",
        "success": "#4CAF50",
        "danger": "#F44336",
        "warning": "#FF9800",
        "purple": "#9C27B0",
        "teal": "#009688",
    }

    def _apply_theme(self, fig: go.Figure) -> go.Figure:
        fig.update_layout(
            template="plotly_white",
            paper_bgcolor=self.THEME["bg_color"],
            plot_bgcolor="#FFFFFF",
            font={"color": self.THEME["font_color"], "family": "DejaVu Sans"},
            hoverlabel={"bgcolor": "#333333", "font_size": 12, "font_color": "#FFFFFF"},
        )
        return fig

4.2 Agent 性能仪表盘

    def create_agent_performance_dashboard(self, stats: dict, title: str = "Agent Performance Dashboard") -> go.Figure:
        """创建四象限 Agent 性能仪表盘"""
        fig = make_subplots(
            rows=2, cols=2,
            subplot_titles=("Daily Requests", "Daily Earnings", "Error Rate", "Balance History"),
            specs=[[{"type": "scatter"}, {"type": "bar"}],
                   [{"type": "scatter"}, {"type": "scatter"}]],
            vertical_spacing=0.12, horizontal_spacing=0.08,
        )

        n = len(stats.get("daily_requests", []))
        requests = stats.get("daily_requests", [])
        fig.add_trace(go.Scatter(y=requests, mode="lines+markers", name="Requests",
                                 line={"color": self.THEME["primary"], "width": 2},
                                 marker={"size": 6},
                                 hovertemplate="Day %{x}<br>Requests: %{y:,}<extra></extra>"),
                      row=1, col=1)

        earnings = stats.get("daily_earnings", [])
        cum_earnings = pd.Series(earnings).cumsum()
        fig.add_trace(go.Bar(y=earnings, name="Earnings",
                             marker={"color": self.THEME["success"], "opacity": 0.7},
                             hovertemplate="Day %{x}<br>Earnings: %{y:.2f} MSG<extra></extra>"),
                      row=1, col=2)
        fig.add_trace(go.Scatter(y=cum_earnings, mode="lines", name="Cumulative",
                                 line={"color": self.THEME["teal"], "width": 2, "dash": "dot"}),
                      row=1, col=2)

        requests_arr = np.array(requests, dtype=float)
        errors_arr = np.array(stats.get("daily_errors", [0] * n), dtype=float)
        error_rate = np.divide(errors_arr, requests_arr, out=np.zeros_like(errors_arr),
                               where=requests_arr > 0) * 100

        fig.add_trace(go.Scatter(y=error_rate, mode="lines", name="Error Rate",
                                 line={"color": self.THEME["danger"], "width": 2},
                                 fill="tozeroy", fillcolor="rgba(244, 67, 54, 0.1)",
                                 hovertemplate="Day %{x}<br>Error Rate: %{y:.2f}%<extra></extra>"),
                      row=2, col=1)
        fig.add_hline(y=5, line={"color": self.THEME["warning"], "dash": "dash", "width": 1},
                      annotation_text="Warning (5%)", row=2, col=1)
        fig.add_hline(y=10, line={"color": self.THEME["danger"], "dash": "dash", "width": 1},
                      annotation_text="Critical (10%)", row=2, col=1)

        balance = stats.get("balance_history", [])
        if balance:
            fig.add_trace(go.Scatter(y=balance, mode="lines", name="Balance",
                                     line={"color": self.THEME["purple"], "width": 2},
                                     fill="tozeroy", fillcolor="rgba(156, 39, 176, 0.1)"),
                          row=2, col=2)

        fig.update_layout(height=800, title_text=title, title_font={"size": 18},
                          showlegend=True,
                          legend={"orientation": "h", "yanchor": "bottom", "y": 1.02,
                                  "xanchor": "right", "x": 1})
        return self._apply_theme(fig)

4.3 收益瀑布图

    def create_earnings_waterfall(self, stats: dict, title: str = "Earnings Waterfall") -> go.Figure:
        """收益瀑布图"""
        earnings = stats.get("daily_earnings", [])
        if not earnings:
            return go.Figure()

        fig = go.Figure(go.Waterfall(
            name="Earnings", orientation="v",
            measure=["relative"] * len(earnings),
            x=[f"Day {i+1}" for i in range(len(earnings))],
            y=earnings,
            text=[f"{v:.2f}" for v in earnings],
            textposition="outside",
            connector={"line": {"color": "#666666", "width": 1}},
            increasing={"marker": {"color": self.THEME["success"]}},
            decreasing={"marker": {"color": self.THEME["danger"]}},
            totals={"marker": {"color": self.THEME["primary"]}},
        ))
        fig.update_layout(title=title, height=500, showlegend=False, yaxis_title="MSG")
        return self._apply_theme(fig)

4.4 声誉 Timeline

    def create_reputation_timeline(self, reputation_data: list, title: str = "Reputation Timeline") -> go.Figure:
        """声誉变化时间线"""
        if not reputation_data:
            return go.Figure()

        df = pd.DataFrame(reputation_data)
        df["datetime"] = pd.to_datetime(df["timestamp"], unit="s")
        df["color"] = df["change"].apply(lambda c: self.THEME["success"] if c >= 0 else self.THEME["danger"])

        fig = go.Figure(go.Scatter(
            x=df["datetime"], y=df["score"],
            mode="lines+markers", name="Score",
            line={"color": self.THEME["purple"], "width": 2},
            marker={"size": 8, "color": df["color"],
                    "line": {"color": "#FFFFFF", "width": 1}},
            text=df["reason"],
            hovertemplate="<b>%{x|%Y-%m-%d %H:%M}</b><br>Score: %{y}<br>Reason: %{text}<extra></extra>",
        ))

        fig.update_layout(title=title, height=500, yaxis={"range": [0, 100], "title": "Score"},
                          xaxis={"title": "Date"}, hovermode="x unified", showlegend=False)
        return self._apply_theme(fig)

4.5 生态系统全局仪表盘

    def create_ecosystem_dashboard(self, market_data: dict, title: str = "MSG Chain Ecosystem Dashboard") -> go.Figure:
        """生态系统全局仪表盘"""
        fig = make_subplots(
            rows=2, cols=3,
            specs=[[{"type": "indicator"}, {"type": "indicator"}, {"type": "indicator"}],
                   [{"colspan": 2, "type": "bar"}, {"type": "pie"}]],
            subplot_titles=("Total Agents", "24h Volume", "Top Agent Earnings",
                           "Top Earners", "Capability Distribution"),
            column_widths=[0.33, 0.33, 0.33], row_heights=[0.4, 0.6],
            vertical_spacing=0.12, horizontal_spacing=0.05,
        )

        fig.add_trace(go.Indicator(mode="number", value=market_data.get("total_agents", 0),
                                   number={"font": {"size": 40, "color": self.THEME["primary"]}}}),
                      row=1, col=1)

        volume = market_data.get("total_volume", 0)
        fig.add_trace(go.Indicator(mode="number+delta", value=volume,
                                   number={"prefix": "MSG ", "font": {"size": 36, "color": self.THEME["success"]}},
                                   delta={"reference": volume * 0.9, "relative": True, "valueformat": ".1%"}),
                      row=1, col=2)

        top_agents = market_data.get("top_agents", [])
        top_earn = float(top_agents[0].get("earnings", 0)) if top_agents else 0
        fig.add_trace(go.Indicator(mode="number", value=top_earn,
                                   number={"prefix": "MSG ", "font": {"size": 36, "color": self.THEME["warning"]}}}),
                      row=1, col=3)

        if top_agents:
            names = [a.get("name", a.get("id", f"Agent {i}"))[:12] for i, a in enumerate(top_agents)]
            earnings_list = [float(a.get("earnings", 0)) for a in top_agents]
            fig.add_trace(go.Bar(x=names, y=earnings_list,
                                 marker={"color": earnings_list, "colorscale": "Viridis", "showscale": True}),
                          row=2, col=1)

        cap_dist = market_data.get("capability_distribution", {})
        if cap_dist:
            fig.add_trace(go.Pie(labels=list(cap_dist.keys()), values=list(cap_dist.values()),
                                 hole=0.4, textinfo="label+percent",
                                 marker={"colors": px.colors.qualitative.Set3[:len(cap_dist)]}),
                          row=2, col=3)

        fig.update_layout(height=700, title_text=title, title_font={"size": 18}, showlegend=False)
        return self._apply_theme(fig)

4.6 热力图

    def create_activity_heatmap(self, activity_data: pd.DataFrame, title: str = "Agent Activity Heatmap") -> go.Figure:
        """Agent 活动热力图(小时 x 星期)"""
        if activity_data.empty:
            return go.Figure()

        if "hour" not in activity_data.columns or "day_of_week" not in activity_data.columns:
            activity_data["hour"] = activity_data.index.hour
            activity_data["day_of_week"] = activity_data.index.dayofweek

        pivot = activity_data.pivot_table(index="hour", columns="day_of_week",
                                          values="count" if "count" in activity_data.columns else activity_data.columns[0],
                                          aggfunc="count", fill_value=0)
        days = ["Mon", "Tue", "Wed", "Thu", "Fri", "Sat", "Sun"]
        pivot.columns = [days[c] if c < 7 else str(c) for c in pivot.columns]

        fig = go.Figure(data=go.Heatmap(
            z=pivot.values, x=list(pivot.columns), y=list(pivot.index),
            colorscale=[[0, "#FFFFFF"], [0.25, "#C8E6C9"], [0.5, "#66BB6A"],
                        [0.75, "#2196F3"], [1, "#1B5E20"]],
            hovertemplate="Hour: %{y}:00<br>Day: %{x}<br>Activity: %{z}<extra></extra>",
            colorbar={"title": "Count", "thickness": 20},
        ))
        fig.update_layout(title=title, height=400,
                          xaxis={"title": "Day of Week", "side": "top"},
                          yaxis={"title": "Hour of Day", "dtick": 2})
        return self._apply_theme(fig)

4.7 图表导出工具

    def export_html(self, fig: go.Figure, path: str):
        html = to_html(fig, include_plotlyjs="cdn", full_html=True,
                       config={"responsive": True, "displayModeBar": True,
                               "modeBarButtonsToRemove": ["lasso2d", "select2d"]})
        with open(path, "w", encoding="utf-8") as f:
            f.write(html)

    def export_image(self, fig: go.Figure, path: str, format: str = "png",
                     width: int = 1200, height: int = 800):
        img_bytes = to_image(fig, format=format, width=width, height=height)
        with open(path, "wb") as f:
            f.write(img_bytes)

    def create_multi_page_report(self, figures: Dict[str, go.Figure], output_path: str = "dashboard_report.html") -> str:
        pages_html = []
        first = True
        for title, fig in figures.items():
            fig.update_layout(title_text=title)
            page = to_html(fig, include_plotlyjs="cdn" if first else False,
                           full_html=first, config={"responsive": True})
            if not first:
                page = page.replace("</body>", "").replace("</html>", "")
            pages_html.append(page)
            first = False

        full_html = "\n<hr style='page-break-after: always;'>\n".join(pages_html)
        if not full_html.endswith("</html>"):
            full_html += "\n</body>\n</html>"
        with open(output_path, "w", encoding="utf-8") as f:
            f.write(full_html)
        return output_path

5. Streamlit 实时仪表盘

5.1 基础配置

"""
Agent Monitor - Streamlit 实时仪表盘
启动: streamlit run agent_monitor.py --server.port 8501
"""

import streamlit as st
import pandas as pd
import numpy as np
import plotly.graph_objects as go
import plotly.express as px
from datetime import datetime, timedelta
import time
import asyncio
import random
from typing import Optional, Dict, List

st.set_page_config(
    page_title="MSG Chain AI Agent Monitor",
    page_icon="🤖",
    layout="wide",
    initial_sidebar_state="expanded",
    menu_items={
        "Get Help": "https://msgchain.org/docs",
        "Report a bug": "https://github.com/msgchain/agent-monitor/issues",
        "About": "### MSG Chain AI Agent Monitor\n实时监控 Agent 链上活动、收益与声誉。",
    },
)

5.2 样式与辅助函数

st.markdown(
    """
    <style>
    .metric-card {
        background: #FFFFFF; border-radius: 10px; padding: 1.5rem 1rem;
        box-shadow: 0 2px 8px rgba(0,0,0,0.08); text-align: center;
        border-left: 4px solid #2196F3;
    }
    .metric-card.green { border-left-color: #4CAF50; }
    .metric-card.red { border-left-color: #F44336; }
    .stMetric {
        background: #FFFFFF; border-radius: 8px; padding: 1rem;
        box-shadow: 0 1px 4px rgba(0,0,0,0.06);
    }
    .report-header { padding: 1rem 0; border-bottom: 2px solid #2196F3; margin-bottom: 1.5rem; }
    .status-badge {
        display: inline-block; padding: 0.25rem 0.75rem; border-radius: 12px;
        font-size: 0.8rem; font-weight: 600;
    }
    .status-badge.online { background: #E8F5E9; color: #2E7D32; }
    .status-badge.offline { background: #FFEBEE; color: #C62828; }
    .status-badge.warning { background: #FFF3E0; color: #E65100; }
    </style>
    """,
    unsafe_allow_html=True,
)


def format_msg(amount: float) -> str:
    if amount >= 1_000_000:
        return f"{amount / 1_000_000:.2f}M MSG"
    elif amount >= 1_000:
        return f"{amount / 1_000:.2f}K MSG"
    return f"{amount:.2f} MSG"


@st.cache_resource
def get_data_collector():
    from data_collector import DataCollector, IndexerClient
    return DataCollector(IndexerClient())


def create_session_state():
    defaults = {
        "agent_id": "msg1agent...",
        "auto_refresh": True,
        "refresh_interval": 15,
        "data_cache": {},
        "last_fetch": 0,
        "show_anomalies": True,
    }
    for key, value in defaults.items():
        if key not in st.session_state:
            st.session_state[key] = value

5.3 侧边栏

def render_sidebar():
    st.sidebar.image("https://msgchain.org/logo.png", width=200)
    st.sidebar.markdown("## AI Agent Monitor")
    st.sidebar.markdown("---")

    st.sidebar.subheader("🔍 Agent Selection")
    agent_id = st.sidebar.text_input("Agent ID", value=st.session_state.agent_id,
                                     placeholder="msg1...")

    use_preset = st.sidebar.checkbox("Use Preset Agents", value=False)
    if use_preset:
        preset = st.sidebar.selectbox("Select Agent",
            ["msg1agent2u9x7qh0u3h", "msg1agent8m3k5p9w1q", "msg1agent4v6n2r8t0s"])
        agent_id = preset
    st.session_state.agent_id = agent_id

    st.sidebar.subheader("📅 Time Range")
    days = st.sidebar.slider("Days of History", 1, 90, 30)

    st.sidebar.subheader("🔄 Refresh")
    st.session_state.auto_refresh = st.sidebar.checkbox("Auto-refresh", value=st.session_state.auto_refresh)
    st.session_state.refresh_interval = st.sidebar.slider("Interval (seconds)", 5, 120,
                                                           st.session_state.refresh_interval, 5,
                                                           disabled=not st.session_state.auto_refresh)

    st.sidebar.subheader("⚙️ Filters")
    st.session_state.show_anomalies = st.sidebar.checkbox("Highlight Anomalies",
                                                           value=st.session_state.show_anomalies)

    st.sidebar.markdown("---")
    st.sidebar.markdown(f"**Last Updated:** {datetime.now().strftime('%H:%M:%S')}")

    manual_refresh = st.sidebar.button("🔄 Refresh Now", type="primary")
    if manual_refresh:
        st.session_state.last_fetch = 0
        st.rerun()

    return agent_id, days

5.4 KPI 指标卡

def render_kpi_metrics(stats: dict):
    earnings = stats.get("daily_earnings", [0])
    requests = stats.get("daily_requests", [0])
    errors = stats.get("daily_errors", [0])
    latency = stats.get("latency_history", [0])

    latest_earnings = earnings[-1] if earnings else 0
    prev_earnings = earnings[-2] if len(earnings) > 1 else 0
    earn_change = ((latest_earnings - prev_earnings) / prev_earnings * 100) if prev_earnings > 0 else 0

    latest_req = requests[-1] if requests else 0
    prev_req = requests[-2] if len(requests) > 1 else 0
    req_change = ((latest_req - prev_req) / prev_req * 100) if prev_req > 0 else 0

    latest_err = errors[-1] if errors else 0
    latest_lat = latency[-1] if latency else 0

    kpis = [
        {"label": "Daily Requests", "value": f"{latest_req:,}",
         "delta": f"{req_change:+.1f}%", "delta_color": "normal" if req_change >= 0 else "inverse"},
        {"label": "Daily Earnings", "value": format_msg(latest_earnings),
         "delta": f"{earn_change:+.1f}%", "delta_color": "normal" if earn_change >= 0 else "inverse"},
        {"label": "Errors", "value": str(latest_err),
         "delta": "✅" if latest_err == 0 else "⚠️", "delta_color": "off"},
        {"label": "Avg Response", "value": f"{latest_lat:.2f}s",
         "delta": f"{latency[-2] - latest_lat:.2f}s" if len(latency) > 1 else "-",
         "delta_color": "normal" if len(latency) > 1 and latest_lat <= latency[-2] else "inverse"},
    ]

    cols = st.columns(len(kpis))
    for col, kpi in zip(cols, kpis):
        with col:
            st.metric(label=kpi["label"], value=kpi["value"],
                      delta=kpi["delta"], delta_color=kpi.get("delta_color", "normal"))

5.5 标签页 - 性能

def render_performance_tab(stats: dict):
    col_a, col_b = st.columns(2)

    with col_a:
        st.subheader("📈 Request Rate")
        requests = stats.get("daily_requests", [])
        if requests:
            st.line_chart(pd.DataFrame({"Requests": requests}), use_container_width=True, height=300)

        st.subheader("⏱️ Latency P95")
        latency = stats.get("latency_history", [])
        if latency:
            st.line_chart(pd.DataFrame({"Latency (s)": latency}), use_container_width=True, height=300)

    with col_b:
        st.subheader("💰 Daily Earnings")
        earnings = stats.get("daily_earnings", [])
        if earnings:
            st.bar_chart(pd.DataFrame({"Earnings (MSG)": earnings}), use_container_width=True, height=300)

        st.subheader("📊 Cumulative Earnings")
        if earnings:
            cumulative = pd.Series(earnings).cumsum()
            st.area_chart(pd.DataFrame({"Cumulative (MSG)": cumulative}), use_container_width=True, height=300)

    if st.session_state.show_anomalies and requests:
        anomalies = []
        mean_req = np.mean(requests)
        std_req = np.std(requests)
        if std_req > 0:
            for i, val in enumerate(requests):
                if abs(val - mean_req) > 2 * std_req:
                    anomalies.append(i)
        if anomalies:
            st.warning(f"🔍 Detected {len(anomalies)} anomalous days in request volume")

5.6 标签页 - 收益

def render_earnings_tab(stats: dict):
    col1, col2 = st.columns([2, 1])

    with col1:
        st.subheader("Earnings Breakdown")
        earnings = stats.get("daily_earnings", [])
        if earnings:
            df = pd.DataFrame({"Day": list(range(len(earnings))), "Earnings": earnings})
            fig = px.bar(df, x="Day", y="Earnings", title="Daily Earnings",
                         labels={"Earnings": "MSG"}, color="Earnings",
                         color_continuous_scale="Viridis")
            fig.update_layout(height=400)
            st.plotly_chart(fig, use_container_width=True)

    with col2:
        st.subheader("Statistics")
        if earnings:
            earnings_arr = np.array(earnings)
            for label, value in [
                ("Total", f"{earnings_arr.sum():.2f} MSG"),
                ("Average", f"{earnings_arr.mean():.2f} MSG"),
                ("Max Day", f"{earnings_arr.max():.2f} MSG"),
                ("Median", f"{np.median(earnings_arr):.2f} MSG"),
            ]:
                st.metric(label, value)

        st.subheader("Top Earning Days")
        if earnings:
            top_indices = np.argsort(earnings)[-5:][::-1]
            st.dataframe(pd.DataFrame({
                "Day": [f"Day {i}" for i in top_indices],
                "Earnings": [f"{earnings[i]:.2f} MSG" for i in top_indices],
            }), use_container_width=True, hide_index=True)

5.7 标签页 - 活动

def render_activity_tab(stats: dict):
    st.subheader("Recent Activity")

    activity_types = ["Request", "Response", "Error", "Reputation Change", "Payment", "Stake", "Unstake"]
    statuses = ["success", "success", "success", "warning", "error"]

    if "activity_log" not in st.session_state:
        activities = []
        now = datetime.now()
        for i in range(20):
            ts = now - timedelta(hours=random.randint(0, 48), minutes=random.randint(0, 59))
            act_type = random.choice(activity_types)
            status = random.choice(statuses)
            activities.append({
                "Time": ts.strftime("%H:%M:%S"), "Type": act_type,
                "Status": status.upper(),
                "Detail": f"{act_type} processed in {random.uniform(0.1, 3.0):.1f}s",
                "Gas": f"{random.randint(50000, 500000):,}",
            })
        st.session_state.activity_log = pd.DataFrame(activities)

    df = st.session_state.activity_log.copy()
    col1, col2, col3 = st.columns(3)
    with col1:
        st.metric("Total Events", len(df))
    with col2:
        st.metric("Successful", len(df[df["Status"] == "SUCCESS"]))
    with col3:
        st.metric("Failed", len(df[df["Status"] == "ERROR"]), delta_color="inverse")

    st.dataframe(df, use_container_width=True, hide_index=True)

    col_left, col_right = st.columns(2)
    with col_left:
        type_counts = df["Type"].value_counts()
        fig = px.pie(values=type_counts.values, names=type_counts.index, title="Activity by Type", hole=0.4)
        st.plotly_chart(fig, use_container_width=True)
    with col_right:
        status_counts = df["Status"].value_counts()
        fig = px.bar(x=status_counts.index, y=status_counts.values, title="Activity by Status",
                     labels={"x": "Status", "y": "Count"}, color=status_counts.index,
                     color_discrete_map={"SUCCESS": "#4CAF50", "WARNING": "#FF9800", "ERROR": "#F44336"})
        st.plotly_chart(fig, use_container_width=True)

5.8 标签页 - 声誉

def render_reputation_tab(stats: dict):
    st.subheader("🏆 Reputation History")
    reputation = stats.get("reputation_history", [])

    if reputation:
        df = pd.DataFrame(reputation)
        df["datetime"] = pd.to_datetime(df["timestamp"], unit="s")

        fig = go.Figure(go.Scatter(
            x=df["datetime"], y=df["score"], mode="lines+markers", name="Score",
            line={"color": "#9C27B0", "width": 2},
            marker={"size": 8, "color": df["change"].apply(lambda c: "#4CAF50" if c >= 0 else "#F44336")},
            text=df["reason"],
        ))
        fig.add_hrect(y0=0, y1=40, fillcolor="red", opacity=0.05, line_width=0, annotation_text="Risk Zone")
        fig.add_hrect(y0=40, y1=70, fillcolor="orange", opacity=0.05, line_width=0, annotation_text="Caution Zone")
        fig.add_hrect(y0=70, y1=100, fillcolor="green", opacity=0.05, line_width=0, annotation_text="Good Zone")
        fig.update_layout(height=500, yaxis_range=[0, 100])
        st.plotly_chart(fig, use_container_width=True)

        latest = df.iloc[-1]
        col1, col2, col3 = st.columns(3)
        with col1:
            st.metric("Current Score", f"{latest['score']:.0f}/100")
        with col2:
            st.metric("Latest Change", f"{'+' if latest['change'] >= 0 else ''}{latest['change']:.0f}")
        with col3:
            st.metric("Last Updated", latest["datetime"].strftime("%Y-%m-%d %H:%M"))
    else:
        st.info("No reputation history available.")

5.9 标签页 - 对比

def render_comparison_tab():
    st.subheader("🆚 Multi-Agent Comparison")
    agents = st.multiselect("Select agents to compare",
        ["msg1agent2u9x7qh0u3h", "msg1agent8m3k5p9w1q", "msg1agent4v6n2r8t0s", "msg1agent0p3l6k9w2r"],
        default=["msg1agent2u9x7qh0u3h", "msg1agent8m3k5p9w1q"])

    if not agents:
        st.info("Select at least two agents.")
        return

    data = {}
    for agent in agents:
        data[agent[:16]] = [random.randint(500, 5000), round(random.uniform(10, 500), 2),
                            round(random.uniform(90, 99.9), 1), round(random.uniform(0.1, 3.0), 2)]

    df = pd.DataFrame(data, index=["Requests (24h)", "Earnings (24h)", "Success Rate", "Avg Latency"]).T
    fig = px.imshow(df.values, x=df.columns, y=df.index, text_auto=".2f",
                    color_continuous_scale="RdYlGn", aspect="auto", title="Agent Comparison Heatmap")
    fig.update_layout(height=400)
    st.plotly_chart(fig, use_container_width=True)
    st.dataframe(df, use_container_width=True)

5.10 主渲染循环

async def fetch_data(agent_id: str, days: int) -> dict:
    now = time.time()
    if (st.session_state.last_fetch > 0
        and now - st.session_state.last_fetch < st.session_state.refresh_interval):
        return st.session_state.data_cache

    try:
        collector = get_data_collector()
        stats = await collector.collect_agent_stats(agent_id, days)
        market = await collector.collect_market_data()
        result = {**stats, "market": market}
        st.session_state.data_cache = result
        st.session_state.last_fetch = now
        return result
    except Exception as e:
        st.error(f"Failed to fetch data: {e}")
        return st.session_state.data_cache or {}


def main():
    create_session_state()
    agent_id, days = render_sidebar()

    st.markdown(f"<div class='report-header'><h1>🤖 Agent Monitor</h1>"
                f"<p>Agent: <code>{agent_id}</code> | Period: {days} days | "
                f"<span class='status-badge online'>● ONLINE</span></p></div>",
                unsafe_allow_html=True)

    with st.spinner("Fetching on-chain data..."):
        stats = asyncio.run(fetch_data(agent_id, days))

    if not stats or not stats.get("daily_requests"):
        st.warning("No data available. Please check the Agent ID.")
        st.stop()

    render_kpi_metrics(stats)
    st.markdown("---")

    tabs = st.tabs(["📈 Performance", "💰 Earnings", "📋 Activity", "🏆 Reputation", "🆚 Comparison"])
    with tabs[0]: render_performance_tab(stats)
    with tabs[1]: render_earnings_tab(stats)
    with tabs[2]: render_activity_tab(stats)
    with tabs[3]: render_reputation_tab(stats)
    with tabs[4]: render_comparison_tab()

    col1, col2, col3 = st.columns(3)
    with col1: st.caption(f"🕐 Last Refresh: {datetime.now().strftime('%H:%M:%S')}")
    with col2: st.caption(f"🔄 Auto-refresh: {'ON' if st.session_state.auto_refresh else 'OFF'} ({st.session_state.refresh_interval}s)")
    with col3: st.caption("🔗 MSG Chain Mainnet | msg-chain-1")

    if st.session_state.auto_refresh:
        time.sleep(st.session_state.refresh_interval)
        st.rerun()


if __name__ == "__main__":
    main()

6. 自动化报告生成

6.1 HTML 报告模板

from datetime import datetime, timedelta
from typing import Optional, List, Dict, Any
import json
import os
import asyncio
import logging

logger = logging.getLogger(__name__)


class HTMLReportTemplate:
    """HTML 报告模板生成器"""

    CSS_STYLE = """
    <style>
        * { margin: 0; padding: 0; box-sizing: border-box; }
        body {
            font-family: 'Segoe UI', -apple-system, sans-serif;
            background: #F5F7FA; color: #333; line-height: 1.6; padding: 20px;
        }
        .container { max-width: 900px; margin: 0 auto; }
        .header {
            background: linear-gradient(135deg, #2196F3, #1565C0);
            color: white; padding: 40px; border-radius: 12px; margin-bottom: 30px;
        }
        .header h1 { font-size: 28px; margin-bottom: 8px; }
        .header p { opacity: 0.9; font-size: 14px; }
        .summary-grid {
            display: grid; grid-template-columns: repeat(auto-fit, minmax(200px, 1fr));
            gap: 16px; margin-bottom: 30px;
        }
        .summary-card {
            background: white; padding: 24px; border-radius: 10px;
            box-shadow: 0 2px 8px rgba(0,0,0,0.06); text-align: center;
        }
        .summary-card .value { font-size: 32px; font-weight: 700; color: #2196F3; margin: 8px 0; }
        .summary-card .label { font-size: 13px; color: #888; text-transform: uppercase; }
        .section {
            background: white; padding: 30px; border-radius: 10px;
            box-shadow: 0 2px 8px rgba(0,0,0,0.06); margin-bottom: 24px;
        }
        .section h2 { font-size: 20px; color: #1565C0; margin-bottom: 16px;
                       padding-bottom: 8px; border-bottom: 2px solid #E3F2FD; }
        table { width: 100%; border-collapse: collapse; margin: 16px 0; }
        th, td { padding: 12px 16px; text-align: left; border-bottom: 1px solid #EEE; }
        th { background: #F5F7FA; font-weight: 600; font-size: 13px; text-transform: uppercase; }
        .chart-container { margin: 20px 0; text-align: center; }
        .chart-container img { max-width: 100%; border-radius: 8px; box-shadow: 0 2px 12px rgba(0,0,0,0.1); }
        .footer { text-align: center; padding: 30px; color: #999; font-size: 12px; }
        .alert { padding: 16px; border-radius: 8px; margin: 16px 0; }
        .alert.warning { background: #FFF3E0; border-left: 4px solid #FF9800; color: #E65100; }
        .alert.info { background: #E3F2FD; border-left: 4px solid #2196F3; color: #1565C0; }
    </style>
    """

    def build_header(self, agent_id: str, report_date: str, report_type: str = "Daily") -> str:
        return f"""
        <div class="header">
            <h1>🤖 {report_type} Report</h1>
            <p>Agent: <code>{agent_id[:20]}...</code> | Date: {report_date}</p>
            <p style="margin-top: 8px;">MSG Chain Mainnet | msg-chain-1</p>
        </div>
        """

    def build_summary_card(self, label: str, value: str, color: str = "#2196F3") -> str:
        return f"""
        <div class="summary-card">
            <div class="label">{label}</div>
            <div class="value" style="color: {color};">{value}</div>
        </div>
        """

    def build_summary_grid(self, cards: List[Dict[str, str]]) -> str:
        cards_html = "".join(self.build_summary_card(c["label"], c["value"], c.get("color", "#2196F3"))
                             for c in cards)
        return f'<div class="summary-grid">{cards_html}</div>'

    def build_section(self, title: str, content: str) -> str:
        return f"""
        <div class="section">
            <h2>{title}</h2>
            {content}
        </div>
        """

    def build_table(self, headers: List[str], rows: List[List[str]]) -> str:
        thead = "<tr>" + "".join(f"<th>{h}</th>" for h in headers) + "</tr>"
        tbody = "".join("<tr>" + "".join(f"<td>{c}</td>" for c in row) + "</tr>" for row in rows)
        return f"<table><thead>{thead}</thead><tbody>{tbody}</tbody></table>"

    def build_chart(self, img_path: str, alt: str = "", caption: str = "") -> str:
        if not os.path.exists(img_path):
            return f'<p style="color: #999;">Chart not available: {alt}</p>'
        caption_html = f"<p style='color: #888; font-size: 12px;'>{caption}</p>" if caption else ""
        return f"""
        <div class="chart-container">
            <img src="{img_path}" alt="{alt}" />
            {caption_html}
        </div>
        """

    def build_footer(self) -> str:
        return f"""
        <div class="footer">
            <p>Generated by MSG Chain AI Agent | {datetime.utcnow().strftime('%Y-%m-%d %H:%M UTC')}</p>
            <p>Data sourced from MSG Chain Indexer</p>
        </div>
        """

6.2 报告生成器

class ReportGenerator:
    """自动报告生成器"""

    def __init__(self, collector, visualizer, template=None, output_dir="./reports"):
        self.collector = collector
        self.visualizer = visualizer
        self.template = template or HTMLReportTemplate()
        self.output_dir = output_dir
        os.makedirs(output_dir, exist_ok=True)

    async def generate_daily_report(self, agent_id: str, output_path=None) -> str:
        logger.info(f"Generating daily report for {agent_id}")

        stats = await self.collector.collect_agent_stats(agent_id, days=30)
        market = await self.collector.collect_market_data()

        report_date = datetime.utcnow().strftime("%Y-%m-%d")
        if not output_path:
            output_path = os.path.join(self.output_dir, f"report_{agent_id[:10]}_{report_date}.html")

        chart_dir = os.path.join(self.output_dir, "charts")
        os.makedirs(chart_dir, exist_ok=True)

        earn_chart = self.visualizer.plot_earnings_trend(stats, os.path.join(chart_dir, "earnings.png"))
        req_chart = self.visualizer.plot_request_volume(stats, os.path.join(chart_dir, "requests.png"))
        err_chart = self.visualizer.plot_error_rate(stats, os.path.join(chart_dir, "errors.png"))
        perf_chart = self.visualizer.plot_performance_panel(stats, os.path.join(chart_dir, "panel.png"))
        rep_chart = self.visualizer.plot_reputation_history(
            stats.get("reputation_history", []), os.path.join(chart_dir, "reputation.png"))

        earnings = stats.get("daily_earnings", [0])
        requests = stats.get("daily_requests", [0])
        errors = stats.get("daily_errors", [0])
        last_day_earnings = earnings[-1] if earnings else 0
        last_day_requests = requests[-1] if requests else 0
        last_day_errors = errors[-1] if errors else 0
        total_earnings = sum(earnings)
        avg_earnings = total_earnings / max(len(earnings), 1)
        error_rate = sum(errors) / max(sum(requests), 1) * 100

        summary_cards = [
            {"label": "Requests (24h)", "value": f"{last_day_requests:,}", "color": "#2196F3"},
            {"label": "Earnings (24h)", "value": f"{last_day_earnings:.2f} MSG", "color": "#4CAF50"},
            {"label": "Errors (24h)", "value": str(last_day_errors), "color": "#F44336"},
            {"label": "Avg Earnings/Day", "value": f"{avg_earnings:.2f} MSG", "color": "#FF9800"},
            {"label": "30-Day Total", "value": f"{total_earnings:.2f} MSG", "color": "#9C27B0"},
            {"label": "Error Rate", "value": f"{error_rate:.1f}%",
             "color": "#F44336" if error_rate > 5 else "#4CAF50"},
        ]

        html_parts = [
            "<!DOCTYPE html>", '<html lang="zh-CN">', "<head>",
            '<meta charset="UTF-8">', '<meta name="viewport" content="width=device-width, initial-scale=1.0">',
            f"<title>Daily Report - {agent_id[:10]}...</title>",
            self.template.CSS_STYLE, "</head>", "<body>", '<div class="container">',
            self.template.build_header(agent_id, report_date, "Daily"),
            self.template.build_summary_grid(summary_cards),
        ]

        html_parts.append(self.template.build_section("Performance Overview",
                            self.template.build_chart(perf_chart, "Performance Panel")))
        html_parts.append(self.template.build_section("Earnings Analysis",
                            self.template.build_chart(earn_chart, "Earnings Trend")))
        html_parts.append(self.template.build_section("Request Analysis",
                            self.template.build_chart(req_chart, "Request Volume")))
        html_parts.append(self.template.build_section("Error Analysis",
                            self.template.build_chart(err_chart, "Error Rate")))
        html_parts.append(self.template.build_section("Reputation",
                            self.template.build_chart(rep_chart, "Reputation History")))

        top_agents = market.get("top_agents", [])
        if top_agents:
            top_table = self.template.build_table(["#", "Agent", "Earnings"],
                [[str(i+1), agent.get("name", agent.get("id", f"Agent {i+1}"))[:20],
                  f"{float(agent.get('earnings', 0)):.2f} MSG"]
                 for i, agent in enumerate(top_agents[:5])])
            html_parts.append(self.template.build_section("Ecosystem Overview", top_table))

        html_parts.append(self.template.build_footer())
        html_parts.append("</div></body></html>")

        html = "\n".join(html_parts)
        with open(output_path, "w", encoding="utf-8") as f:
            f.write(html)

        logger.info(f"Report saved to {output_path}")
        return output_path

6.3 IPFS 存储

class IPFSStorage:
    """将报告上传到 IPFS"""

    def __init__(self, api_url: str = "http://localhost:5001/api/v0",
                 gateway_url: str = "https://ipfs.msgchain.org/ipfs"):
        self.api_url = api_url
        self.gateway_url = gateway_url

    async def add_file(self, file_path: str) -> str:
        import aiohttp
        if not os.path.exists(file_path):
            raise FileNotFoundError(f"File not found: {file_path}")
        async with aiohttp.ClientSession() as session:
            with open(file_path, "rb") as f:
                data = aiohttp.FormData()
                data.add_field("file", f, filename=os.path.basename(file_path))
                async with session.post(f"{self.api_url}/add", data=data) as resp:
                    result = await resp.json()
                    return result["Hash"]

    async def add_html(self, html_content: str) -> str:
        import aiohttp
        async with aiohttp.ClientSession() as session:
            data = aiohttp.FormData()
            data.add_field("file", html_content.encode("utf-8"), filename="report.html")
            async with session.post(f"{self.api_url}/add", data=data) as resp:
                result = await resp.json()
                return result["Hash"]

    def get_url(self, cid: str) -> str:
        return f"{self.gateway_url}/{cid}"

    async def pin(self, cid: str) -> bool:
        import aiohttp
        async with aiohttp.ClientSession() as session:
            async with session.post(f"{self.api_url}/pin/add?arg={cid}") as resp:
                return resp.status == 200

6.4 定时调度器

class ReportScheduler:
    """报告定时调度器"""

    def __init__(self, generator, distributor, ipfs):
        self.generator = generator
        self.distributor = distributor
        self.ipfs = ipfs
        self._tasks = {}

    async def generate_and_distribute(self, agent_id: str, channels=None) -> str:
        channels = channels or ["ipfs"]
        html_path = await self.generator.generate_daily_report(agent_id)
        cid = await self.ipfs.add_file(html_path)
        await self.ipfs.pin(cid)
        for channel in channels:
            await self.distributor.distribute_report(cid, [channel])
        return cid

    async def schedule_recurring(self, agent_id: str, interval_hours: int = 24, channels=None):
        channels = channels or ["ipfs"]
        while True:
            cid = await self.generate_and_distribute(agent_id, channels)
            logger.info(f"Report generated: {self.ipfs.get_url(cid)}")
            await asyncio.sleep(interval_hours * 3600)

    async def start_agent_schedule(self, agent_id: str, mode: str = "daily", channels=None):
        if mode == "daily":
            task = asyncio.create_task(self.schedule_recurring(agent_id, 24, channels))
        elif mode == "hourly":
            task = asyncio.create_task(self.schedule_recurring(agent_id, 1, channels))
        elif mode == "weekly":
            task = asyncio.create_task(self.schedule_recurring(agent_id, 168, channels))
        else:
            raise ValueError(f"Unknown mode: {mode}")
        self._tasks[agent_id] = task

    def stop_agent_schedule(self, agent_id: str):
        if agent_id in self._tasks:
            self._tasks[agent_id].cancel()
            del self._tasks[agent_id]

6.5 告警规则引擎

class AlertRule:
    """告警规则定义"""

    def __init__(self, name: str, metric: str, condition: str, threshold: float,
                 severity: str = "warning", message_template: str = None):
        self.name = name
        self.metric = metric
        self.condition = condition
        self.threshold = threshold
        self.severity = severity
        self.message_template = message_template or f"Alert: {metric} {condition} {threshold}"

    def evaluate(self, value: float) -> bool:
        if self.condition == ">": return value > self.threshold
        elif self.condition == "<": return value < self.threshold
        elif self.condition == ">=": return value >= self.threshold
        elif self.condition == "<=": return value <= self.threshold
        elif self.condition == "==": return value == self.threshold
        return False


class AlertEngine:
    """告警引擎"""

    DEFAULT_RULES = [
        AlertRule("High Error Rate", "error_rate", ">", 10.0, severity="critical",
                  message_template="Critical: Error rate {value:.1f}% exceeds 10%"),
        AlertRule("Elevated Error Rate", "error_rate", ">", 5.0, severity="warning",
                  message_template="Warning: Error rate {value:.1f}% exceeds 5%"),
        AlertRule("Low Reputation", "reputation", "<", 40, severity="critical",
                  message_template="Critical: Reputation score {value:.0f} below 40"),
        AlertRule("No Requests", "daily_requests", "==", 0, severity="warning",
                  message_template="Warning: No requests processed today"),
        AlertRule("High Latency", "avg_latency", ">", 5.0, severity="warning",
                  message_template="Warning: Avg latency {value:.1f}s exceeds 5s"),
    ]

    def __init__(self, rules: list = None):
        self.rules = rules or self.DEFAULT_RULES

    def check_stats(self, stats: dict) -> list:
        alerts = []
        requests = stats.get("daily_requests", [0])
        errors = stats.get("daily_errors", [0])
        reputation = stats.get("reputation_history", [])

        current_error_rate = (errors[-1] / requests[-1] * 100 if requests[-1] > 0 else 0) if errors and requests else 0
        current_reputation = reputation[-1]["score"] if reputation else 100
        current_latency = stats.get("latency_history", [0])[-1]

        metrics = {
            "error_rate": current_error_rate,
            "reputation": current_reputation,
            "daily_requests": requests[-1] if requests else 0,
            "avg_latency": current_latency,
        }

        for rule in self.rules:
            if rule.metric in metrics:
                value = metrics[rule.metric]
                if rule.evaluate(value):
                    alerts.append({
                        "rule": rule.name, "severity": rule.severity,
                        "message": rule.message_template.format(value=value),
                        "value": value, "threshold": rule.threshold,
                        "timestamp": datetime.utcnow().isoformat(),
                    })
        return alerts

7. 报告分发

7.1 分发器基类

from abc import ABC, abstractmethod
from typing import List, Optional, Dict, Any
import aiohttp
import smtplib
from email.mime.text import MIMEText
from email.mime.multipart import MIMEMultipart
import logging

logger = logging.getLogger(__name__)


class ReportChannel(ABC):
    """报告分发渠道基类"""

    @abstractmethod
    async def send(self, report_url: str, **kwargs) -> bool:
        ...


class ReportDistributor:
    """多通道分发器"""

    def __init__(self, ipfs: IPFSStorage, owner_email: str = None,
                 discord_webhook: str = None, telegram_bot_token: str = None,
                 telegram_chat_id: str = None):
        self.ipfs = ipfs
        self.owner_email = owner_email
        self.discord_webhook = discord_webhook
        self.telegram_bot_token = telegram_bot_token
        self.telegram_chat_id = telegram_chat_id
        self.peers = []
        self.channels = {}

    def register_channel(self, name: str, channel: ReportChannel):
        self.channels[name] = channel

7.2 邮件分发

class EmailChannel(ReportChannel):
    """SMTP 邮件分发"""

    def __init__(self, smtp_host: str = "smtp.gmail.com", smtp_port: int = 587,
                 username: str = None, password: str = None,
                 from_addr: str = None, to_addr: str = None, use_tls: bool = True):
        self.smtp_host = smtp_host
        self.smtp_port = smtp_port
        self.username = username
        self.password = password
        self.from_addr = from_addr or username
        self.to_addr = to_addr
        self.use_tls = use_tls

    async def send(self, report_url: str, subject: str = "AI Agent Report",
                   body: str = None, **kwargs) -> bool:
        if not body:
            body = f"Your AI Agent report is ready.\n\nView Report: {report_url}\n\nGenerated by MSG Chain AI Agent"
        to = kwargs.get("to", self.to_addr)
        if not to:
            logger.error("No recipient address configured")
            return False
        try:
            msg = MIMEMultipart("alternative")
            msg["Subject"] = subject
            msg["From"] = self.from_addr
            msg["To"] = to
            text_part = MIMEText(body, "plain", "utf-8")
            html_part = MIMEText(f"<p>{body.replace(chr(10), '<br>')}</p><p><a href='{report_url}'>Open Report</a></p>",
                                "html", "utf-8")
            msg.attach(text_part)
            msg.attach(html_part)
            with smtplib.SMTP(self.smtp_host, self.smtp_port) as server:
                if self.use_tls:
                    server.starttls()
                if self.username and self.password:
                    server.login(self.username, self.password)
                server.send_message(msg)
            logger.info(f"Email sent to {to}")
            return True
        except Exception as e:
            logger.error(f"Failed to send email: {e}")
            return False

7.3 Discord 分发

class DiscordChannel(ReportChannel):
    """Discord Webhook 分发"""

    def __init__(self, webhook_url: str = None):
        self.webhook_url = webhook_url

    async def send(self, report_url: str, message: str = None, **kwargs) -> bool:
        webhook = kwargs.get("webhook_url", self.webhook_url)
        if not webhook:
            logger.error("No Discord webhook URL configured")
            return False
        if not message:
            message = "AI Agent Report\nYour daily report is ready!"
        payload = {
            "content": message,
            "embeds": [{
                "title": "View Full Report",
                "url": report_url,
                "color": 33758,
            }],
            "username": "Agent Reporter",
        }
        try:
            async with aiohttp.ClientSession() as session:
                async with session.post(webhook, json=payload) as resp:
                    if resp.status != 204:
                        return False
            return True
        except Exception as e:
            logger.error(f"Discord send error: {e}")
            return False

7.4 Telegram 分发

class TelegramChannel(ReportChannel):
    """Telegram Bot 分发"""

    API = "https://api.telegram.org/bot"

    def __init__(self, bot_token: str = None, chat_id: str = None):
        self.bot_token = bot_token
        self.chat_id = chat_id

    async def send(self, report_url: str, message: str = None, **kwargs) -> bool:
        token = kwargs.get("bot_token", self.bot_token)
        chat_id = kwargs.get("chat_id", self.chat_id)
        if not token or not chat_id:
            logger.error("Telegram not configured")
            return False
        if not message:
            message = f"AI Agent Report\n{report_url}"
        payload = {"chat_id": chat_id, "text": message, "disable_web_page_preview": False}
        try:
            async with aiohttp.ClientSession() as session:
                async with session.post(f"{self.API}{token}/sendMessage", json=payload) as resp:
                    data = await resp.json()
                    return data.get("ok", False)
        except Exception as e:
            logger.error(f"Telegram error: {e}")
            return False

7.5 A2A Agent 间通信

class A2AChannel(ReportChannel):
    """Agent-to-Agent 协议分发"""

    def __init__(self, peers=None):
        self.peers = peers or []

    def add_peer(self, agent_id: str, endpoint: str):
        self.peers.append({"id": agent_id, "endpoint": endpoint})

    async def send(self, report_url: str, message: str = None, **kwargs) -> bool:
        peers = kwargs.get("peers", self.peers)
        if not peers:
            return True
        payload = {
            "type": "report", "url": report_url,
            "period": kwargs.get("period", "daily"),
            "generated_at": datetime.utcnow().isoformat(),
            "summary": message or "Daily report available",
        }
        success = True
        async with aiohttp.ClientSession() as session:
            for peer in peers:
                try:
                    async with session.post(f"{peer['endpoint']}/a2a/report", json=payload,
                                            timeout=aiohttp.ClientTimeout(total=10)) as resp:
                        if resp.status != 200:
                            success = False
                except Exception as e:
                    logger.error(f"A2A to {peer['id']} failed: {e}")
                    success = False
        return success

7.6 分发协调器

class DistributionCoordinator:
    """分发协调器"""

    def __init__(self, distributor: ReportDistributor):
        self.distributor = distributor

    async def distribute_to_all(self, cid: str, channels=None, custom_message: str = None) -> dict:
        report_url = self.distributor.ipfs.get_url(cid)
        channels = channels or ["email", "discord", "telegram", "a2a"]
        results = {}
        for channel_name in channels:
            channel = self.distributor.channels.get(channel_name)
            if channel:
                try:
                    results[channel_name] = await channel.send(report_url, message=custom_message)
                except Exception as e:
                    logger.error(f"Channel {channel_name} failed: {e}")
                    results[channel_name] = False
            else:
                results[channel_name] = False
        return results

    async def distribute_with_priority(self, cid: str, alert_level: str = "info") -> dict:
        if alert_level == "critical":
            channels = ["telegram", "discord", "email", "a2a"]
        elif alert_level == "warning":
            channels = ["email", "discord"]
        else:
            channels = ["email"]
        return await self.distribute_to_all(cid, channels)

7.7 Grafana/Prometheus 数据源

class GrafanaDatasource:
    """将链上数据暴露为 Prometheus 格式"""

    def __init__(self, indexer, grafana_url: str = "http://localhost:3000"):
        self.indexer = indexer
        self.grafana_url = grafana_url

    def generate_prometheus_metrics(self, stats: dict, agent_id: str) -> str:
        lines = ['# HELP msg_agent_requests_total Total requests',
                 '# TYPE msg_agent_requests_total counter']
        for i, val in enumerate(stats.get("daily_requests", [])):
            lines.append(f'msg_agent_requests_total{{agent_id="{agent_id[:16]}",day="{i}"}} {val}')

        lines.extend(['# HELP msg_agent_earnings_total Total earnings',
                      '# TYPE msg_agent_earnings_total counter'])
        cum = 0
        for i, val in enumerate(stats.get("daily_earnings", [])):
            cum += val
            lines.append(f'msg_agent_earnings_total{{agent_id="{agent_id[:16]}",day="{i}"}} {cum}')

        err_rate = (stats["daily_errors"][-1] / stats["daily_requests"][-1] * 100
                    if stats.get("daily_requests") and stats["daily_requests"][-1] > 0 else 0)
        lines.extend(['# HELP msg_agent_error_rate Error rate', '# TYPE msg_agent_error_rate gauge'])
        lines.append(f'msg_agent_error_rate{{agent_id="{agent_id[:16]}"}} {err_rate:.2f}')

        latency = stats.get("latency_history", [0])
        lines.extend(['# HELP msg_agent_latency_seconds Latency', '# TYPE msg_agent_latency_seconds gauge'])
        lines.append(f'msg_agent_latency_seconds{{agent_id="{agent_id[:16]}",quantile="avg"}} {latency[-1]:.2f}')
        return "\n".join(lines)

8. 完整示例:监控与报告 Agent

#!/usr/bin/env python3
"""""""
MSG Chain AI Agent - 完整的监控与报告 Agent
集成了数据收集、可视化、仪表盘、报告生成与分发功能。

环境变量:
  MSGCHAIN_INDEXER_URL     Indexer API 地址
  MSGCHAIN_IPFS_API        IPFS API 地址
  MSGCHAIN_DISCORD_WEBHOOK Discord Webhook URL
  MSGCHAIN_TELEGRAM_TOKEN  Telegram Bot Token
  MSGCHAIN_TELEGRAM_CHAT   Telegram Chat ID
  MSGCHAIN_EMAIL_TO        报告接收邮箱
  MSGCHAIN_SMTP_HOST       SMTP 服务器
  MSGCHAIN_SMTP_USER       SMTP 用户名
  MSGCHAIN_SMTP_PASS       SMTP 密码
""""""

import os
import sys
import json
import asyncio
import logging
from datetime import datetime, timedelta
from typing import Optional, List, Dict, Any

logging.basicConfig(level=logging.INFO, format="%(asctime)s [%(levelname)s] %(name)s: %(message)s")
logger = logging.getLogger("agent-monitor")


class MonitoringAgent:
    """完整的监控与报告 AI Agent"""

    def __init__(self, agent_id: str):
        self.agent_id = agent_id
        self.config = self._load_config()
        self.indexer = IndexerClient(self.config.get("indexer_url", "https://indexer.msgchain.org/api/v1"))
        self.collector = DataCollector(self.indexer)
        self.visualizer = AgentVisualizer(output_dir=self.config.get("chart_dir", "./charts"))
        self.template = HTMLReportTemplate()
        self.ipfs = IPFSStorage(api_url=self.config.get("ipfs_api", "http://localhost:5001/api/v0"))
        self.alert_engine = AlertEngine()
        self.report_generator = ReportGenerator(self.collector, self.visualizer, self.template,
                                                output_dir=self.config.get("report_dir", "./reports"))
        self.distributor = self._setup_distributor()
        self.scheduler = ReportScheduler(self.report_generator, self.distributor, self.ipfs)
        self._running = False

    def _load_config(self) -> dict:
        config = {
            "indexer_url": os.getenv("MSGCHAIN_INDEXER_URL"),
            "ipfs_api": os.getenv("MSGCHAIN_IPFS_API"),
            "discord_webhook": os.getenv("MSGCHAIN_DISCORD_WEBHOOK"),
            "telegram_token": os.getenv("MSGCHAIN_TELEGRAM_TOKEN"),
            "telegram_chat": os.getenv("MSGCHAIN_TELEGRAM_CHAT"),
            "email_to": os.getenv("MSGCHAIN_EMAIL_TO"),
            "smtp_host": os.getenv("MSGCHAIN_SMTP_HOST", "smtp.gmail.com"),
            "smtp_user": os.getenv("MSGCHAIN_SMTP_USER"),
            "smtp_pass": os.getenv("MSGCHAIN_SMTP_PASS"),
            "chart_dir": "./charts", "report_dir": "./reports", "check_interval": 3600,
        }
        config_path = os.path.join(os.path.dirname(__file__), "agent_config.json")
        if os.path.exists(config_path):
            with open(config_path) as f:
                config.update(json.load(f))
        return config

    def _setup_distributor(self):
        distributor = ReportDistributor(
            self.ipfs, owner_email=self.config.get("email_to"),
            discord_webhook=self.config.get("discord_webhook"),
            telegram_bot_token=self.config.get("telegram_token"),
            telegram_chat_id=self.config.get("telegram_chat"))
        if self.config.get("email_to") and self.config.get("smtp_user"):
            distributor.register_channel("email", EmailChannel(
                smtp_host=self.config["smtp_host"],
                username=self.config["smtp_user"],
                password=self.config["smtp_pass"],
                to_addr=self.config["email_to"]))
        if self.config.get("discord_webhook"):
            distributor.register_channel("discord", DiscordChannel(self.config["discord_webhook"]))
        if self.config.get("telegram_token") and self.config.get("telegram_chat"):
            distributor.register_channel("telegram", TelegramChannel(
                self.config["telegram_token"], self.config["telegram_chat"]))
        distributor.register_channel("a2a", A2AChannel())
        return distributor

    async def run_once(self) -> dict:
        """单次运行:采集 -> 检查 -> 可视化 -> 报告 -> 分发"""
        logger.info(f"Starting monitoring cycle for {self.agent_id}")
        stats = await self.collector.collect_agent_stats(self.agent_id, 30)
        logger.info("Data collected")

        alerts = self.alert_engine.check_stats(stats)
        if alerts:
            logger.warning(f"Alerts triggered: {len(alerts)}")

        html_path = await self.report_generator.generate_daily_report(self.agent_id)
        logger.info(f"Report generated: {html_path}")

        cid = await self.ipfs.add_file(html_path)
        await self.ipfs.pin(cid)
        report_url = self.ipfs.get_url(cid)
        logger.info(f"Report uploaded to IPFS: {report_url}")

        alert_level = "info"
        if any(a["severity"] == "critical" for a in alerts):
            alert_level = "critical"
        elif any(a["severity"] == "warning" for a in alerts):
            alert_level = "warning"

        coordinator = DistributionCoordinator(self.distributor)
        if alerts:
            results = await coordinator.distribute_with_priority(cid, alert_level)
        else:
            results = await coordinator.distribute_to_all(cid, ["email"])

        return {
            "agent_id": self.agent_id, "timestamp": datetime.utcnow().isoformat(),
            "alerts": alerts, "report_url": report_url,
            "distribution_results": results,
            "stats_summary": {
                "requests_24h": stats["daily_requests"][-1] if stats.get("daily_requests") else 0,
                "earnings_24h": stats["daily_earnings"][-1] if stats.get("daily_earnings") else 0,
                "errors_24h": stats["daily_errors"][-1] if stats.get("daily_errors") else 0,
            },
        }

    async def run_loop(self, interval_seconds: int = None):
        interval = interval_seconds or self.config.get("check_interval", 3600)
        self._running = True
        logger.info(f"Starting monitoring loop, interval={interval}s")
        while self._running:
            try:
                result = await self.run_once()
                logger.info(f"Cycle complete. Report: {result['report_url']}")
            except Exception as e:
                logger.error(f"Monitoring cycle failed: {e}")
            await asyncio.sleep(interval)

    def stop(self):
        self._running = False
        logger.info("Monitoring agent stopped")


async def main():
    agent_id = os.getenv("MSGCHAIN_AGENT_ID", "msg1agent2u9x7qh0u3h")
    agent = MonitoringAgent(agent_id)
    logger.info(f"Starting MonitoringAgent: {agent_id}")
    logger.info(f"Report channels: {list(agent.distributor.channels.keys())}")
    await agent.run_once()


if __name__ == "__main__":
    asyncio.run(main())

9. 附录:常见问题

9.1 如何获取 Agent ID?

Agent ID 是 Agent 在 MSG Chain 上的合约地址或钱包地址,以 msg1 开头。

9.2 Streamlit 仪表盘部署

pip install streamlit pandas numpy plotly
streamlit run agent_monitor.py --server.port 8501

9.3 常见 Matplotlib 问题

问题 解决方案
中文乱码 安装中文字体: plt.rcParams['font.sans-serif'] = ['SimHei']
无头环境报错 matplotlib.use('Agg') 必须在 import pyplot 之前设置
图片模糊 增加 figure.dpi 和 savefig.dpi
内存泄漏 每次保存后调用 plt.close(fig)

9.4 数据采集调优

9.5 MSG Chain 主网信息

参数 值
Chain ID msg-chain-1
bech32 前缀 msg
原生代币 MSG (umsg)
精度 18 decimals
Indexer API https://indexer.msgchain.org/api/v1
RPC https://rpc.msgchain.org
REST/LCD https://rest.msgchain.org
WebSocket wss://rpc.msgchain.org/websocket
IPFS Gateway https://ipfs.msgchain.org/ipfs
浏览器 https://explorer.msgchain.org

9.6 推荐的工具链组合

场景 推荐组合
个人开发调试 Streamlit + DataCollector + Plotly
生产环境监控 Grafana + Prometheus + AlertManager
合规报告 Matplotlib + HTML Report + IPFS + Email
社区分享 Plotly + Streamlit Cloud + Twitter/Discord
多 Agent 对比 Streamlit Multi-page + A2A 数据聚合

本文档由 MSG Chain AI Agent 数据可视化指南生成


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