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. 概述
1.1 为什么 Agent 需要数据可视化
AI Agent 在 MSG Chain 上运行时产生大量链上数据:请求日志、收益记录、声誉评分变化、执行错误等。原始数据难以洞察趋势,可视化将数字转化为直观图形,帮助开发者、持有者和利益相关者快速理解 Agent 状态。
一个具备可视化能力的 Agent 可以:
- 自我监控:实时追踪自身指标,异常时自动告警
- 透明可信:将运行数据以图表形式公开展示,增强用户信任
- 决策支持:通过趋势分析优化 Gas 策略、调用频率和任务分配
- 合规报告:自动生成运营报告供社区或 DAO 审阅
1.2 工具栈选择
| 工具 | 用途 | 优势 |
|---|---|---|
| Matplotlib | 静态图表生成 | 成熟稳定,适合 PDF/PNG 导出 |
| Plotly | 交互式图表 | 鼠标悬停、缩放、点击交互 |
| Grafana | 生产仪表盘 | 实时数据源、告警规则、团队共享 |
| Streamlit | 快速应用搭建 | 纯 Python、组件丰富、部署简单 |
1.3 报告类型
- 性能报告:TPS、响应时间、Gas 消耗、成功率
- 收益报告:日/周/月收益、累计收益、收益来源分析
- 声誉报告:声誉变化曲线、评级事件时间线
- 活动报告:调用频率、活跃时段、请求分布
- 综合运营报告:以上所有维度的汇总
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 可以通过以下方式获取链上数据:
- RPC 节点:直接查询链状态 (
msgd query ...) - Indexer API:专用索引服务,提供历史聚合数据
- 事件日志:订阅链上事件,实时接收更新
- 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 开头。
- 合约部署输出:
msgd tx wasm store ...返回的合约地址 - 区块浏览器:访问 https://explorer.msgchain.org 搜索 Agent 名称
- Indexer API:
GET /api/v1/agents列出所有注册 Agent
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 数据采集调优
- 缓存策略:DataCollector 内置 cache 字典,避免重复请求
- 并发控制:使用 asyncio.gather 并行获取多 Agent 数据
- 限速:IndexerClient 可加 asyncio.Semaphore 控制并发数
- 重试机制:建议使用 tenacity 库添加指数退避重试
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/
