dApp Docs/链上数据分析与仪表盘构建指南
Development reference. Not independently verified for production.

MSG Chain 链上数据分析与仪表盘构建完全指南

数据来源:MSG Chain 代码库核实

主网状态: No-Go — 当前 MSGChain 主网裁决为 No-Go,以下内容反映代码实际状态,不代表生产可用。


目录

  1. 概述
  2. 数据索引管道
  3. 数据库模型
  4. SQL 分析查询
  5. 仪表盘 API
  6. React 前端仪表盘
  7. 部署与运维

1. 概述

1.1 为什么需要链上数据分析

MSG Chain 是一条基于 Cosmos SDK + CosmWasm 的应用链,专为 AI Agent 经济设计。链上运行着 Agent 注册、支付结算、DID 管理、治理规则和微支付通道等核心合约。随着生态发展,以下问题日益突出:

链上数据分析仪表盘(类似 Dune Analytics、Etherscan Analytics)正是为解决这些问题而生。它将原始区块链数据转化为可交互的可视化洞察。

1.2 MSG Chain 数据源

MSG Chain 提供多层级数据访问接口:

接口 协议 端点示例 用途 局限性
Tendermint RPC HTTP/WebSocket http://localhost:26657 区块、交易、验证人查询 无聚合能力
Cosmos REST HTTP http://localhost:1317 Cosmos SDK 标准查询 分页有限
gRPC gRPC localhost:9090 高性能批量查询 工具链要求高
Tendermint WS WebSocket ws://localhost:26657/websocket 实时事件订阅 无持久化
LCD HTTP http://localhost:1317/cosmos/tx/v1beta1 交易查询 单次查询限制

1.3 架构总览

+----------------------------------------------------+
|                    MSG Chain                         |
|  +-----------------------------------------------+  |
|  |  Tendermint RPC (26657) | CosmWasm Events     |  |
|  +-----------------------------------------------+  |
+----------------------------------------------------+
                        |
                        v
+----------------------------------------------------+
|          Analytics Indexer (Python async)           |
|  Block Fetch -> Tx Parser -> Event Xt -> DB Writer  |
|  Backfill -> State Snapshots -> Aggregator          |
+----------------------------------------------------+
                        |
                        v
+----------------------------------------------------+
|           PostgreSQL 15+ / TimescaleDB              |
|  blocks | transactions | events | materialized views|
+----------------------------------------------------+
                        |
                        v
+----------------------------------------------------+
|     FastAPI + Strawberry GraphQL + Redis Cache      |
+----------------------------------------------------+
                        |
                        v
+----------------------------------------------------+
|          React 18 + TypeScript + recharts           |
|  Overview | Charts | Tables | Agent Detail View     |
+----------------------------------------------------+

1.4 与传统区块浏览器的区别

特性 传统区块浏览器 数据分析仪表盘
查询粒度 单一条目 聚合统计
时间范围 实时快照 历史趋势 (7d/30d/90d)
数据模型 原始链上数据 清洗聚合建模后
可视化 表格为主 图表 + 表格 + 看板
目标用户 日常查询用户 分析师、产品经理
查询语言 界面操作 SQL / GraphQL

1.5 技术栈选择

组件 选型 理由
索引器 Python 3.12+ (asyncio + aiohttp + asyncpg) 生态丰富、AI/ML 集成方便
数据库 PostgreSQL 15+ + TimescaleDB 时序数据优化、成熟稳定
API 层 FastAPI + Strawberry GraphQL 高性能、自动文档、类型安全
缓存 Redis 7+ 热点数据加速、会话管理
前端 React 18 + TypeScript + recharts 组件化、图表丰富
容器化 Docker Compose 一键部署、环境一致
调度 APScheduler 定时聚合任务

2. 数据索引管道

2.1 概述

数据索引管道是仪表盘的基础设施,负责从 MSG Chain 获取原始数据,解析、清洗、规范化后存入分析数据库。本章使用 Python asyncio 实现完整的索引器,包含 RPC 客户端、事件解析器、区块索引器、历史回填器、状态快照器和数据聚合器。

整体数据流如下:

Tendermint RPC / REST API
        |
        v
  MsgChainIndexer (asyncio)
   - 并行获取区块和交易
   - 解析 wasm 事件
   - 提取 sender/contract/action
        |
        v
  EventParser + TxParser
   - 标准化为 BlockData/TxData/EventData
   - 分类事件 (agent_payment/registry/did/...)
        |
        v
  Database (asyncpg)
   - 事务写入 blocks + transactions + events
   - 更新 indexer_state
   - 批量提交优化

索引器支持两种运行模式:

2.2 RPC 客户端

# indexer/rpc_client.py
"""MSG Chain RPC 客户端 — 支持同步和异步模式"""

import asyncio
import time
from typing import Any, Optional
from dataclasses import dataclass

import aiohttp


@dataclass
class ChainStatus:
    """链状态"""
    latest_block_height: int
    latest_block_hash: str
    chain_id: str
    node_info: dict
    validator_count: int


class RPCClient:
    """Cosmos/Tendermint RPC 客户端"""

    def __init__(
        self,
        rpc_endpoint: str = "http://localhost:26657",
        rest_endpoint: str = "http://localhost:1317",
        timeout: int = 30,
    ):
        self.rpc_endpoint = rpc_endpoint.rstrip("/")
        self.rest_endpoint = rest_endpoint.rstrip("/")
        self.timeout = timeout
        self._session: Optional[aiohttp.ClientSession] = None

    async def get_session(self) -> aiohttp.ClientSession:
        if self._session is None or self._session.closed:
            self._session = aiohttp.ClientSession(
                timeout=aiohttp.ClientTimeout(total=self.timeout),
                headers={"Content-Type": "application/json"},
            )
        return self._session

    async def rpc_call(self, method: str, params: list[Any] | None = None) -> dict:
        """通用 JSON-RPC 调用"""
        session = await self.get_session()
        payload = {
            "jsonrpc": "2.0",
            "id": int(time.time() * 1000),
            "method": method,
            "params": params or [],
        }
        async with session.post(f"{self.rpc_endpoint}/", json=payload) as resp:
            resp.raise_for_status()
            data = await resp.json()
            if "error" in data and data["error"]:
                raise ValueError(f"RPC error: {data['error']}")
            return data["result"]

    async def get_status(self) -> ChainStatus:
        """获取链状态"""
        result = await self.rpc_call("status")
        sync_info = result["sync_info"]
        return ChainStatus(
            latest_block_height=int(sync_info["latest_block_height"]),
            latest_block_hash=sync_info["latest_block_hash"],
            chain_id=result["node_info"]["network"],
            node_info=result["node_info"],
            validator_count=int(result.get("validator_info", {}).get("count", 0)),
        )

    async def get_block(self, height: int | None = None) -> dict:
        """获取区块"""
        params = [str(height)] if height is not None else []
        return await self.rpc_call("block", params)

    async def get_block_results(self, height: int) -> dict:
        """获取区块执行结果(包含事件)"""
        return await self.rpc_call("block_results", [str(height)])

    async def get_txs_by_height(self, height: int) -> list[dict]:
        """获取指定区块的所有交易(REST API 分页)"""
        session = await self.get_session()
        txs: list[dict] = []
        offset = 0
        while True:
            url = (
                f"{self.rest_endpoint}/cosmos/tx/v1beta1/txs"
                f"?events=tx.height={height}&pagination.offset={offset}&pagination.limit=100"
            )
            async with session.get(url) as resp:
                resp.raise_for_status()
                data = await resp.json()
                tx_responses = data.get("tx_responses", [])
                txs.extend(tx_responses)
                if len(tx_responses) < 100:
                    break
                offset += 100
        return txs

    async def close(self):
        if self._session and not self._session.closed:
            await self._session.close()

2.3 数据模型定义

# indexer/models.py
"""索引数据模型 — 使用 dataclasses 类型安全"""

from dataclasses import dataclass, field
from datetime import datetime
from decimal import Decimal
from typing import Any


@dataclass
class BlockData:
    """区块数据"""
    height: int
    hash: str
    time: datetime
    proposer: str
    tx_count: int
    gas_used: int = 0
    gas_wanted: int = 0
    total_fee: Decimal = Decimal("0")
    app_hash: str = ""
    num_events: int = 0
    transactions: list["TxData"] = field(default_factory=list)


@dataclass
class TxData:
    """交易数据"""
    hash: str
    block_height: int
    index: int
    status: str  # SUCCESS / FAILED
    code: int = 0
    gas_used: int = 0
    gas_wanted: int = 0
    fee_amount: Decimal = Decimal("0")
    fee_denom: str = "umsg"
    memo: str = ""
    sender: str = ""
    contract: str = ""
    action: str = ""
    messages: list[dict] = field(default_factory=list)
    events: list["EventData"] = field(default_factory=list)
    raw_log: str = ""


@dataclass
class EventData:
    """合约事件数据"""
    tx_hash: str
    block_height: int
    event_type: str  # wasm / coin / transfer / message
    action: str
    contract_address: str = ""
    attributes: dict[str, str] = field(default_factory=dict)
    success: bool = True


@dataclass
class AggregatedStats:
    """聚合统计数据"""
    timestamp: datetime
    period: str  # hour / day / week / month
    total_blocks: int = 0
    total_txs: int = 0
    total_events: int = 0
    unique_senders: int = 0
    unique_contracts: int = 0
    avg_gas_price: Decimal = Decimal("0")
    avg_gas_used: int = 0
    total_fees: Decimal = Decimal("0")
    success_rate: float = 1.0
    active_agents: int = 0
    new_agents: int = 0
    new_sessions: int = 0
    session_volume: Decimal = Decimal("0")


@dataclass
class IndexerState:
    """索引器状态"""
    last_processed_height: int
    last_processed_hash: str = ""
    chain_id: str = "msg-chain-1"
    blocks_indexed: int = 0
    txs_indexed: int = 0
    events_indexed: int = 0

2.4 事件解析器

# indexer/event_parser.py
"""CosmWasm 事件解析器 — 原始 wasm 事件 → 结构化 EventData"""

import json
from decimal import Decimal
from typing import Any

from indexer.models import EventData, TxData


class EventParser:
    """事件解析器"""

    @classmethod
    def parse_raw_events(
        cls,
        raw_events: list[dict[str, Any]],
        tx_hash: str = "",
        block_height: int = 0,
        success: bool = True,
    ) -> list[EventData]:
        """解析原始 wasm 事件列表"""
        results: list[EventData] = []
        for evt in raw_events:
            evt_type = evt.get("type", "")
            if evt_type != "wasm":
                continue
            attributes = evt.get("attributes", [])
            attrs_dict = {}
            for attr in attributes:
                key = attr.get("key", "")
                value = attr.get("value", "")
                if key:
                    attrs_dict[key] = value
            action = attrs_dict.pop("action", "")
            if not action:
                continue
            contract_addr = attrs_dict.pop("_contract_address", "")
            event = EventData(
                tx_hash=tx_hash,
                block_height=block_height,
                event_type=evt_type,
                action=action,
                contract_address=contract_addr,
                attributes=attrs_dict,
                success=success,
            )
            results.append(event)
        return results

    @classmethod
    def get_event_category(cls, action: str) -> str:
        """事件分类"""
        payment = {"CreateSession", "FundSession", "ReleasePayment",
                   "DisputePayment", "CloseSession"}
        registry = {"RegisterAgent", "UpdateAgent", "DeregisterAgent", "SetStatus"}
        did = {"CreateDID", "UpdateDID", "DeactivateDID"}
        constitution = {"AddRule", "RemoveRule", "UpdateRule", "SetActive"}
        micropayment = {"OpenChannel", "Deposit", "Claim", "CloseChannel"}

        if action in payment:
            return "agent_payment"
        elif action in registry:
            return "agent_registry"
        elif action in did:
            return "did_registry"
        elif action in constitution:
            return "constitution"
        elif action in micropayment:
            return "micropayment"
        return "other"


class TxParser:
    """交易解析器"""

    @classmethod
    def extract_sender(cls, tx: dict) -> str:
        tx_body = tx.get("tx", {}).get("body", {})
        msgs = tx_body.get("messages", [])
        if msgs:
            sender = msgs[0].get("sender", "")
            if sender:
                return sender
        return ""

    @classmethod
    def extract_contract(cls, tx: dict) -> str:
        tx_body = tx.get("tx", {}).get("body", {})
        for msg in tx_body.get("messages", []):
            if msg.get("@type") == "/cosmwasm.wasm.v1.MsgExecuteContract":
                return msg.get("contract", "")
        return ""

    @classmethod
    def extract_action(cls, tx: dict) -> str:
        for evt in tx.get("events", []):
            if evt.get("type") == "wasm":
                for attr in evt.get("attributes", []):
                    if attr.get("key") == "action":
                        return attr.get("value", "")
        return ""

    @classmethod
    def extract_fee(cls, tx: dict) -> tuple[Decimal, str]:
        fee_info = (
            tx.get("tx", {}).get("auth_info", {}).get("fee", {})
        )
        amounts = fee_info.get("amount", [])
        if amounts:
            return (Decimal(amounts[0]["amount"]), amounts[0]["denom"])
        return (Decimal("0"), "umsg")

    @classmethod
    def parse_tx(cls, tx: dict, block_height: int, index: int) -> dict:
        """完全解析交易"""
        tx_result = tx.get("tx_response", tx)
        return {
            "hash": tx_result.get("txhash", ""),
            "block_height": block_height,
            "index": index,
            "code": tx_result.get("code", 0),
            "status": "SUCCESS" if tx_result.get("code", 0) == 0 else "FAILED",
            "gas_used": int(tx_result.get("gas_used", 0)),
            "gas_wanted": int(tx_result.get("gas_wanted", 0)),
            "fee_amount": cls.extract_fee(tx)[0],
            "fee_denom": cls.extract_fee(tx)[1],
            "memo": tx.get("tx", {}).get("body", {}).get("memo", ""),
            "sender": cls.extract_sender(tx),
            "contract": cls.extract_contract(tx),
            "action": cls.extract_action(tx),
        }

2.5 区块索引器

# indexer/block_indexer.py
"""MSG Chain 数据索引器 — 核心索引逻辑"""

import asyncio
import logging
from datetime import datetime
from typing import Callable, Coroutine

from indexer.rpc_client import RPCClient
from indexer.event_parser import EventParser, TxParser
from indexer.models import BlockData, TxData, EventData

logger = logging.getLogger(__name__)


class MsgChainIndexer:
    """MSG Chain 数据索引器

    支持两种模式:
    - realtime: 通过轮询监听新区块,实时索引
    - backfill: 从指定高度开始批量回填历史数据

    Usage:
        indexer = MsgChainIndexer(rpc_endpoint="http://localhost:26657")
        indexer.on_block(callback)
        await indexer.index_blocks(100, 200)
    """

    def __init__(
        self,
        rpc_endpoint: str = "http://localhost:26657",
        rest_endpoint: str = "http://localhost:1317",
        batch_size: int = 50,
        concurrency: int = 10,
    ):
        self.rpc = RPCClient(rpc_endpoint, rest_endpoint)
        self.batch_size = batch_size
        self.concurrency = concurrency
        self._on_block: list[Callable[[BlockData], Coroutine]] = []
        self._on_tx: list[Callable[[TxData], Coroutine]] = []
        self._on_event: list[Callable[[EventData], Coroutine]] = []
        self._running = False
        self._last_height = 0
        self._semaphore = asyncio.Semaphore(concurrency)

    def on_block(self, cb: Callable[[BlockData], Coroutine]):
        self._on_block.append(cb)
        return cb

    def on_tx(self, cb: Callable[[TxData], Coroutine]):
        self._on_tx.append(cb)
        return cb

    def on_event(self, cb: Callable[[EventData], Coroutine]):
        self._on_event.append(cb)
        return cb

    async def index_blocks(
        self,
        start_height: int,
        end_height: int | None = None,
        on_progress: Callable[[int, int], None] | None = None,
    ) -> int:
        """批量索引区块

        Args:
            start_height: 起始高度(包含)
            end_height: 结束高度(包含),None 表示索引到最新

        Returns:
            已索引的区块数量
        """
        if end_height is None:
            status = await self.rpc.get_status()
            end_height = status.latest_block_height

        indexed = 0
        height = start_height
        logger.info("开始批量索引: %d → %d", start_height, end_height)

        while height <= end_height:
            batch_end = min(height + self.batch_size - 1, end_height)
            tasks = [self._process_block(h) for h in range(height, batch_end + 1)]
            results = await asyncio.gather(*tasks, return_exceptions=True)
            for r in results:
                if isinstance(r, Exception):
                    logger.error("索引错误: %s", r)
                elif r:
                    indexed += 1
            height = batch_end + 1
            if on_progress:
                on_progress(indexed, end_height - start_height + 1)
            await asyncio.sleep(0.1)

        logger.info("索引完成: %d 区块", indexed)
        return indexed

    async def _process_block(self, height: int) -> bool:
        """处理单个区块"""
        async with self._semaphore:
            try:
                block_raw = await self.rpc.get_block(height)
                block_meta = block_raw.get("block", {})
                header = block_meta.get("header", {})
                block_id = block_raw.get("block_id", {})

                time_str = header.get("time", "")
                if time_str:
                    block_time = datetime.fromisoformat(time_str.replace("Z", "+00:00"))
                else:
                    block_time = datetime.utcnow()

                txs_raw = await self.rpc.get_txs_by_height(height)
                txs: list[TxData] = []

                for idx, tx_raw in enumerate(txs_raw):
                    td = TxParser.parse_tx(tx_raw, height, idx)
                    tx_res = tx_raw.get("tx_response", tx_raw)
                    events = EventParser.parse_raw_events(
                        tx_res.get("result", {}).get("events", []),
                        tx_hash=td["hash"],
                        block_height=height,
                        success=td["code"] == 0,
                    )
                    tx = TxData(
                        hash=td["hash"], block_height=height, index=idx,
                        status=td["status"], code=td["code"],
                        gas_used=td["gas_used"], gas_wanted=td["gas_wanted"],
                        fee_amount=td["fee_amount"], fee_denom=td["fee_denom"],
                        memo=td["memo"], sender=td["sender"],
                        contract=td["contract"], action=td["action"],
                        events=events,
                    )
                    txs.append(tx)

                block = BlockData(
                    height=height,
                    hash=block_id.get("hash", ""),
                    time=block_time,
                    proposer=header.get("proposer_address", ""),
                    tx_count=len(txs),
                    gas_used=sum(t.gas_used for t in txs),
                    gas_wanted=sum(t.gas_wanted for t in txs),
                    total_fee=sum(t.fee_amount for t in txs),
                    num_events=sum(len(t.events) for t in txs),
                    transactions=txs,
                )

                for cb in self._on_block:
                    await cb(block)
                for t in txs:
                    for cb in self._on_tx:
                        await cb(t)
                    for e in t.events:
                        for cb in self._on_event:
                            await cb(e)

                self._last_height = height
                return True

            except Exception as e:
                logger.error("处理区块 %d 失败: %s", height, e)
                return False

    async def start_realtime(self, poll_interval: float = 3.0):
        """启动实时索引模式"""
        self._running = True
        status = await self.rpc.get_status()
        self._last_height = status.latest_block_height
        logger.info("实时索引启动, 当前高度: %d", self._last_height)

        while self._running:
            try:
                status = await self.rpc.get_status()
                latest = status.latest_block_height
                if latest > self._last_height:
                    for h in range(self._last_height + 1, latest + 1):
                        await self._process_block(h)
                        self._last_height = h
                await asyncio.sleep(poll_interval)
            except asyncio.CancelledError:
                break
            except Exception as e:
                logger.error("实时索引错误: %s", e)
                await asyncio.sleep(poll_interval * 2)

    def stop(self):
        self._running = False

    async def close(self):
        await self.rpc.close()

2.6 历史数据回填器

# indexer/backfill.py
"""历史数据回填引擎 — 支持全量/增量/断点续传"""

import asyncio
import logging
from datetime import datetime
from pathlib import Path

from indexer.block_indexer import MsgChainIndexer

logger = logging.getLogger(__name__)


class BackfillEngine:
    """历史数据回填引擎"""

    def __init__(
        self,
        rpc_endpoint: str = "http://localhost:26657",
        rest_endpoint: str = "http://localhost:1317",
    ):
        self.rpc_endpoint = rpc_endpoint
        self.rest_endpoint = rest_endpoint

    async def backfill(
        self,
        start_height: int,
        end_height: int | None = None,
        workers: int = 5,
        checkpoint_file: str | None = None,
        on_block=None,
    ) -> int:
        """并行回填历史数据"""
        if end_height is None:
            idx = MsgChainIndexer(self.rpc_endpoint, self.rest_endpoint)
            status = await idx.rpc.get_status()
            end_height = status.latest_block_height
            await idx.close()

        if checkpoint_file:
            cp = await self._load_checkpoint(checkpoint_file)
            if cp and cp > start_height:
                start_height = cp + 1
                logger.info("从断点恢复: %d", start_height)

        total = end_height - start_height + 1
        processed = 0
        start_time = datetime.now()

        indexer = MsgChainIndexer(
            rpc_endpoint=self.rpc_endpoint,
            rest_endpoint=self.rest_endpoint,
            batch_size=workers * 10, concurrency=workers,
        )
        if on_block:
            indexer.on_block(on_block)

        current = start_height
        while current <= end_height:
            batch_end = min(current + 500 - 1, end_height)
            count = await indexer.index_blocks(current, batch_end)
            processed += count
            if checkpoint_file and processed % 1000 == 0:
                await self._save_checkpoint(checkpoint_file, batch_end)
            elapsed = (datetime.now() - start_time).total_seconds()
            rate = processed / elapsed if elapsed > 0 else 0
            pct = (processed / total) * 100
            eta = (total - processed) / rate if rate > 0 else 0
            logger.info("进度: %d/%d (%.1f%%) %.1f blk/s ETA: %.0fs", processed, total, pct, rate, eta)
            current = batch_end + 1

        if checkpoint_file:
            await self._save_checkpoint(checkpoint_file, end_height)

        elapsed = (datetime.now() - start_time).total_seconds()
        logger.info("回填完成: %d 块, %.0fs", processed, elapsed)
        await indexer.close()
        return processed

    async def _load_checkpoint(self, path: str) -> int | None:
        p = Path(path)
        if p.exists():
            try:
                return int(p.read_text().strip())
            except (ValueError, OSError):
                pass
        return None

    async def _save_checkpoint(self, path: str, height: int):
        try:
            Path(path).write_text(str(height))
        except OSError as e:
            logger.warning("保存断点失败: %s", e)

2.7 状态快照

# indexer/state_snapshot.py
"""合约状态快照 — 定期抓取链上合约状态"""

import asyncio
import json
import logging
from datetime import datetime

from indexer.rpc_client import RPCClient

logger = logging.getLogger(__name__)


class StateSnapshotter:
    """状态快照器 — 定期抓取链上合约状态"""

    QUERY_PATHS = {
        "agents": "/custom/agent_registry/agents",
        "sessions": "/custom/agent_payment/sessions",
        "dids": "/custom/did_registry/dids",
        "rules": "/custom/constitution/rules",
        "channels": "/custom/micropayment/channels",
    }

    def __init__(self, rpc_endpoint: str = "http://localhost:26657"):
        self.rpc = RPCClient(rpc_endpoint)

    async def query_contract_state(self, contract_address: str, query_msg: dict):
        """查询特定合约状态"""
        query_hex = json.dumps(query_msg).encode().hex()
        result = await self.rpc.get_abci_query(
            path="/cosmwasm.wasm.v1.Query/SmartContractState",
            data=query_hex,
        )
        value = result.get("response", {}).get("value", "")
        if value:
            return json.loads(bytes.fromhex(value).decode())
        return {}

    async def _decode_value(self, result: dict):
        value = result.get("response", {}).get("value", "")
        if value:
            try:
                return json.loads(bytes.fromhex(value).decode())
            except (ValueError, json.JSONDecodeError):
                return []
        return []

    async def periodic_snapshot(self, interval: int = 3600, callback=None):
        """定时状态快照"""
        while True:
            try:
                height = (await self.rpc.get_status()).latest_block_height
                snapshot = {
                    "height": height,
                    "timestamp": datetime.utcnow().isoformat(),
                    "agents": await self._decode_value(
                        await self.rpc.get_abci_query(path=self.QUERY_PATHS["agents"])
                    ),
                    "sessions": await self._decode_value(
                        await self.rpc.get_abci_query(path=self.QUERY_PATHS["sessions"])
                    ),
                }
                if callback:
                    await callback(snapshot)
                logger.info("快照完成 (高度 %d)", height)
            except Exception as e:
                logger.error("快照失败: %s", e)
            await asyncio.sleep(interval)

    async def close(self):
        await self.rpc.close()

2.8 数据聚合器

# indexer/aggregator.py
"""数据聚合器 — 计算小时/日聚合统计"""

import logging
from datetime import datetime, timedelta
from decimal import Decimal

from indexer.models import AggregatedStats

logger = logging.getLogger(__name__)


class DataAggregator:
    """数据聚合器 — 从原始数据计算聚合统计"""

    def __init__(self):
        self._buffer: list[dict] = []

    def add_block(self, block: dict):
        self._buffer.append({
            "height": block.get("height"),
            "time": block.get("time"),
            "tx_count": block.get("tx_count", 0),
            "gas_used": block.get("gas_used", 0),
            "total_fee": block.get("total_fee", Decimal("0")),
            "senders": set(),
            "contracts": set(),
        })

    def add_tx(self, tx: dict):
        if self._buffer:
            self._buffer[-1]["senders"].add(tx.get("sender", ""))
            if tx.get("contract"):
                self._buffer[-1]["contracts"].add(tx.get("contract", ""))

    async def aggregate_hourly(self, dt: datetime | None = None,
                                store_callback=None) -> AggregatedStats:
        if dt is None:
            dt = datetime.utcnow().replace(minute=0, second=0, microsecond=0)
        hour_end = dt + timedelta(hours=1)
        blocks = [b for b in self._buffer if b["time"] and dt <= b["time"] < hour_end]
        if not blocks:
            return AggregatedStats(timestamp=dt, period="hour")

        total_blocks = len(blocks)
        total_txs = sum(b["tx_count"] for b in blocks)
        all_senders = set().union(*[b["senders"] for b in blocks]) if blocks else set()
        all_contracts = set().union(*[b["contracts"] for b in blocks]) if blocks else set()
        total_fees = sum(b.get("total_fee", Decimal("0")) for b in blocks)
        total_gas = sum(b["gas_used"] for b in blocks)

        stats = AggregatedStats(
            timestamp=dt, period="hour",
            total_blocks=total_blocks, total_txs=total_txs,
            unique_senders=len(all_senders),
            unique_contracts=len(all_contracts),
            avg_gas_used=int(total_gas / total_txs) if total_txs else 0,
            total_fees=total_fees,
        )
        if store_callback:
            await store_callback(stats)
        return stats

    def clear(self):
        self._buffer.clear()

2.9 数据库访问层

# indexer/database.py
"""异步数据库访问层 — 使用 asyncpg"""

import json
import logging

import asyncpg

from indexer.models import BlockData, TxData, EventData

logger = logging.getLogger(__name__)


class Database:
    """异步 PostgreSQL 数据库操作"""

    def __init__(self, dsn: str, pool_size: int = 20):
        self.dsn = dsn
        self.pool_size = pool_size
        self.pool: asyncpg.Pool | None = None

    async def connect(self):
        self.pool = await asyncpg.create_pool(
            dsn=self.dsn, min_size=4, max_size=self.pool_size, command_timeout=30,
        )
        logger.info("数据库连接池已创建 (max=%d)", self.pool_size)

    async def close(self):
        if self.pool:
            await self.pool.close()

    async def store_block(self, block: BlockData):
        """存储区块及关联交易和事件"""
        async with self.pool.acquire() as conn:
            async with conn.transaction():
                await conn.execute(
                    """INSERT INTO blocks (height, hash, time, proposer, tx_count,
                       gas_used, gas_wanted, total_fee, app_hash, num_events)
                       VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10)
                       ON CONFLICT (height) DO UPDATE SET
                           hash = EXCLUDED.hash, time = EXCLUDED.time""",
                    block.height, block.hash, block.time, block.proposer,
                    block.tx_count, block.gas_used, block.gas_wanted,
                    str(block.total_fee), block.app_hash, block.num_events,
                )
                for tx in block.transactions:
                    await self._store_tx(conn, tx)
                    for event in tx.events:
                        await self._store_event(conn, event)
                await self._update_state(conn, block)

    async def _store_tx(self, conn: asyncpg.Connection, tx: TxData):
        await conn.execute(
            """INSERT INTO transactions (hash, block_height, index, status, code,
               gas_used, gas_wanted, fee_amount, fee_denom, memo, sender, contract, action, raw_log)
               VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12, $13, $14)
               ON CONFLICT (hash) DO NOTHING""",
            tx.hash, tx.block_height, tx.index, tx.status, tx.code,
            tx.gas_used, tx.gas_wanted, str(tx.fee_amount), tx.fee_denom,
            tx.memo, tx.sender, tx.contract, tx.action, tx.raw_log,
        )

    async def _store_event(self, conn: asyncpg.Connection, event: EventData):
        await conn.execute(
            """INSERT INTO events (tx_hash, block_height, event_type, action,
               contract_address, attributes, success)
               VALUES ($1, $2, $3, $4, $5, $6, $7)""",
            event.tx_hash, event.block_height, event.event_type,
            event.action, event.contract_address,
            json.dumps(event.attributes), event.success,
        )

    async def _update_state(self, conn: asyncpg.Connection, block: BlockData):
        txs = len(block.transactions)
        evts = sum(len(t.events) for t in block.transactions)
        await conn.execute(
            """UPDATE indexer_state SET last_processed_height = $1,
               last_processed_hash = $2, blocks_indexed = blocks_indexed + 1,
               txs_indexed = txs_indexed + $3, events_indexed = events_indexed + $4,
               updated_at = NOW() WHERE id = 1""",
            block.height, block.hash, txs, evts,
        )

    async def get_last_processed_height(self) -> int:
        row = await self.pool.fetchrow(
            "SELECT last_processed_height FROM indexer_state WHERE id = 1"
        )
        return row["last_processed_height"] if row else 0

    async def get_blocks(self, limit=20, offset=0):
        return await self.pool.fetch(
            "SELECT * FROM blocks ORDER BY height DESC LIMIT $1 OFFSET $2", limit, offset
        )

    async def get_tx_by_hash(self, tx_hash: str):
        return await self.pool.fetchrow(
            "SELECT * FROM transactions WHERE hash = $1", tx_hash
        )

    async def get_chain_overview(self) -> dict:
        row = await self.pool.fetchrow("""
            SELECT
                (SELECT MAX(height) FROM blocks) AS latest_height,
                (SELECT COUNT(*) FROM blocks WHERE time >= NOW() - INTERVAL '24 hours') AS blocks_24h,
                (SELECT COUNT(*) FROM transactions WHERE created_at >= NOW() - INTERVAL '24 hours') AS txs_24h,
                (SELECT COUNT(*) FROM agents WHERE status = 'ACTIVE') AS active_agents,
                (SELECT COUNT(*) FROM agents) AS total_agents
        """)
        return dict(row) if row else {}

    async def execute_raw(self, query: str, *args):
        return await self.pool.fetch(query, *args)

2.10 索引器主程序入口

# indexer/main.py
"""索引器主程序入口"""

import asyncio
import logging
import os

from indexer.block_indexer import MsgChainIndexer
from indexer.database import Database
from indexer.backfill import BackfillEngine
from indexer.aggregator import DataAggregator

logging.basicConfig(level=logging.INFO)
logger = logging.getLogger(__name__)


async def main():
    rpc = os.environ.get("RPC_ENDPOINT", "http://localhost:26657")
    rest = os.environ.get("REST_ENDPOINT", "http://localhost:1317")
    dsn = os.environ.get(
        "DATABASE_URL",
        "postgresql://msgchain:msgchain_secret@localhost:5432/msgchain_analytics",
    )
    mode = os.environ.get("INDEXER_MODE", "realtime")
    start = int(os.environ.get("START_HEIGHT", "0"))

    db = Database(dsn)
    await db.connect()

    aggregator = DataAggregator()
    indexer = MsgChainIndexer(rpc_endpoint=rpc, rest_endpoint=rest)

    @indexer.on_block
    async def handle_block(block):
        await db.store_block(block)
        aggregator.add_block(block.__dict__)

    @indexer.on_tx
    async def handle_tx(tx):
        aggregator.add_tx(tx.__dict__)

    if mode == "backfill":
        engine = BackfillEngine(rpc, rest)
        await engine.backfill(start, workers=10, on_block=handle_block)
    else:
        await indexer.start_realtime(poll_interval=3.0)

    await indexer.close()
    await db.close()


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

3. 数据库模型

3.1 PostgreSQL 基础 Schema

-- migrations/analytics_v1/001_base_schema.sql
-- MSG Chain 分析数据库 Schema (PostgreSQL 15+)

CREATE EXTENSION IF NOT EXISTS "uuid-ossp";
CREATE EXTENSION IF NOT EXISTS "pg_trgm";
CREATE EXTENSION IF NOT EXISTS "btree_gin";

-- 枚举类型
CREATE TYPE tx_status AS ENUM ('SUCCESS', 'FAILED');
CREATE TYPE agent_status AS ENUM ('ACTIVE', 'INACTIVE', 'SUSPENDED', 'DEREGISTERED');
CREATE TYPE session_status AS ENUM ('OPEN', 'FUNDED', 'RELEASED', 'DISPUTED', 'CLOSED', 'EXPIRED');
CREATE TYPE channel_status AS ENUM ('OPEN', 'SETTLING', 'CLOSED');
CREATE TYPE indexer_mode AS ENUM ('realtime', 'backfill', 'idle');
CREATE TYPE agg_period AS ENUM ('hour', 'day', 'week', 'month');

-- blocks 表 — 区块头数据
CREATE TABLE blocks (
    height          BIGINT PRIMARY KEY,
    hash            VARCHAR(64) NOT NULL UNIQUE,
    time            TIMESTAMPTZ NOT NULL,
    proposer        VARCHAR(64) NOT NULL,
    tx_count        INTEGER NOT NULL DEFAULT 0,
    gas_used        BIGINT NOT NULL DEFAULT 0,
    gas_wanted      BIGINT NOT NULL DEFAULT 0,
    total_fee       NUMERIC(50, 0) NOT NULL DEFAULT 0,
    app_hash        VARCHAR(64),
    validators_hash VARCHAR(64),
    consensus_hash  VARCHAR(64),
    num_events      INTEGER NOT NULL DEFAULT 0,
    created_at      TIMESTAMPTZ NOT NULL DEFAULT NOW(),
    CONSTRAINT blocks_hash_format CHECK (hash ~ '^[A-Fa-f0-9]{64}$')
);

COMMENT ON TABLE blocks IS '区块头数据 — 每个区块对应一行';
COMMENT ON COLUMN blocks.time IS '出块时间 (UTC)';
COMMENT ON COLUMN blocks.proposer IS '提议人地址 (bech32 msg 前缀)';

CREATE INDEX idx_blocks_time_desc ON blocks (time DESC);
CREATE INDEX idx_blocks_proposer ON blocks (proposer);
CREATE INDEX idx_blocks_time_date ON blocks ((time::date));
CREATE INDEX idx_blocks_gas_used ON blocks (gas_used DESC);

-- transactions 表 — 交易数据
CREATE TABLE transactions (
    hash            VARCHAR(64) PRIMARY KEY,
    block_height    BIGINT NOT NULL REFERENCES blocks(height) ON DELETE CASCADE,
    index           INTEGER NOT NULL,
    status          tx_status NOT NULL,
    code            INTEGER NOT NULL DEFAULT 0,
    gas_used        BIGINT NOT NULL DEFAULT 0,
    gas_wanted      BIGINT NOT NULL DEFAULT 0,
    fee_amount      NUMERIC(50, 0) NOT NULL DEFAULT 0,
    fee_denom       VARCHAR(16) NOT NULL DEFAULT 'umsg',
    memo            TEXT NOT NULL DEFAULT '',
    sender          VARCHAR(64) NOT NULL DEFAULT '',
    contract        VARCHAR(64) NOT NULL DEFAULT '',
    action          VARCHAR(128) NOT NULL DEFAULT '',
    raw_log         TEXT,
    created_at      TIMESTAMPTZ NOT NULL DEFAULT NOW(),
    CONSTRAINT unique_tx_in_block UNIQUE (block_height, index)
);

COMMENT ON TABLE transactions IS '交易数据 — 每笔交易对应一行';
COMMENT ON COLUMN transactions.sender IS '交易发送者地址 (msg1...)';
COMMENT ON COLUMN transactions.contract IS '交互的合约地址';
COMMENT ON COLUMN transactions.action IS '交易动作 (CreateSession, RegisterAgent 等)';

CREATE INDEX idx_transactions_block_height ON transactions (block_height DESC);
CREATE INDEX idx_transactions_sender ON transactions (sender);
CREATE INDEX idx_transactions_contract ON transactions (contract);
CREATE INDEX idx_transactions_action ON transactions (action);
CREATE INDEX idx_transactions_status ON transactions (status);
CREATE INDEX idx_transactions_created_at ON transactions (created_at DESC);
CREATE INDEX idx_transactions_sender_contract ON transactions (sender, contract);
CREATE INDEX idx_transactions_block_action ON transactions (block_height, action);
CREATE INDEX idx_transactions_hash_trgm ON transactions USING GIN (hash gin_trgm_ops);

-- events 表 — 合约事件日志(核心分析表)
CREATE TABLE events (
    id              BIGSERIAL PRIMARY KEY,
    event_uuid      UUID NOT NULL DEFAULT uuid_generate_v4(),
    tx_hash         VARCHAR(64) NOT NULL REFERENCES transactions(hash) ON DELETE CASCADE,
    block_height    BIGINT NOT NULL REFERENCES blocks(height) ON DELETE CASCADE,
    event_type      VARCHAR(64) NOT NULL DEFAULT 'wasm',
    action          VARCHAR(128) NOT NULL,
    contract_address VARCHAR(64) NOT NULL DEFAULT '',
    attributes      JSONB NOT NULL DEFAULT '{}',
    success         BOOLEAN NOT NULL DEFAULT true,
    created_at      TIMESTAMPTZ NOT NULL DEFAULT NOW()
);

COMMENT ON TABLE events IS '合约事件日志 — 类似 Etherscan Event Logs';
COMMENT ON COLUMN events.action IS '事件动作 (CreateSession, RegisterAgent 等)';
COMMENT ON COLUMN events.attributes IS '事件属性 JSONB — 包含 agent_id, session_id, amount 等';

CREATE INDEX idx_events_tx_hash ON events (tx_hash);
CREATE INDEX idx_events_block_height ON events (block_height DESC);
CREATE INDEX idx_events_action ON events (action);
CREATE INDEX idx_events_contract_address ON events (contract_address);
CREATE INDEX idx_events_contract_action ON events (contract_address, action);
CREATE INDEX idx_events_attributes_gin ON events USING GIN (attributes jsonb_path_ops);
CREATE INDEX idx_events_action_block ON events (action, block_height DESC);
CREATE INDEX idx_events_created_at ON events (created_at DESC);
CREATE INDEX idx_events_contract_action_time ON events (contract_address, action, block_height DESC);

-- messages 表 — CosmWasm 消息
CREATE TABLE messages (
    id              BIGSERIAL PRIMARY KEY,
    tx_hash         VARCHAR(64) NOT NULL REFERENCES transactions(hash) ON DELETE CASCADE,
    block_height    BIGINT NOT NULL REFERENCES blocks(height) ON DELETE CASCADE,
    index           INTEGER NOT NULL,
    message_type    VARCHAR(128) NOT NULL,
    sender          VARCHAR(64) NOT NULL DEFAULT '',
    contract_address VARCHAR(64) NOT NULL DEFAULT '',
    msg_type        VARCHAR(128) NOT NULL DEFAULT '',
    data            JSONB NOT NULL DEFAULT '{}',
    funds           JSONB NOT NULL DEFAULT '[]',
    created_at      TIMESTAMPTZ NOT NULL DEFAULT NOW()
);

CREATE INDEX idx_messages_tx_hash ON messages (tx_hash);
CREATE INDEX idx_messages_block_height ON messages (block_height DESC);
CREATE INDEX idx_messages_sender ON messages (sender);
CREATE INDEX idx_messages_contract ON messages (contract_address);
CREATE INDEX idx_messages_msg_type ON messages (msg_type);

-- contracts 表 — 合约元信息
CREATE TABLE contracts (
    address        VARCHAR(64) PRIMARY KEY,
    code_id        BIGINT NOT NULL,
    creator        VARCHAR(64) NOT NULL,
    admin          VARCHAR(64),
    label          VARCHAR(256) NOT NULL,
    contract_type  VARCHAR(64),
    created_at     TIMESTAMPTZ NOT NULL,
    CONSTRAINT unique_contract_address UNIQUE (address)
);

CREATE INDEX idx_contracts_code_id ON contracts (code_id);
CREATE INDEX idx_contracts_creator ON contracts (creator);
CREATE INDEX idx_contracts_type ON contracts (contract_type);

-- agents 表 — AI Agent 当前状态(由事件驱动更新)
CREATE TABLE agents (
    agent_id        VARCHAR(128) PRIMARY KEY,
    owner           VARCHAR(64) NOT NULL,
    name            VARCHAR(256) NOT NULL,
    status          agent_status NOT NULL DEFAULT 'ACTIVE',
    agent_type      VARCHAR(32) NOT NULL DEFAULT 'ai',
    description     TEXT NOT NULL DEFAULT '',
    endpoint        VARCHAR(512) NOT NULL DEFAULT '',
    capabilities    JSONB NOT NULL DEFAULT '[]',
    fee             VARCHAR(64) NOT NULL DEFAULT '0',
    contract_address VARCHAR(64),
    created_height  BIGINT NOT NULL,
    created_at      TIMESTAMPTZ NOT NULL,
    updated_at      TIMESTAMPTZ NOT NULL DEFAULT NOW()
);

COMMENT ON TABLE agents IS 'AI Agent 注册信息 — 由 RegisterAgent/SetStatus/DeregisterAgent 事件驱动更新';

CREATE INDEX idx_agents_owner ON agents (owner);
CREATE INDEX idx_agents_status ON agents (status);
CREATE INDEX idx_agents_type ON agents (agent_type);
CREATE INDEX idx_agents_created_at ON agents (created_at DESC);
CREATE INDEX idx_agents_capabilities_gin ON agents USING GIN (capabilities);
CREATE INDEX idx_agents_name_trgm ON agents USING GIN (name gin_trgm_ops);

-- sessions 表 — 支付会话
CREATE TABLE sessions (
    session_id      VARCHAR(128) PRIMARY KEY,
    agent_id        VARCHAR(128) NOT NULL REFERENCES agents(agent_id) ON DELETE CASCADE,
    client          VARCHAR(64) NOT NULL,
    status          session_status NOT NULL DEFAULT 'OPEN',
    amount          NUMERIC(78, 0) NOT NULL DEFAULT 0,
    deposit         NUMERIC(78, 0) NOT NULL DEFAULT 0,
    balance         NUMERIC(78, 0) NOT NULL DEFAULT 0,
    deadline        TIMESTAMPTZ,
    terms_hash      VARCHAR(64),
    contract_address VARCHAR(64),
    created_height  BIGINT NOT NULL,
    created_at      TIMESTAMPTZ NOT NULL,
    updated_at      TIMESTAMPTZ NOT NULL DEFAULT NOW(),
    closed_at       TIMESTAMPTZ,
    closed_reason   VARCHAR(64)
);

COMMENT ON TABLE sessions IS 'AI Agent 支付会话 — 由 CreateSession/FundSession/ReleasePayment/CloseSession 事件驱动';

CREATE INDEX idx_sessions_agent_id ON sessions (agent_id);
CREATE INDEX idx_sessions_client ON sessions (client);
CREATE INDEX idx_sessions_status ON sessions (status);
CREATE INDEX idx_sessions_agent_status ON sessions (agent_id, status);
CREATE INDEX idx_sessions_created_at ON sessions (created_at DESC);
CREATE INDEX idx_sessions_client_status ON sessions (client, status);

-- did_documents 表 — 去中心化身份文档
CREATE TABLE did_documents (
    did                 VARCHAR(256) PRIMARY KEY,
    controller          VARCHAR(64) NOT NULL,
    public_key          TEXT NOT NULL,
    authentication_methods JSONB NOT NULL DEFAULT '[]',
    service_endpoints   JSONB NOT NULL DEFAULT '[]',
    active              BOOLEAN NOT NULL DEFAULT true,
    contract_address    VARCHAR(64),
    created_height      BIGINT NOT NULL,
    created_at          TIMESTAMPTZ NOT NULL,
    updated_at          TIMESTAMPTZ NOT NULL DEFAULT NOW(),
    deactivated_at      TIMESTAMPTZ
);

CREATE INDEX idx_did_controller ON did_documents (controller);
CREATE INDEX idx_did_active ON did_documents (active);

-- constitution_rules 表 — AI Agent 治理规则
CREATE TABLE constitution_rules (
    rule_id         VARCHAR(128) PRIMARY KEY,
    rule_type       VARCHAR(32) NOT NULL,
    description     TEXT NOT NULL DEFAULT '',
    parameters      JSONB NOT NULL DEFAULT '{}',
    active          BOOLEAN NOT NULL DEFAULT true,
    version         VARCHAR(32) NOT NULL DEFAULT '1.0',
    contract_address VARCHAR(64),
    created_height  BIGINT NOT NULL,
    created_at      TIMESTAMPTZ NOT NULL,
    updated_at      TIMESTAMPTZ NOT NULL DEFAULT NOW(),
    removed_at      TIMESTAMPTZ
);

CREATE INDEX idx_rules_type ON constitution_rules (rule_type);
CREATE INDEX idx_rules_active ON constitution_rules (active);

-- micropayment_channels 表 — 微支付状态通道
CREATE TABLE micropayment_channels (
    channel_id      VARCHAR(128) PRIMARY KEY,
    participant_a   VARCHAR(64) NOT NULL,
    participant_b   VARCHAR(64) NOT NULL,
    capacity        NUMERIC(78, 0) NOT NULL,
    balance_a       NUMERIC(78, 0) NOT NULL DEFAULT 0,
    balance_b       NUMERIC(78, 0) NOT NULL DEFAULT 0,
    status          channel_status NOT NULL DEFAULT 'OPEN',
    expiry          TIMESTAMPTZ,
    contract_address VARCHAR(64),
    created_height  BIGINT NOT NULL,
    created_at      TIMESTAMPTZ NOT NULL,
    updated_at      TIMESTAMPTZ NOT NULL DEFAULT NOW(),
    closed_at       TIMESTAMPTZ
);

CREATE INDEX idx_channels_participant_a ON micropayment_channels (participant_a);
CREATE INDEX idx_channels_participant_b ON micropayment_channels (participant_b);
CREATE INDEX idx_channels_status ON micropayment_channels (status);

-- ============================================================
-- 聚合统计表
-- ============================================================

-- hourly_stats — 小时级网络统计(仪表盘时间序列数据源)
CREATE TABLE hourly_stats (
    time            TIMESTAMPTZ NOT NULL,
    period          agg_period NOT NULL DEFAULT 'hour',
    total_blocks    INTEGER NOT NULL DEFAULT 0,
    total_txs       INTEGER NOT NULL DEFAULT 0,
    total_events    INTEGER NOT NULL DEFAULT 0,
    unique_senders  INTEGER NOT NULL DEFAULT 0,
    unique_contracts INTEGER NOT NULL DEFAULT 0,
    avg_gas_price   NUMERIC(30, 18) NOT NULL DEFAULT 0,
    avg_gas_used    INTEGER NOT NULL DEFAULT 0,
    total_fees      NUMERIC(50, 0) NOT NULL DEFAULT 0,
    success_rate    NUMERIC(5, 4) NOT NULL DEFAULT 1.0000,
    active_agents   INTEGER NOT NULL DEFAULT 0,
    new_agents      INTEGER NOT NULL DEFAULT 0,
    new_sessions    INTEGER NOT NULL DEFAULT 0,
    session_volume  NUMERIC(78, 0) NOT NULL DEFAULT 0,
    created_at      TIMESTAMPTZ NOT NULL DEFAULT NOW(),
    PRIMARY KEY (time, period)
);

CREATE INDEX idx_hourly_stats_time ON hourly_stats (time DESC);

-- daily_stats — 日级网络统计
CREATE TABLE daily_stats (
    date            DATE NOT NULL PRIMARY KEY,
    total_blocks    INTEGER NOT NULL DEFAULT 0,
    total_txs       INTEGER NOT NULL DEFAULT 0,
    total_events    INTEGER NOT NULL DEFAULT 0,
    unique_senders  INTEGER NOT NULL DEFAULT 0,
    unique_contracts INTEGER NOT NULL DEFAULT 0,
    avg_gas_price   NUMERIC(30, 18) NOT NULL DEFAULT 0,
    avg_gas_used    INTEGER NOT NULL DEFAULT 0,
    total_fees      NUMERIC(50, 0) NOT NULL DEFAULT 0,
    success_rate    NUMERIC(5, 4) NOT NULL DEFAULT 1.0000,
    active_agents   INTEGER NOT NULL DEFAULT 0,
    new_agents      INTEGER NOT NULL DEFAULT 0,
    new_sessions    INTEGER NOT NULL DEFAULT 0,
    session_volume  NUMERIC(78, 0) NOT NULL DEFAULT 0,
    dau             INTEGER NOT NULL DEFAULT 0,
    mau             INTEGER NOT NULL DEFAULT 0,
    avg_block_time  NUMERIC(10, 2) NOT NULL DEFAULT 0,
    created_at      TIMESTAMPTZ NOT NULL DEFAULT NOW()
);

CREATE INDEX idx_daily_stats_date ON daily_stats (date DESC);

-- agent_daily_stats — Agent 日级统计
CREATE TABLE agent_daily_stats (
    agent_id        VARCHAR(128) NOT NULL REFERENCES agents(agent_id) ON DELETE CASCADE,
    date            DATE NOT NULL,
    tx_count        INTEGER NOT NULL DEFAULT 0,
    session_count   INTEGER NOT NULL DEFAULT 0,
    total_earned    NUMERIC(78, 0) NOT NULL DEFAULT 0,
    total_gas_used  BIGINT NOT NULL DEFAULT 0,
    unique_clients  INTEGER NOT NULL DEFAULT 0,
    event_count     INTEGER NOT NULL DEFAULT 0,
    created_at      TIMESTAMPTZ NOT NULL DEFAULT NOW(),
    PRIMARY KEY (agent_id, date)
);

CREATE INDEX idx_agent_daily_stats_date ON agent_daily_stats (date DESC);

-- contract_daily_stats — 合约日级统计
CREATE TABLE contract_daily_stats (
    contract_address VARCHAR(64) NOT NULL,
    date            DATE NOT NULL,
    tx_count        INTEGER NOT NULL DEFAULT 0,
    event_count     INTEGER NOT NULL DEFAULT 0,
    unique_senders  INTEGER NOT NULL DEFAULT 0,
    total_gas_used  BIGINT NOT NULL DEFAULT 0,
    total_fees      NUMERIC(50, 0) NOT NULL DEFAULT 0,
    created_at      TIMESTAMPTZ NOT NULL DEFAULT NOW(),
    PRIMARY KEY (contract_address, date)
);

CREATE INDEX idx_contract_daily_stats_date ON contract_daily_stats (date DESC);
CREATE INDEX idx_contract_daily_stats_tx ON contract_daily_stats (tx_count DESC);

-- indexer_state 表 — 索引器同步状态
CREATE TABLE indexer_state (
    id                    INTEGER PRIMARY KEY DEFAULT 1,
    last_processed_height BIGINT NOT NULL DEFAULT 0,
    last_processed_hash   VARCHAR(64),
    chain_id              VARCHAR(64) NOT NULL DEFAULT 'msg-chain-1',
    mode                  indexer_mode NOT NULL DEFAULT 'idle',
    blocks_indexed        BIGINT NOT NULL DEFAULT 0,
    txs_indexed           BIGINT NOT NULL DEFAULT 0,
    events_indexed        BIGINT NOT NULL DEFAULT 0,
    errors_count          BIGINT NOT NULL DEFAULT 0,
    last_error_message    TEXT,
    started_at            TIMESTAMPTZ NOT NULL DEFAULT NOW(),
    updated_at            TIMESTAMPTZ NOT NULL DEFAULT NOW(),
    CONSTRAINT single_row CHECK (id = 1)
);

COMMENT ON TABLE indexer_state IS '索引器同步状态 — 断点续传和监控';

INSERT INTO indexer_state (id, last_processed_height) VALUES (1, 0);

3.2 物化视图

-- migrations/analytics_v1/002_materialized_views.sql
-- 物化视图 — 预聚合高频查询

-- 日活跃用户 MV
CREATE MATERIALIZED VIEW mv_daily_active_users AS
SELECT
    DATE(b.time) AS date,
    COUNT(DISTINCT t.sender) AS dau,
    COUNT(DISTINCT t.contract) AS active_contracts,
    COUNT(*) AS total_txs,
    SUM(t.gas_used) AS total_gas_used,
    SUM(t.fee_amount) AS total_fees
FROM blocks b
JOIN transactions t ON t.block_height = b.height
WHERE t.status = 'SUCCESS'
GROUP BY DATE(b.time)
ORDER BY date DESC
WITH DATA;

CREATE UNIQUE INDEX idx_mv_dau_date ON mv_daily_active_users (date);

COMMENT ON MATERIALIZED VIEW mv_daily_active_users IS '日活跃用户 — 按天聚合';

-- 合约排名 MV
CREATE MATERIALIZED VIEW mv_contract_rankings AS
SELECT
    t.contract AS contract_address,
    c.label AS contract_label,
    c.contract_type,
    COUNT(*) AS total_txs,
    COUNT(DISTINCT t.sender) AS unique_users,
    SUM(t.gas_used) AS total_gas_used,
    SUM(t.fee_amount) AS total_fees,
    MIN(b.time) AS first_seen,
    MAX(b.time) AS last_seen,
    ROUND(100.0 * COUNT(*) FILTER (WHERE t.status = 'SUCCESS')
          / NULLIF(COUNT(*), 0), 2) AS success_rate
FROM transactions t
JOIN blocks b ON b.height = t.block_height
LEFT JOIN contracts c ON c.address = t.contract
WHERE t.contract != ''
GROUP BY t.contract, c.label, c.contract_type
ORDER BY total_txs DESC
WITH DATA;

CREATE UNIQUE INDEX idx_mv_contract_rankings_addr ON mv_contract_rankings (contract_address);
CREATE INDEX idx_mv_contract_rankings_tx ON mv_contract_rankings (total_txs DESC);

COMMENT ON MATERIALIZED VIEW mv_contract_rankings IS '合约排名 — 按交易量排序';

-- Agent 经济 MV
CREATE MATERIALIZED VIEW mv_agent_economics AS
SELECT
    a.agent_id, a.name, a.owner, a.status, a.agent_type,
    COUNT(DISTINCT s.session_id) AS total_sessions,
    COUNT(DISTINCT s.client) AS unique_clients,
    COALESCE(SUM(s.amount), 0) AS total_volume,
    COALESCE(AVG(s.amount) FILTER (WHERE s.status = 'RELEASED'), 0) AS avg_session_value,
    COUNT(*) FILTER (WHERE s.status = 'OPEN') AS active_sessions,
    COUNT(*) FILTER (WHERE s.status = 'DISPUTED') AS disputed_sessions,
    MAX(s.created_at) AS last_session_at,
    a.created_at AS registered_at
FROM agents a
LEFT JOIN sessions s ON s.agent_id = a.agent_id
GROUP BY a.agent_id, a.name, a.owner, a.status, a.agent_type, a.created_at
ORDER BY total_volume DESC
WITH DATA;

CREATE UNIQUE INDEX idx_mv_agent_economics_id ON mv_agent_economics (agent_id);
CREATE INDEX idx_mv_agent_economics_volume ON mv_agent_economics (total_volume DESC);

COMMENT ON MATERIALIZED VIEW mv_agent_economics IS 'Agent 经济指标 — 会话量、收入、活跃度';

-- Gas 价格分布 MV
CREATE MATERIALIZED VIEW mv_gas_price_distribution AS
SELECT
    DATE(b.time) AS date,
    PERCENTILE_CONT(0.5) WITHIN GROUP (
        ORDER BY t.fee_amount::NUMERIC / NULLIF(t.gas_used, 1)
    ) AS median_gas_price,
    PERCENTILE_CONT(0.95) WITHIN GROUP (
        ORDER BY t.fee_amount::NUMERIC / NULLIF(t.gas_used, 1)
    ) AS p95_gas_price,
    AVG(t.fee_amount::NUMERIC / NULLIF(t.gas_used, 1)) AS mean_gas_price
FROM transactions t
JOIN blocks b ON b.height = t.block_height
WHERE t.status = 'SUCCESS' AND t.gas_used > 0
GROUP BY DATE(b.time)
ORDER BY date DESC
WITH DATA;

CREATE UNIQUE INDEX idx_mv_gas_price_date ON mv_gas_price_distribution (date);

COMMENT ON MATERIALIZED VIEW mv_gas_price_distribution IS 'Gas 价格百分位数分布';

-- 刷新所有物化视图的函数
CREATE OR REPLACE FUNCTION refresh_analytics_views()
RETURNS void AS $$
BEGIN
    REFRESH MATERIALIZED VIEW CONCURRENTLY mv_daily_active_users;
    REFRESH MATERIALIZED VIEW CONCURRENTLY mv_contract_rankings;
    REFRESH MATERIALIZED VIEW CONCURRENTLY mv_agent_economics;
    REFRESH MATERIALIZED VIEW CONCURRENTLY mv_gas_price_distribution;
END;
$$ LANGUAGE plpgsql;

COMMENT ON FUNCTION refresh_analytics_views IS '并发刷新所有分析物化视图(不阻塞读)';

3.3 表分区策略

-- migrations/analytics_v1/003_partitions.sql

-- 方式一: 原生 PostgreSQL 分区
-- events 表是最大的表,按 block_height 范围分区
CREATE TABLE events_partitioned (
    LIKE events INCLUDING ALL
) PARTITION BY RANGE (block_height);

-- 每 100,000 块一个分区 (~1-2 周数据)
CREATE TABLE events_0_100k PARTITION OF events_partitioned
    FOR VALUES FROM (0) TO (100000);
CREATE TABLE events_100k_200k PARTITION OF events_partitioned
    FOR VALUES FROM (100000) TO (200000);
CREATE TABLE events_200k_300k PARTITION OF events_partitioned
    FOR VALUES FROM (200000) TO (300000);
CREATE TABLE events_300k_400k PARTITION OF events_partitioned
    FOR VALUES FROM (300000) TO (400000);
CREATE TABLE events_default PARTITION OF events_partitioned DEFAULT;

-- 自动创建新分区的函数
CREATE OR REPLACE FUNCTION create_events_partition(
    p_start BIGINT, p_end BIGINT
) RETURNS void AS $$
BEGIN
    EXECUTE format(
        'CREATE TABLE events_%s_%s PARTITION OF events_partitioned '
        'FOR VALUES FROM (%s) TO (%s)',
        p_start, p_end, p_start, p_end
    );
END;
$$ LANGUAGE plpgsql;

-- 方式二: TimescaleDB hypertable(推荐)
-- CREATE EXTENSION IF NOT EXISTS timescaledb;
-- SELECT create_hypertable('events', 'block_height', chunk_time_interval => 100000);
-- SELECT create_hypertable('transactions', 'block_height', chunk_time_interval => 100000);
-- SELECT create_hypertable('hourly_stats', 'time', chunk_time_interval => INTERVAL '7 days');

3.4 数据保留策略

-- migrations/analytics_v1/004_retention.sql

-- 保留策略配置
CREATE TABLE retention_policies (
    table_name      VARCHAR(128) PRIMARY KEY,
    retention_days  INTEGER NOT NULL DEFAULT 365,
    enabled         BOOLEAN NOT NULL DEFAULT true,
    created_at      TIMESTAMPTZ NOT NULL DEFAULT NOW()
);

COMMENT ON TABLE retention_policies IS '数据保留策略 — 按表配置';

INSERT INTO retention_policies (table_name, retention_days) VALUES
    ('events', 365),
    ('transactions', 730),
    ('blocks', 730),
    ('messages', 365),
    ('hourly_stats', 1095),
    ('daily_stats', 7300),
    ('agent_daily_stats', 730),
    ('contract_daily_stats', 730);

-- 数据清理函数 — 返回每张表的删除行数
CREATE OR REPLACE FUNCTION apply_retention()
RETURNS TABLE (tbl TEXT, deleted BIGINT) AS $$
DECLARE
    r RECORD;
    cutoff TIMESTAMPTZ;
    d BIGINT;
BEGIN
    FOR r IN SELECT * FROM retention_policies WHERE enabled = true LOOP
        cutoff := NOW() - (r.retention_days || ' days')::INTERVAL;
        CASE r.table_name
            WHEN 'events' THEN DELETE FROM events WHERE created_at < cutoff;
            WHEN 'transactions' THEN DELETE FROM transactions WHERE created_at < cutoff;
            WHEN 'blocks' THEN DELETE FROM blocks WHERE time < cutoff;
            WHEN 'hourly_stats' THEN DELETE FROM hourly_stats WHERE time < cutoff;
            WHEN 'daily_stats' THEN DELETE FROM daily_stats WHERE date < cutoff::date;
            ELSE EXECUTE format('DELETE FROM %I WHERE created_at < $1', r.table_name) USING cutoff;
        END CASE;
        GET DIAGNOSTICS d = ROW_COUNT;
        tbl := r.table_name; deleted := d; RETURN NEXT;
    END LOOP;
END;
$$ LANGUAGE plpgsql;

3.5 索引策略总结

高频查询路径 → 对应的索引策略

按高度查询区块             → PRIMARY KEY (height)
按哈希查询交易             → PRIMARY KEY (hash) + GIN trigram
按 action 查事件           → idx_events_action (B-tree)
按合约地址查事件           → idx_events_contract_address (B-tree)
按合约地址 + action       → idx_events_contract_action (复合)
JSONB 属性查询             → idx_events_attributes_gin (GIN jsonb_path_ops)
最新区块列表               → idx_blocks_time_desc (DESC)
Agent 的所有会话           → idx_sessions_agent_status (复合)
按发送者查交易             → idx_transactions_sender (B-tree)
按时间范围扫描             → idx_blocks_time_date (基于函数的索引)
模糊搜索 (hash/address)   → GIN trigram 索引

4. SQL 分析查询

4.1 网络概览查询

-- 4.1.1 实时网络核心指标
SELECT
    (SELECT MAX(height) FROM blocks) AS latest_block_height,
    (SELECT COUNT(*) FROM blocks) AS total_blocks,
    (SELECT COUNT(*) FROM transactions) AS total_transactions,
    (SELECT COUNT(*) FROM events) AS total_events,
    (SELECT COUNT(*) FROM agents) AS total_agents,
    (SELECT COUNT(*) FROM agents WHERE status = 'ACTIVE') AS active_agents,
    (SELECT COUNT(*) FROM sessions WHERE status IN ('OPEN', 'FUNDED')) AS active_sessions,
    ROUND(
        (SELECT EXTRACT(EPOCH FROM AVG(time_diff)) FROM (
            SELECT time - LAG(time) OVER (ORDER BY height) AS time_diff
            FROM blocks WHERE time >= NOW() - INTERVAL '100 blocks'
        ) sub WHERE time_diff IS NOT NULL)
    ) AS avg_block_time_seconds;

-- 4.1.2 24 小时活动统计
SELECT
    COUNT(DISTINCT t.sender) AS unique_active_addresses,
    COUNT(DISTINCT t.contract) AS unique_contracts_interacted,
    COUNT(*) AS total_transactions,
    SUM(t.gas_used) AS total_gas_used,
    COALESCE(SUM(t.fee_amount) FILTER (WHERE t.fee_denom = 'umsg'), 0) AS total_fees_umsg,
    COUNT(*) FILTER (WHERE t.status = 'FAILED') AS failed_transactions,
    ROUND(
        100.0 * COUNT(*) FILTER (WHERE t.status = 'SUCCESS') / NULLIF(COUNT(*), 0), 2
    ) AS success_rate_pct,
    COUNT(*) FILTER (WHERE e.action = 'CreateSession') AS new_sessions,
    COUNT(*) FILTER (WHERE e.action = 'RegisterAgent') AS new_agents
FROM transactions t
LEFT JOIN events e ON e.tx_hash = t.hash
WHERE t.created_at >= NOW() - INTERVAL '24 hours';

-- 4.1.3 区块生产率趋势(30天)
SELECT
    DATE(time) AS block_date,
    COUNT(*) AS blocks_produced,
    COUNT(DISTINCT proposer) AS unique_proposers,
    ROUND(86400.0 / NULLIF(COUNT(*), 0), 2) AS avg_block_time_seconds,
    SUM(tx_count) AS total_tx_in_blocks,
    AVG(tx_count)::NUMERIC(10, 2) AS avg_tx_per_block,
    AVG(gas_used)::BIGINT AS avg_gas_per_block,
    MAX(gas_used) AS max_gas_in_block
FROM blocks
WHERE time >= NOW() - INTERVAL '30 days'
GROUP BY DATE(time)
ORDER BY block_date DESC;

4.2 交易分析

-- 4.2.1 日交易量趋势(90天)
SELECT
    DATE(t.created_at) AS tx_date,
    COUNT(*) AS tx_count,
    COUNT(DISTINCT t.sender) AS unique_senders,
    COUNT(DISTINCT t.contract) AS unique_contracts,
    SUM(t.gas_used) AS total_gas,
    AVG(t.gas_used)::INTEGER AS avg_gas_per_tx,
    COALESCE(SUM(t.fee_amount), 0) AS total_fees,
    AVG(t.fee_amount)::NUMERIC(20, 2) AS avg_fee_per_tx,
    COUNT(*) FILTER (WHERE t.status = 'FAILED') AS failed_count,
    ROUND(100.0 * COUNT(*) FILTER (WHERE t.status = 'SUCCESS') / NULLIF(COUNT(*), 0), 2) AS success_rate
FROM transactions t
WHERE t.created_at >= NOW() - INTERVAL '90 days'
GROUP BY DATE(t.created_at)
ORDER BY tx_date DESC;

-- 4.2.2 最活跃发送者排名
SELECT
    t.sender,
    COUNT(*) AS tx_count,
    COUNT(DISTINCT t.contract) AS contracts_interacted,
    COUNT(DISTINCT t.action) AS unique_actions,
    SUM(t.gas_used) AS total_gas_used,
    SUM(t.fee_amount) AS total_fees_paid,
    MIN(b.time)::date AS first_tx_date,
    MAX(b.time)::date AS last_tx_date,
    COUNT(*) FILTER (WHERE t.status = 'FAILED') AS failed_count
FROM transactions t
JOIN blocks b ON b.height = t.block_height
WHERE t.sender != ''
GROUP BY t.sender
ORDER BY tx_count DESC
LIMIT 20;

-- 4.2.3 交易动作分布
SELECT
    t.action,
    COUNT(*) AS tx_count,
    COUNT(DISTINCT t.sender) AS unique_senders,
    SUM(t.gas_used) AS total_gas,
    COALESCE(SUM(t.fee_amount), 0) AS total_fees,
    MIN(b.time)::date AS first_seen,
    MAX(b.time)::date AS last_seen
FROM transactions t
JOIN blocks b ON b.height = t.block_height
WHERE t.action != ''
GROUP BY t.action
ORDER BY tx_count DESC;

-- 4.2.4 小时交易热力图(30天)
SELECT
    EXTRACT(DOW FROM b.time) AS day_of_week,
    EXTRACT(HOUR FROM b.time) AS hour_of_day,
    COUNT(*) AS tx_count,
    COUNT(DISTINCT t.sender) AS unique_senders,
    SUM(t.gas_used) AS total_gas
FROM transactions t
JOIN blocks b ON b.height = t.block_height
WHERE t.created_at >= NOW() - INTERVAL '30 days'
GROUP BY day_of_week, hour_of_day
ORDER BY day_of_week, hour_of_day;

4.3 Gas 分析

-- 4.3.1 Gas 价格分布(7天)
SELECT
    CASE
        WHEN t.fee_amount = 0 THEN 'free'
        WHEN t.fee_amount::NUMERIC / NULLIF(t.gas_used, 1) < 0.001 THEN 'very_low'
        WHEN t.fee_amount::NUMERIC / NULLIF(t.gas_used, 1) < 0.01 THEN 'low'
        WHEN t.fee_amount::NUMERIC / NULLIF(t.gas_used, 1) < 0.025 THEN 'medium'
        WHEN t.fee_amount::NUMERIC / NULLIF(t.gas_used, 1) < 0.05 THEN 'high'
        ELSE 'very_high'
    END AS gas_price_tier,
    COUNT(*) AS tx_count,
    ROUND(AVG(t.fee_amount::NUMERIC / NULLIF(t.gas_used, 1)), 6) AS avg_gas_price,
    SUM(t.gas_used) AS total_gas_used,
    COUNT(DISTINCT t.sender) AS unique_senders
FROM transactions t
WHERE t.created_at >= NOW() - INTERVAL '7 days' AND t.gas_used > 0
GROUP BY gas_price_tier
ORDER BY avg_gas_price;

-- 4.3.2 合约 Gas 消耗排名(30天)
SELECT
    t.contract AS contract_address,
    c.label AS contract_name,
    c.contract_type,
    COUNT(*) AS tx_count,
    SUM(t.gas_used) AS total_gas_consumed,
    AVG(t.gas_used)::INTEGER AS avg_gas_per_tx,
    COALESCE(SUM(t.fee_amount), 0) AS total_fees,
    AVG(t.gas_used)::NUMERIC / NULLIF(AVG(t.gas_wanted), 0) AS avg_gas_efficiency,
    COUNT(DISTINCT t.sender) AS unique_users
FROM transactions t
LEFT JOIN contracts c ON c.address = t.contract
WHERE t.contract != '' AND t.created_at >= NOW() - INTERVAL '30 days'
GROUP BY t.contract, c.label, c.contract_type
ORDER BY total_gas_consumed DESC
LIMIT 20;

-- 4.3.3 Gas 日趋势(90天)
SELECT
    DATE(t.created_at) AS tx_date,
    COUNT(*) AS total_tx,
    SUM(gas_used) AS total_gas_used,
    AVG(gas_used)::INTEGER AS avg_gas_per_tx,
    PERCENTILE_CONT(0.5) WITHIN GROUP (ORDER BY gas_used) AS median_gas_per_tx,
    MAX(gas_used) AS max_gas_per_tx,
    SUM(t.fee_amount) AS total_fees,
    SUM(t.fee_amount)::NUMERIC / NULLIF(SUM(t.gas_used), 0) AS implied_avg_gas_price
FROM transactions t
WHERE t.created_at >= NOW() - INTERVAL '90 days'
GROUP BY DATE(t.created_at)
ORDER BY tx_date DESC;

4.4 Agent 经济分析

-- 4.4.1 Agent 综合排名
SELECT
    a.agent_id, a.name, a.owner, a.status, a.agent_type,
    COUNT(DISTINCT s.session_id) AS total_sessions,
    COUNT(DISTINCT s.client) AS unique_clients,
    COALESCE(SUM(s.amount), 0) AS total_volume_umsg,
    COALESCE(AVG(s.amount) FILTER (WHERE s.status = 'RELEASED'), 0)::NUMERIC(30, 0) AS avg_session_value,
    COUNT(*) FILTER (WHERE s.status = 'RELEASED') AS completed_sessions,
    COUNT(*) FILTER (WHERE s.status = 'OPEN') AS active_sessions,
    COUNT(*) FILTER (WHERE s.status = 'DISPUTED') AS disputed_sessions,
    ROUND(
        100.0 * COUNT(*) FILTER (WHERE s.status IN ('RELEASED', 'CLOSED'))
        / NULLIF(COUNT(*) FILTER (WHERE s.status NOT IN ('OPEN', 'FUNDED')), 0), 2
    ) AS completion_rate,
    MAX(s.created_at) AS last_activity
FROM agents a
LEFT JOIN sessions s ON s.agent_id = a.agent_id
GROUP BY a.agent_id, a.name, a.owner, a.status, a.agent_type
ORDER BY total_volume_umsg DESC
LIMIT 20;

-- 4.4.2 Agent 会话转化漏斗
SELECT
    agent_id, name,
    COUNT(*) AS total_created,
    COUNT(*) FILTER (WHERE status = 'FUNDED') AS funded,
    ROUND(100.0 * COUNT(*) FILTER (WHERE status = 'FUNDED') / NULLIF(COUNT(*), 0), 2) AS funded_pct,
    COUNT(*) FILTER (WHERE status = 'RELEASED') AS released,
    ROUND(100.0 * COUNT(*) FILTER (WHERE status = 'RELEASED') / NULLIF(COUNT(*), 0), 2) AS released_pct,
    COUNT(*) FILTER (WHERE status = 'DISPUTED') AS disputed,
    COUNT(*) FILTER (WHERE status = 'CLOSED') AS closed
FROM sessions
GROUP BY agent_id, name
ORDER BY total_created DESC
LIMIT 20;

-- 4.4.3 Agent 日收入趋势
SELECT
    ads.date, ads.agent_id, a.name,
    ads.tx_count, ads.session_count,
    ads.total_earned, ads.unique_clients
FROM agent_daily_stats ads
JOIN agents a ON a.agent_id = ads.agent_id
WHERE ads.date >= NOW() - INTERVAL '30 days'
ORDER BY ads.date DESC, ads.total_earned DESC;

4.5 用户分析

-- 4.5.1 DAU / WAU / MAU 趋势(90天)
WITH daily_users AS (
    SELECT
        DATE(created_at) AS dt,
        COUNT(DISTINCT sender) AS dau
    FROM transactions
    WHERE status = 'SUCCESS'
    GROUP BY DATE(created_at)
)
SELECT
    dt, dau,
    ROUND(AVG(dau) OVER (ORDER BY dt ROWS 6 PRECEDING), 2) AS wau_avg,
    ROUND(AVG(dau) OVER (ORDER BY dt ROWS 29 PRECEDING), 2) AS mau_avg,
    ROUND(100.0 * (dau - LAG(dau, 1) OVER (ORDER BY dt))
          / NULLIF(LAG(dau, 1) OVER (ORDER BY dt), 0), 2) AS dau_change_pct
FROM daily_users
WHERE dt >= NOW() - INTERVAL '90 days'
ORDER BY dt DESC;

-- 4.5.2 用户留存(按周同期群)
WITH first_activity AS (
    SELECT sender, MIN(DATE(created_at)) AS first_date
    FROM transactions WHERE status = 'SUCCESS'
    GROUP BY sender
),
weekly_cohort AS (
    SELECT
        DATE_TRUNC('week', fa.first_date) AS cohort_week,
        DATE_TRUNC('week', t.created_at) AS activity_week,
        COUNT(DISTINCT t.sender) AS active_users
    FROM first_activity fa
    JOIN transactions t ON t.sender = fa.sender AND t.status = 'SUCCESS'
    GROUP BY cohort_week, activity_week
)
SELECT
    cohort_week, activity_week, active_users,
    FIRST_VALUE(active_users) OVER (PARTITION BY cohort_week ORDER BY activity_week) AS week_0_users,
    ROUND(100.0 * active_users / NULLIF(
        FIRST_VALUE(active_users) OVER (PARTITION BY cohort_week ORDER BY activity_week), 0), 2
    ) AS retention_pct,
    EXTRACT(WEEK FROM activity_week) - EXTRACT(WEEK FROM cohort_week) AS week_number
FROM weekly_cohort
WHERE cohort_week >= NOW() - INTERVAL '180 days'
ORDER BY cohort_week DESC, activity_week;

-- 4.5.3 用户活跃度分段
SELECT
    sender,
    COUNT(*) AS total_txs,
    COUNT(DISTINCT DATE(created_at)) AS active_days,
    COUNT(DISTINCT contract) AS contracts_used,
    SUM(gas_used) AS total_gas,
    SUM(fee_amount) AS total_fees,
    MIN(created_at)::date AS first_active,
    MAX(created_at)::date AS last_active,
    CASE
        WHEN MAX(created_at) >= NOW() - INTERVAL '7 days' THEN 'active_7d'
        WHEN MAX(created_at) >= NOW() - INTERVAL '30 days' THEN 'active_30d'
        ELSE 'dormant'
    END AS user_segment
FROM transactions
WHERE status = 'SUCCESS' AND sender != ''
GROUP BY sender
ORDER BY total_txs DESC
LIMIT 50;

4.6 合约分析

-- 4.6.1 合约全局行为摘要
SELECT
    t.contract, c.label, c.contract_type,
    c.creator, c.created_at AS deploy_time,
    COUNT(*) AS total_interactions,
    COUNT(DISTINCT t.sender) AS total_users,
    COUNT(DISTINCT t.action) AS actions_supported,
    MIN(b.time)::date AS first_interaction,
    MAX(b.time)::date AS last_interaction,
    AVG(t.gas_used)::INTEGER AS avg_gas_per_tx,
    SUM(t.fee_amount) AS total_fees_generated
FROM transactions t
JOIN blocks b ON b.height = t.block_height
LEFT JOIN contracts c ON c.address = t.contract
WHERE t.contract != ''
GROUP BY t.contract, c.label, c.contract_type, c.creator, c.created_at
ORDER BY total_interactions DESC
LIMIT 30;

-- 4.6.2 合约日交互趋势(30天)
SELECT
    DATE(t.created_at) AS dt,
    t.contract, COUNT(*) AS interactions,
    COUNT(DISTINCT t.sender) AS unique_users,
    SUM(t.gas_used) AS gas_used,
    SUM(t.fee_amount) AS fees
FROM transactions t
WHERE t.contract != '' AND t.created_at >= NOW() - INTERVAL '30 days'
GROUP BY DATE(t.created_at), t.contract
ORDER BY dt DESC, interactions DESC;

-- 4.6.3 合约事件类型分布
SELECT
    e.contract_address, e.action,
    COUNT(*) AS event_count,
    COUNT(DISTINCT e.tx_hash) AS unique_tx_count,
    MIN(b.time)::date AS first_seen,
    MAX(b.time)::date AS last_seen
FROM events e
JOIN blocks b ON b.height = e.block_height
WHERE e.contract_address != ''
GROUP BY e.contract_address, e.action
ORDER BY e.contract_address, event_count DESC;

4.7 高级时间序列分析

-- 4.7.1 7 日和 30 日移动平均交易量
SELECT
    dt, tx_count,
    ROUND(AVG(tx_count) OVER (ORDER BY dt ROWS 6 PRECEDING), 2) AS sma_7d,
    ROUND(AVG(tx_count) OVER (ORDER BY dt ROWS 29 PRECEDING), 2) AS sma_30d,
    ROUND(STDDEV_SAMP(tx_count) OVER (ORDER BY dt ROWS 29 PRECEDING), 2) AS volatility_30d
FROM (
    SELECT DATE(created_at) AS dt, COUNT(*) AS tx_count
    FROM transactions WHERE created_at >= NOW() - INTERVAL '180 days'
    GROUP BY DATE(created_at)
) daily
ORDER BY dt;

-- 4.7.2 Gas 价格百分位数趋势
SELECT
    DATE(created_at) AS dt,
    COUNT(*) AS total_tx,
    ROUND(PERCENTILE_CONT(0.1) WITHIN GROUP (ORDER BY fee_amount::NUMERIC / NULLIF(gas_used, 1))::NUMERIC, 8) AS p10,
    ROUND(PERCENTILE_CONT(0.25) WITHIN GROUP (ORDER BY fee_amount::NUMERIC / NULLIF(gas_used, 1))::NUMERIC, 8) AS p25,
    ROUND(PERCENTILE_CONT(0.5) WITHIN GROUP (ORDER BY fee_amount::NUMERIC / NULLIF(gas_used, 1))::NUMERIC, 8) AS p50,
    ROUND(PERCENTILE_CONT(0.75) WITHIN GROUP (ORDER BY fee_amount::NUMERIC / NULLIF(gas_used, 1))::NUMERIC, 8) AS p75,
    ROUND(PERCENTILE_CONT(0.9) WITHIN GROUP (ORDER BY fee_amount::NUMERIC / NULLIF(gas_used, 1))::NUMERIC, 8) AS p90,
    ROUND(PERCENTILE_CONT(0.99) WITHIN GROUP (ORDER BY fee_amount::NUMERIC / NULLIF(gas_used, 1))::NUMERIC, 8) AS p99
FROM transactions
WHERE created_at >= NOW() - INTERVAL '90 days' AND gas_used > 0
GROUP BY DATE(created_at)
ORDER BY dt DESC;

-- 4.7.3 累计统计
SELECT
    DATE(created_at) AS dt,
    SUM(COUNT(*)) OVER (ORDER BY DATE(created_at)) AS cumulative_txs,
    SUM(COUNT(DISTINCT sender)) OVER (ORDER BY DATE(created_at)) AS cumulative_addresses,
    SUM(SUM(gas_used)) OVER (ORDER BY DATE(created_at)) AS cumulative_gas,
    SUM(SUM(fee_amount)) OVER (ORDER BY DATE(created_at)) AS cumulative_fees
FROM transactions
WHERE created_at >= NOW() - INTERVAL '365 days'
GROUP BY DATE(created_at)
ORDER BY dt;

4.8 Python 查询执行器

# indexer/analytics_queries.py
"""SQL 查询执行器 — Python 环境中执行分析查询"""

from datetime import datetime, timedelta
from typing import Any

from indexer.database import Database


class AnalyticsQueries:
    """分析查询执行器 — 封装常用分析查询"""

    def __init__(self, db: Database):
        self.db = db

    async def get_overview(self) -> dict[str, Any]:
        """获取网络概览"""
        sql = """
        SELECT
            (SELECT MAX(height) FROM blocks) AS latest_block_height,
            (SELECT COUNT(*) FROM blocks) AS total_blocks,
            (SELECT COUNT(*) FROM transactions) AS total_transactions,
            (SELECT COUNT(*) FROM events) AS total_events,
            (SELECT COUNT(*) FROM agents WHERE status = 'ACTIVE') AS active_agents,
            (SELECT COUNT(*) FROM sessions WHERE status IN ('OPEN', 'FUNDED')) AS active_sessions,
            (SELECT COUNT(*) FROM transactions WHERE created_at >= NOW() - INTERVAL '24 hours') AS txs_24h,
            (SELECT COUNT(DISTINCT sender) FROM transactions WHERE created_at >= NOW() - INTERVAL '24 hours') AS users_24h
        """
        rows = await self.db.execute_raw(sql)
        return dict(rows[0]) if rows else {}

    async def get_daily_tx_volume(self, days: int = 30) -> list[dict]:
        sql = """
        SELECT DATE(created_at) AS tx_date, COUNT(*) AS tx_count,
               COUNT(DISTINCT sender) AS unique_senders, SUM(gas_used) AS total_gas,
               COALESCE(SUM(fee_amount), 0) AS total_fees
        FROM transactions
        WHERE created_at >= NOW() - $1::INTERVAL
        GROUP BY DATE(created_at)
        ORDER BY tx_date ASC
        """
        rows = await self.db.execute_raw(sql, timedelta(days=days))
        return [dict(r) for r in rows]

    async def get_top_contracts(self, limit: int = 20, days: int = 30) -> list[dict]:
        sql = """
        SELECT t.contract, c.label, c.contract_type,
               COUNT(*) AS tx_count, COUNT(DISTINCT t.sender) AS unique_users,
               SUM(t.gas_used) AS total_gas, COALESCE(SUM(t.fee_amount), 0) AS total_fees,
               COUNT(*) FILTER (WHERE t.status = 'FAILED') AS failed_count
        FROM transactions t
        LEFT JOIN contracts c ON c.address = t.contract
        WHERE t.contract != '' AND t.created_at >= NOW() - $1::INTERVAL
        GROUP BY t.contract, c.label, c.contract_type
        ORDER BY tx_count DESC LIMIT $2
        """
        rows = await self.db.execute_raw(sql, timedelta(days=days), limit)
        return [dict(r) for r in rows]

    async def get_dau(self, days: int = 90) -> list[dict]:
        sql = """
        SELECT DATE(created_at) AS dt, COUNT(DISTINCT sender) AS dau, COUNT(*) AS total_txs
        FROM transactions WHERE status = 'SUCCESS' AND created_at >= NOW() - $1::INTERVAL
        GROUP BY DATE(created_at) ORDER BY dt ASC
        """
        rows = await self.db.execute_raw(sql, timedelta(days=days))
        return [dict(r) for r in rows]

    async def get_gas_price_trend(self, days: int = 30) -> list[dict]:
        sql = """
        SELECT DATE(created_at) AS dt, COUNT(*) AS count,
               ROUND(AVG(fee_amount::NUMERIC / NULLIF(gas_used, 1)), 8) AS avg_gas_price,
               ROUND(PERCENTILE_CONT(0.5) WITHIN GROUP (ORDER BY fee_amount::NUMERIC / NULLIF(gas_used, 1))::NUMERIC, 8) AS median_gas_price,
               ROUND(PERCENTILE_CONT(0.95) WITHIN GROUP (ORDER BY fee_amount::NUMERIC / NULLIF(gas_used, 1))::NUMERIC, 8) AS p95_gas_price
        FROM transactions
        WHERE created_at >= NOW() - $1::INTERVAL AND gas_used > 0 AND fee_amount > 0
        GROUP BY DATE(created_at) ORDER BY dt ASC
        """
        rows = await self.db.execute_raw(sql, timedelta(days=days))
        return [dict(r) for r in rows]

    async def get_agent_rankings(self, limit: int = 20) -> list[dict]:
        sql = """
        SELECT a.agent_id, a.name, a.status, a.agent_type,
               COUNT(DISTINCT s.session_id) AS total_sessions,
               COUNT(DISTINCT s.client) AS unique_clients,
               COALESCE(SUM(s.amount), 0) AS total_volume,
               COUNT(*) FILTER (WHERE s.status = 'DISPUTED') AS disputed_sessions,
               MAX(s.created_at) AS last_active
        FROM agents a
        LEFT JOIN sessions s ON s.agent_id = a.agent_id
        GROUP BY a.agent_id, a.name, a.status, a.agent_type
        ORDER BY total_volume DESC LIMIT $1
        """
        rows = await self.db.execute_raw(sql, limit)
        return [dict(r) for r in rows]

    async def get_event_distribution(self, days: int = 30) -> list[dict]:
        sql = """
        SELECT e.action, COUNT(*) AS count,
               COUNT(DISTINCT e.tx_hash) AS unique_txs,
               COUNT(DISTINCT e.contract_address) AS unique_contracts
        FROM events e
        WHERE e.created_at >= NOW() - $1::INTERVAL
        GROUP BY e.action ORDER BY count DESC
        """
        rows = await self.db.execute_raw(sql, timedelta(days=days))
        return [dict(r) for r in rows]

    async def custom_query(self, sql: str, **params) -> list[dict]:
        """自定义 SQL 查询"""
        rows = await self.db.execute_raw(sql, *params.values()) if params else await self.db.execute_raw(sql)
        return [dict(r) for r in rows]

5. 仪表盘 API

5.1 FastAPI 服务

# api/main.py
"""FastAPI 仪表盘 API 服务 — REST 端点"""

from datetime import datetime, timedelta
from typing import Any, Optional

from fastapi import FastAPI, HTTPException, Query, Depends
from fastapi.middleware.cors import CORSMiddleware

from indexer.database import Database
from indexer.analytics_queries import AnalyticsQueries

app = FastAPI(title="MSG Chain Analytics API", version="1.0.0")

app.add_middleware(
    CORSMiddleware,
    allow_origins=["*"],
    allow_credentials=True,
    allow_methods=["*"],
    allow_headers=["*"],
)

_db: Database | None = None
_queries: AnalyticsQueries | None = None


async def get_db() -> Database:
    if _db is None:
        raise HTTPException(status_code=503, detail="Database not initialized")
    return _db


async def get_queries() -> AnalyticsQueries:
    if _queries is None:
        raise HTTPException(status_code=503, detail="AnalyticsQueries not initialized")
    return _queries


@app.on_event("startup")
async def startup():
    global _db, _queries
    import os
    dsn = os.environ.get(
        "DATABASE_URL",
        "postgresql://msgchain:msgchain_secret@localhost:5432/msgchain_analytics",
    )
    _db = Database(dsn)
    await _db.connect()
    _queries = AnalyticsQueries(_db)


@app.on_event("shutdown")
async def shutdown():
    if _db:
        await _db.close()


# ── 健康检查 ──

@app.get("/api/v1/health")
async def health():
    return {"status": "ok", "timestamp": datetime.utcnow().isoformat(), "chain_id": "msg-chain-1"}


# ── 概览端点 ──

@app.get("/api/v1/dashboard/overview")
async def get_overview(queries: AnalyticsQueries = Depends(get_queries)):
    """仪表盘概览 — 核心指标"""
    data = await queries.get_overview()
    return {"success": True, "data": data, "timestamp": datetime.utcnow().isoformat()}


@app.get("/api/v1/dashboard/tx-volume")
async def get_tx_volume(
    days: int = Query(30, ge=1, le=365),
    queries: AnalyticsQueries = Depends(get_queries),
):
    """交易量时间序列"""
    data = await queries.get_daily_tx_volume(days=days)
    return {
        "success": True,
        "data": [
            {"date": str(r["tx_date"]), "tx_count": r["tx_count"],
             "unique_senders": r["unique_senders"], "total_gas": r["total_gas"],
             "total_fees": str(r["total_fees"])}
            for r in data
        ],
    }


# ── 合约端点 ──

@app.get("/api/v1/dashboard/top-contracts")
async def get_top_contracts(
    limit: int = Query(20, ge=1, le=100),
    days: int = Query(30, ge=1, le=365),
    queries: AnalyticsQueries = Depends(get_queries),
):
    """最活跃合约排名"""
    data = await queries.get_top_contracts(limit=limit, days=days)
    return {"success": True, "data": data}


# ── Agent 端点 ──

@app.get("/api/v1/dashboard/top-agents")
async def get_top_agents(
    limit: int = Query(20, ge=1, le=100),
    queries: AnalyticsQueries = Depends(get_queries),
):
    """Top Agent 排名"""
    data = await queries.get_agent_rankings(limit=limit)
    return {"success": True, "data": data}


@app.get("/api/v1/dashboard/agent/{agent_id}")
async def get_agent_detail(
    agent_id: str,
    db: Database = Depends(get_db),
):
    """Agent 详情 — 基本信息 + 会话统计 + 近期事件"""
    agent = await db.pool.fetchrow("SELECT * FROM agents WHERE agent_id = $1", agent_id)
    if not agent:
        raise HTTPException(status_code=404, detail=f"Agent {agent_id} not found")

    stats = await db.pool.fetchrow("""
        SELECT COUNT(*) AS total,
               COUNT(*) FILTER (WHERE status = 'RELEASED') AS completed,
               COUNT(*) FILTER (WHERE status = 'OPEN') AS active,
               COUNT(*) FILTER (WHERE status = 'DISPUTED') AS disputed,
               COALESCE(SUM(amount), 0) AS total_volume,
               COUNT(DISTINCT client) AS unique_clients,
               MAX(created_at) AS last_active
        FROM sessions WHERE agent_id = $1
    """, agent_id)

    events = await db.pool.fetch("""
        SELECT e.*, b.time FROM events e
        JOIN blocks b ON b.height = e.block_height
        WHERE e.attributes->>'agent_id' = $1
        ORDER BY e.block_height DESC LIMIT 20
    """, agent_id)

    return {
        "success": True,
        "data": {
            "agent": dict(agent),
            "session_stats": dict(stats) if stats else {},
            "recent_events": [dict(r) for r in events],
        },
    }


# ── Gas 分析端点 ──

@app.get("/api/v1/dashboard/gas-trend")
async def get_gas_trend(
    days: int = Query(30, ge=1, le=365),
    queries: AnalyticsQueries = Depends(get_queries),
):
    """Gas 价格趋势"""
    data = await queries.get_gas_price_trend(days=days)
    return {"success": True, "data": data}


@app.get("/api/v1/dashboard/gas-distribution")
async def get_gas_distribution(queries: AnalyticsQueries = Depends(get_queries)):
    """Gas 价格分布(7天)"""
    sql = """
    SELECT
        CASE
            WHEN fee_amount = 0 THEN 'free'
            WHEN fee_amount::NUMERIC / NULLIF(gas_used, 1) < 0.001 THEN 'very_low'
            WHEN fee_amount::NUMERIC / NULLIF(gas_used, 1) < 0.01 THEN 'low'
            WHEN fee_amount::NUMERIC / NULLIF(gas_used, 1) < 0.025 THEN 'medium'
            WHEN fee_amount::NUMERIC / NULLIF(gas_used, 1) < 0.05 THEN 'high'
            ELSE 'very_high'
        END AS tier,
        COUNT(*) AS count,
        ROUND(AVG(fee_amount::NUMERIC / NULLIF(gas_used, 1))::NUMERIC, 6) AS avg_price
    FROM transactions
    WHERE created_at >= NOW() - INTERVAL '7 days' AND gas_used > 0
    GROUP BY tier ORDER BY avg_price
    """
    data = await queries.custom_query(sql)
    return {"success": True, "data": data}


# ── 用户分析端点 ──

@app.get("/api/v1/dashboard/dau")
async def get_dau(
    days: int = Query(90, ge=1, le=365),
    queries: AnalyticsQueries = Depends(get_queries),
):
    """日活跃用户趋势"""
    data = await queries.get_dau(days=days)
    return {"success": True, "data": data}


# ── 事件分析端点 ──

@app.get("/api/v1/dashboard/event-distribution")
async def get_event_distribution(
    days: int = Query(30, ge=1, le=365),
    queries: AnalyticsQueries = Depends(get_queries),
):
    """事件类型分布"""
    data = await queries.get_event_distribution(days=days)
    return {"success": True, "data": data}


# ── 区块/交易浏览端点 ──

@app.get("/api/v1/blocks")
async def get_blocks(
    limit: int = Query(20, ge=1, le=100),
    offset: int = Query(0, ge=0),
    db: Database = Depends(get_db),
):
    """获取区块列表"""
    blocks = await db.get_blocks(limit=limit, offset=offset)
    return {"success": True, "data": [dict(r) for r in blocks]}


@app.get("/api/v1/transactions/{tx_hash}")
async def get_transaction(tx_hash: str, db: Database = Depends(get_db)):
    """获取交易详情(含事件)"""
    tx = await db.get_tx_by_hash(tx_hash)
    if not tx:
        raise HTTPException(status_code=404, detail="Transaction not found")
    events = await db.pool.fetch("SELECT * FROM events WHERE tx_hash = $1 ORDER BY id", tx_hash)
    return {"success": True, "data": {"transaction": dict(tx), "events": [dict(e) for e in events]}}


# ── 自定义查询端点 ──

@app.post("/api/v1/query")
async def custom_query(body: dict, queries: AnalyticsQueries = Depends(get_queries)):
    """自定义 SQL 查询(仅允许 SELECT)"""
    sql = body.get("sql", "")
    if not sql.strip().upper().startswith("SELECT"):
        raise HTTPException(status_code=400, detail="Only SELECT allowed")
    try:
        data = await queries.custom_query(sql, **(body.get("params", {})))
        return {"success": True, "data": data, "row_count": len(data)}
    except Exception as e:
        raise HTTPException(status_code=400, detail=str(e))

5.2 GraphQL 集成 (Strawberry)

# api/graphql_schema.py
"""Strawberry GraphQL Schema — 可选的 GraphQL 接口"""

from datetime import datetime
from typing import List, Optional

import strawberry
from indexer.database import Database


@strawberry.type
class BlockType:
    height: int
    hash: str
    time: datetime
    proposer: str
    tx_count: int
    gas_used: int


@strawberry.type
class TransactionType:
    hash: str
    block_height: int
    status: str
    gas_used: int
    fee_amount: str
    sender: str
    contract: str
    action: str


@strawberry.type
class EventType:
    id: int
    tx_hash: str
    block_height: int
    action: str
    contract_address: str
    attributes: str  # JSON string
    created_at: datetime


@strawberry.type
class AgentType:
    agent_id: str
    name: str
    owner: str
    status: str
    agent_type: str
    total_sessions: int
    total_volume: str


@strawberry.type
class ChainSummary:
    latest_block_height: int
    total_blocks: int
    total_transactions: int
    total_events: int
    total_agents: int
    active_agents: int


@strawberry.type
class Query:
    @strawberry.field
    async def blocks(self, limit: int = 20, offset: int = 0,
                     db: Database = strawberry.dependency()) -> List[BlockType]:
        rows = await db.get_blocks(limit=limit, offset=offset)
        return [BlockType(**dict(r)) for r in rows]

    @strawberry.field
    async def transaction(self, hash: str,
                          db: Database = strawberry.dependency()) -> Optional[TransactionType]:
        row = await db.get_tx_by_hash(hash)
        return TransactionType(**dict(row)) if row else None

    @strawberry.field
    async def agents(self, limit: int = 20,
                     queries: AnalyticsQueries = strawberry.dependency()) -> List[AgentType]:
        data = await queries.get_agent_rankings(limit=limit)
        return [AgentType(**r) for r in data]

    @strawberry.field
    async def chain_summary(self, db: Database = strawberry.dependency()) -> ChainSummary:
        overview = await db.get_chain_overview()
        return ChainSummary(**overview)


schema = strawberry.Schema(query=Query)

5.3 TypeScript API 客户端

// frontend/src/api/client.ts
/* TypeScript API 客户端 — 类型安全 */

export interface DashboardOverview {
  latest_block_height: number;
  total_blocks: number;
  total_transactions: number;
  total_events: number;
  active_agents: number;
  total_agents: number;
  active_sessions: number;
  txs_24h: number;
  users_24h: number;
}

export interface DailyTxVolume {
  date: string;
  tx_count: number;
  unique_senders: number;
  total_gas: number;
  total_fees: string;
}

export interface ContractRanking {
  contract: string;
  label: string;
  contract_type: string;
  tx_count: number;
  unique_users: number;
  total_gas: number;
  total_fees: string;
}

export interface AgentRanking {
  agent_id: string;
  name: string;
  owner: string;
  status: string;
  agent_type: string;
  total_sessions: number;
  unique_clients: number;
  total_volume: string;
  disputed_sessions: number;
  last_active: string | null;
}

export interface AgentDetail {
  agent: Record<string, unknown>;
  session_stats: {
    total: number;
    completed: number;
    active: number;
    disputed: number;
    total_volume: string;
    unique_clients: number;
    last_active: string | null;
  };
  recent_events: Array<Record<string, unknown>>;
}

export interface DauData {
  dt: string;
  dau: number;
  total_txs: number;
}

export interface GasTrend {
  dt: string;
  avg_gas_price: string;
  median_gas_price: string;
  p95_gas_price: string;
}

class AnalyticsAPI {
  private baseUrl: string;

  constructor(baseUrl = "http://localhost:8000") {
    this.baseUrl = baseUrl;
  }

  private async fetch<T>(path: string): Promise<T> {
    const res = await fetch(`${this.baseUrl}${path}`);
    if (!res.ok) {
      throw new Error(`API error: ${res.status} ${res.statusText}`);
    }
    const json = await res.json();
    if (!json.success) {
      throw new Error(`API error: ${JSON.stringify(json)}`);
    }
    return json.data as T;
  }

  async getOverview(): Promise<DashboardOverview> {
    return this.fetch("/api/v1/dashboard/overview");
  }

  async getTxVolume(days = 30): Promise<DailyTxVolume[]> {
    return this.fetch(`/api/v1/dashboard/tx-volume?days=${days}`);
  }

  async getTopContracts(limit = 20, days = 30): Promise<ContractRanking[]> {
    return this.fetch(`/api/v1/dashboard/top-contracts?limit=${limit}&days=${days}`);
  }

  async getTopAgents(limit = 20): Promise<AgentRanking[]> {
    return this.fetch(`/api/v1/dashboard/top-agents?limit=${limit}`);
  }

  async getAgentDetail(agentId: string): Promise<AgentDetail> {
    return this.fetch(`/api/v1/dashboard/agent/${encodeURIComponent(agentId)}`);
  }

  async getDau(days = 90): Promise<DauData[]> {
    return this.fetch(`/api/v1/dashboard/dau?days=${days}`);
  }

  async getGasTrend(days = 30): Promise<GasTrend[]> {
    return this.fetch(`/api/v1/dashboard/gas-trend?days=${days}`);
  }

  async getGasDistribution(): Promise<Array<{ tier: string; count: number; avg_price: string }>> {
    return this.fetch("/api/v1/dashboard/gas-distribution");
  }

  async getEventDistribution(days = 30): Promise<Array<{ action: string; count: number }>> {
    return this.fetch(`/api/v1/dashboard/event-distribution?days=${days}`);
  }
}

export const analyticsApi = new AnalyticsAPI();

6. React 前端仪表盘

6.1 项目结构与配置

frontend/
├── package.json
├── tsconfig.json
├── vite.config.ts
├── index.html
└── src/
    ├── main.tsx
    ├── App.tsx
    ├── api/
    │   └── client.ts          # API 客户端(§5.3)
    ├── components/
    │   ├── Layout.tsx          # 侧栏 + 顶栏布局
    │   ├── StatCard.tsx        # KPI 卡片
    │   ├── ChartContainer.tsx  # 图表容器包装
    │   ├── ContractTable.tsx   # 合约排名表格
    │   ├── AgentTable.tsx      # Agent 排名表格
    │   ├── TransactionPanel.tsx# 交易详情面板
    │   └── LoadingSpinner.tsx  # 加载指示器
    ├── charts/
    │   ├── TxVolumeChart.tsx   # 交易量折线图
    │   ├── GasPriceChart.tsx   # Gas 价格趋势图
    │   ├── DauChart.tsx        # DAU 面积图
    │   ├── ContractChart.tsx   # 合约活动柱状图
    │   ├── EventPieChart.tsx   # 事件类型饼图
    │   └── GasDistributionChart.tsx # Gas 分布柱状图
    ├── pages/
    │   ├── Overview.tsx        # 概览页
    │   ├── Agents.tsx          # Agent 列表页
    │   ├── AgentDetail.tsx     # Agent 详情页
    │   ├── Transactions.tsx    # 交易浏览页
    │   ├── Contracts.tsx       # 合约分析页
    │   └── Gas.tsx             # Gas 分析页
    └── hooks/
        ├── useApi.ts           # 通用数据请求 Hook
        └── usePolling.ts       # 轮询 Hook

6.2 入口与基础组件

// frontend/src/main.tsx
import React from "react";
import ReactDOM from "react-dom/client";
import App from "./App";

ReactDOM.createRoot(document.getElementById("root")!).render(
  <React.StrictMode>
    <App />
  </React.StrictMode>
);
// frontend/src/App.tsx
import { BrowserRouter, Routes, Route, Navigate } from "react-router-dom";
import Layout from "./components/Layout";
import Overview from "./pages/Overview";
import Agents from "./pages/Agents";
import AgentDetail from "./pages/AgentDetail";
import Transactions from "./pages/Transactions";
import Contracts from "./pages/Contracts";
import Gas from "./pages/Gas";

export default function App() {
  return (
    <BrowserRouter>
      <Routes>
        <Route path="/" element={<Layout />}>
          <Route index element={<Overview />} />
          <Route path="agents" element={<Agents />} />
          <Route path="agents/:id" element={<AgentDetail />} />
          <Route path="transactions" element={<Transactions />} />
          <Route path="contracts" element={<Contracts />} />
          <Route path="gas" element={<Gas />} />
          <Route path="*" element={<Navigate to="/" replace />} />
        </Route>
      </Routes>
    </BrowserRouter>
  );
}
// frontend/src/components/Layout.tsx
import { Outlet, NavLink } from "react-router-dom";

const navItems = [
  { path: "/", label: "概览", icon: "📊" },
  { path: "/agents", label: "AI Agent", icon: "🤖" },
  { path: "/transactions", label: "交易", icon: "🔗" },
  { path: "/contracts", label: "合约", icon: "📝" },
  { path: "/gas", label: "Gas", icon: "⛽" },
];

export default function Layout() {
  return (
    <div style={{ display: "flex", minHeight: "100vh", background: "#f5f7fa" }}>
      <nav style={{ width: 220, background: "#1a1d29", color: "#fff", padding: "1rem" }}>
        <h2 style={{ fontSize: "1.1rem", marginBottom: "1.5rem", padding: "0 0.5rem" }}>
          MSG Chain Analytics
        </h2>
        {navItems.map(({ path, label, icon }) => (
          <NavLink
            key={path}
            to={path}
            end={path === "/"}
            style={({ isActive }) => ({
              display: "block",
              padding: "0.6rem 0.8rem",
              borderRadius: 8,
              color: isActive ? "#fff" : "#8b8fa3",
              background: isActive ? "#2d6cdf" : "transparent",
              textDecoration: "none",
              marginBottom: 4,
              fontSize: "0.9rem",
            })}
          >
            {icon} {label}
          </NavLink>
        ))}
      </nav>
      <main style={{ flex: 1, padding: "1.5rem", overflow: "auto" }}>
        <Outlet />
      </main>
    </div>
  );
}
// frontend/src/components/StatCard.tsx
interface StatCardProps {
  title: string;
  value: string | number;
  subtitle?: string;
  trend?: { direction: "up" | "down"; value: string };
  icon?: string;
}

export default function StatCard({ title, value, subtitle, trend, icon }: StatCardProps) {
  return (
    <div style={{ background: "#fff", borderRadius: 12, padding: "1.2rem", boxShadow: "0 1px 3px rgba(0,0,0,0.08)" }}>
      <div style={{ display: "flex", justifyContent: "space-between", marginBottom: 8 }}>
        <span style={{ color: "#6b7280", fontSize: "0.85rem" }}>{title}</span>
        {icon && <span style={{ fontSize: "1.5rem" }}>{icon}</span>}
      </div>
      <div style={{ fontSize: "1.8rem", fontWeight: 700, marginBottom: 4 }}>{value}</div>
      {subtitle && <div style={{ color: "#9ca3af", fontSize: "0.8rem" }}>{subtitle}</div>}
      {trend && (
        <div style={{ color: trend.direction === "up" ? "#10b981" : "#ef4444", fontSize: "0.8rem" }}>
          {trend.direction === "up" ? "↑" : "↓"} {trend.value}
        </div>
      )}
    </div>
  );
}

6.3 数据 Hooks

// frontend/src/hooks/useApi.ts
import { useState, useEffect, useCallback } from "react";

interface UseApiResult<T> {
  data: T | null;
  loading: boolean;
  error: Error | null;
  refetch: () => void;
}

export function useApi<T>(fetcher: () => Promise<T>, deps: unknown[] = []): UseApiResult<T> {
  const [data, setData] = useState<T | null>(null);
  const [loading, setLoading] = useState(true);
  const [error, setError] = useState<Error | null>(null);

  const refetch = useCallback(async () => {
    setLoading(true);
    setError(null);
    try {
      const result = await fetcher();
      setData(result);
    } catch (e) {
      setError(e instanceof Error ? e : new Error(String(e)));
    } finally {
      setLoading(false);
    }
  }, deps); // eslint-disable-line react-hooks/exhaustive-deps

  useEffect(() => { refetch(); }, [refetch]);

  return { data, loading, error, refetch };
}
// frontend/src/hooks/usePolling.ts
import { useEffect, useRef } from "react";

export function usePolling(callback: () => void, intervalMs: number, enabled = true) {
  const savedCallback = useRef(callback);

  useEffect(() => { savedCallback.current = callback; }, [callback]);

  useEffect(() => {
    if (!enabled) return;
    const id = setInterval(() => savedCallback.current(), intervalMs);
    return () => clearInterval(id);
  }, [intervalMs, enabled]);
}

6.4 图表组件 (Recharts)

// frontend/src/charts/TxVolumeChart.tsx
import {
  ResponsiveContainer, AreaChart, Area, XAxis, YAxis, CartesianGrid, Tooltip, Legend,
} from "recharts";
import type { DailyTxVolume } from "../api/client";
import ChartContainer from "../components/ChartContainer";

interface Props { data: DailyTxVolume[]; loading: boolean; }

export default function TxVolumeChart({ data, loading }: Props) {
  return (
    <ChartContainer title="每日交易量" loading={loading}>
      <ResponsiveContainer width="100%" height={350}>
        <AreaChart data={data}>
          <defs>
            <linearGradient id="txGrad" x1="0" y1="0" x2="0" y2="1">
              <stop offset="5%" stopColor="#2d6cdf" stopOpacity={0.3} />
              <stop offset="95%" stopColor="#2d6cdf" stopOpacity={0} />
            </linearGradient>
          </defs>
          <CartesianGrid strokeDasharray="3 3" stroke="#eee" />
          <XAxis dataKey="date" tick={{ fontSize: 12 }} />
          <YAxis tick={{ fontSize: 12 }} />
          <Tooltip />
          <Legend />
          <Area type="monotone" dataKey="tx_count" stroke="#2d6cdf" fill="url(#txGrad)" name="交易数" />
          <Area type="monotone" dataKey="unique_senders" stroke="#10b981" fill="none" name="独立发送者" />
        </AreaChart>
      </ResponsiveContainer>
    </ChartContainer>
  );
}
// frontend/src/charts/GasPriceChart.tsx
import {
  ResponsiveContainer, LineChart, Line, XAxis, YAxis, CartesianGrid, Tooltip, Legend,
} from "recharts";
import type { GasTrend } from "../api/client";
import ChartContainer from "../components/ChartContainer";

interface Props { data: GasTrend[]; loading: boolean; }

export default function GasPriceChart({ data, loading }: Props) {
  return (
    <ChartContainer title="Gas 价格趋势" loading={loading}>
      <ResponsiveContainer width="100%" height={350}>
        <LineChart data={data}>
          <CartesianGrid strokeDasharray="3 3" stroke="#eee" />
          <XAxis dataKey="dt" tick={{ fontSize: 12 }} />
          <YAxis tick={{ fontSize: 12 }} />
          <Tooltip />
          <Legend />
          <Line type="monotone" dataKey="avg_gas_price" stroke="#2d6cdf" name="均价" dot={false} strokeWidth={2} />
          <Line type="monotone" dataKey="median_gas_price" stroke="#f59e0b" name="中位数" dot={false} strokeWidth={2} />
          <Line type="monotone" dataKey="p95_gas_price" stroke="#ef4444" name="P95" dot={false} strokeWidth={2} />
        </LineChart>
      </ResponsiveContainer>
    </ChartContainer>
  );
}
// frontend/src/charts/DauChart.tsx
import {
  ResponsiveContainer, AreaChart, Area, XAxis, YAxis, CartesianGrid, Tooltip, Legend,
} from "recharts";
import type { DauData } from "../api/client";
import ChartContainer from "../components/ChartContainer";

interface Props { data: DauData[]; loading: boolean; }

export default function DauChart({ data, loading }: Props) {
  return (
    <ChartContainer title="日活跃用户 (DAU)" loading={loading}>
      <ResponsiveContainer width="100%" height={250}>
        <AreaChart data={data}>
          <defs>
            <linearGradient id="dauGrad" x1="0" y1="0" x2="0" y2="1">
              <stop offset="5%" stopColor="#10b981" stopOpacity={0.3} />
              <stop offset="95%" stopColor="#10b981" stopOpacity={0} />
            </linearGradient>
          </defs>
          <CartesianGrid strokeDasharray="3 3" stroke="#eee" />
          <XAxis dataKey="dt" tick={{ fontSize: 12 }} />
          <YAxis tick={{ fontSize: 12 }} />
          <Tooltip />
          <Area type="monotone" dataKey="dau" stroke="#10b981" fill="url(#dauGrad)" name="DAU" strokeWidth={2} />
        </AreaChart>
      </ResponsiveContainer>
    </ChartContainer>
  );
}
// frontend/src/charts/EventPieChart.tsx
import { ResponsiveContainer, PieChart, Pie, Cell, Tooltip, Legend } from "recharts";
import ChartContainer from "../components/ChartContainer";

const COLORS = ["#2d6cdf", "#10b981", "#f59e0b", "#ef4444", "#8b5cf6", "#ec4899", "#14b8a6", "#f97316"];

interface Props {
  data: Array<{ action: string; count: number }>;
  loading: boolean;
}

export default function EventPieChart({ data, loading }: Props) {
  return (
    <ChartContainer title="事件分布" loading={loading}>
      <ResponsiveContainer width="100%" height={300}>
        <PieChart>
          <Pie data={data} dataKey="count" nameKey="action" cx="50%" cy="50%" outerRadius={100} label>
            {data.map((_, idx) => <Cell key={idx} fill={COLORS[idx % COLORS.length]} />)}
          </Pie>
          <Tooltip />
          <Legend />
        </PieChart>
      </ResponsiveContainer>
    </ChartContainer>
  );
}
// frontend/src/components/ChartContainer.tsx
import { ReactNode } from "react";
import LoadingSpinner from "./LoadingSpinner";

interface Props { title: string; loading: boolean; children: ReactNode; }

export default function ChartContainer({ title, loading, children }: Props) {
  return (
    <div style={{ background: "#fff", borderRadius: 12, padding: "1.2rem", boxShadow: "0 1px 3px rgba(0,0,0,0.08)", marginBottom: "1.2rem" }}>
      <h3 style={{ fontSize: "1rem", marginBottom: "1rem", color: "#374151" }}>{title}</h3>
      {loading ? <LoadingSpinner /> : children}
    </div>
  );
}
// frontend/src/components/LoadingSpinner.tsx
export default function LoadingSpinner() {
  return (
    <div style={{ display: "flex", justifyContent: "center", padding: "3rem" }}>
      <div style={{
        width: 36, height: 36, border: "3px solid #e5e7eb",
        borderTopColor: "#2d6cdf", borderRadius: "50%",
        animation: "spin 0.8s linear infinite",
      }} />
      <style>{`@keyframes spin { to { transform: rotate(360deg); } }`}</style>
    </div>
  );
}

6.5 表格组件

// frontend/src/components/ContractTable.tsx
import type { ContractRanking } from "../api/client";

interface Props { data: ContractRanking[]; loading: boolean; }

export default function ContractTable({ data, loading }: Props) {
  if (loading) return <div>加载中...</div>;
  return (
    <div style={{ background: "#fff", borderRadius: 12, overflow: "hidden", boxShadow: "0 1px 3px rgba(0,0,0,0.08)" }}>
      <table style={{ width: "100%", borderCollapse: "collapse" }}>
        <thead>
          <tr style={{ background: "#f9fafb", textAlign: "left" }}>
            <th style={thStyle}>#</th>
            <th style={thStyle}>合约地址</th>
            <th style={thStyle}>类型</th>
            <th style={thStyle}>交易数</th>
            <th style={thStyle}>独立用户</th>
            <th style={thStyle}>总 Gas</th>
          </tr>
        </thead>
        <tbody>
          {data.map((c, i) => (
            <tr key={c.contract} style={{ borderBottom: "1px solid #f3f4f6" }}>
              <td style={tdStyle}>{i + 1}</td>
              <td style={{ ...tdStyle, fontFamily: "monospace", fontSize: "0.85rem" }}>
                {c.contract.slice(0, 20)}...
              </td>
              <td style={tdStyle}>{c.contract_type}</td>
              <td style={tdStyle}>{c.tx_count.toLocaleString()}</td>
              <td style={tdStyle}>{c.unique_users.toLocaleString()}</td>
              <td style={tdStyle}>{c.total_gas.toLocaleString()}</td>
            </tr>
          ))}
        </tbody>
      </table>
    </div>
  );
}

const thStyle: React.CSSProperties = { padding: "0.75rem 1rem", fontSize: "0.8rem", color: "#6b7280", fontWeight: 600 };
const tdStyle: React.CSSProperties = { padding: "0.75rem 1rem", fontSize: "0.9rem" };
// frontend/src/components/AgentTable.tsx
import { useNavigate } from "react-router-dom";
import type { AgentRanking } from "../api/client";

interface Props { data: AgentRanking[]; loading: boolean; }

export default function AgentTable({ data, loading }: Props) {
  const navigate = useNavigate();
  if (loading) return <div>加载中...</div>;
  return (
    <div style={{ background: "#fff", borderRadius: 12, overflow: "hidden", boxShadow: "0 1px 3px rgba(0,0,0,0.08)" }}>
      <table style={{ width: "100%", borderCollapse: "collapse" }}>
        <thead>
          <tr style={{ background: "#f9fafb", textAlign: "left" }}>
            <th style={thStyle}>Agent ID</th>
            <th style={thStyle}>名称</th>
            <th style={thStyle}>状态</th>
            <th style={thStyle}>会话数</th>
            <th style={thStyle}>交易量</th>
            <th style={thStyle}>纠纷</th>
            <th style={thStyle}>最后活跃</th>
          </tr>
        </thead>
        <tbody>
          {data.map((a) => (
            <tr key={a.agent_id} onClick={() => navigate(`/agents/${a.agent_id}`)}
                style={{ borderBottom: "1px solid #f3f4f6", cursor: "pointer" }}>
              <td style={{ ...tdStyle, fontFamily: "monospace", fontSize: "0.85rem" }}>
                {a.agent_id.slice(0, 16)}...
              </td>
              <td style={tdStyle}>{a.name || "-"}</td>
              <td style={tdStyle}>
                <span style={{
                  padding: "2px 8px", borderRadius: 999, fontSize: "0.75rem",
                  background: a.status === "ACTIVE" ? "#d1fae5" : a.status === "DISPUTED" ? "#fee2e2" : "#fef3c7",
                  color: a.status === "ACTIVE" ? "#065f46" : a.status === "DISPUTED" ? "#991b1b" : "#92400e",
                }}>{a.status}</span>
              </td>
              <td style={tdStyle}>{a.total_sessions}</td>
              <td style={tdStyle}>{Number(a.total_volume).toFixed(2)} MSG</td>
              <td style={tdStyle}>{a.disputed_sessions}</td>
              <td style={tdStyle}>
                {a.last_active
                  ? new Date(a.last_active).toLocaleDateString("zh-CN")
                  : "-"}
              </td>
            </tr>
          ))}
        </tbody>
      </table>
    </div>
  );
}

6.6 页面组件

// frontend/src/pages/Overview.tsx
import { analyticsApi, type DashboardOverview } from "../api/client";
import { useApi } from "../hooks/useApi";
import { usePolling } from "../hooks/usePolling";
import StatCard from "../components/StatCard";
import TxVolumeChart from "../charts/TxVolumeChart";
import GasPriceChart from "../charts/GasPriceChart";
import DauChart from "../charts/DauChart";
import EventPieChart from "../charts/EventPieChart";
import ContractTable from "../components/ContractTable";

export default function Overview() {
  const overview = useApi<DashboardOverview>(() => analyticsApi.getOverview());
  const txVolume = useApi(() => analyticsApi.getTxVolume(30));
  const gasTrend = useApi(() => analyticsApi.getGasTrend(7));
  const dau = useApi(() => analyticsApi.getDau(90));
  const topContracts = useApi(() => analyticsApi.getTopContracts(10, 30));
  const eventDist = useApi(() => analyticsApi.getEventDistribution(30));

  usePolling(() => overview.refetch(), 10_000);

  return (
    <div>
      <h1 style={{ fontSize: "1.5rem", marginBottom: "1.5rem" }}>网络概览</h1>

      {overview.data && (
        <div style={{ display: "grid", gridTemplateColumns: "repeat(auto-fill, minmax(200px, 1fr))", gap: "1rem", marginBottom: "1.5rem" }}>
          <StatCard title="最新区块" value={`#${overview.data.latest_block_height}`} subtitle="当前高度" icon="⛓️" />
          <StatCard title="总交易数" value={overview.data.total_transactions.toLocaleString()} icon="🔗" />
          <StatCard title="24h 交易" value={overview.data.txs_24h.toLocaleString()} icon="📈" />
          <StatCard title="活跃 Agent" value={overview.data.active_agents} subtitle={`/ ${overview.data.total_agents} 总计`} icon="🤖" />
          <StatCard title="活跃会话" value={overview.data.active_sessions} icon="💬" />
          <StatCard title="24h 用户" value={overview.data.users_24h.toLocaleString()} icon="👤" />
        </div>
      )}

      <div style={{ display: "grid", gridTemplateColumns: "2fr 1fr", gap: "1.2rem" }}>
        <TxVolumeChart data={txVolume.data ?? []} loading={txVolume.loading} />
        <EventPieChart data={eventDist.data ?? []} loading={eventDist.loading} />
      </div>

      <DauChart data={dau.data ?? []} loading={dau.loading} />

      <div style={{ display: "grid", gridTemplateColumns: "1fr 1fr", gap: "1.2rem" }}>
        <GasPriceChart data={gasTrend.data ?? []} loading={gasTrend.loading} />
        <div>
          <h3 style={{ fontSize: "1rem", marginBottom: "0.75rem" }}>Top 10 合约</h3>
          <ContractTable data={topContracts.data ?? []} loading={topContracts.loading} />
        </div>
      </div>
    </div>
  );
}
// frontend/src/pages/Agents.tsx
import { analyticsApi, type AgentRanking } from "../api/client";
import { useApi } from "../hooks/useApi";
import AgentTable from "../components/AgentTable";

export default function Agents() {
  const { data, loading } = useApi<AgentRanking[]>(() => analyticsApi.getTopAgents(50));
  return (
    <div>
      <h1 style={{ fontSize: "1.5rem", marginBottom: "1.5rem" }}>AI Agent 排名</h1>
      <AgentTable data={data ?? []} loading={loading} />
    </div>
  );
}
// frontend/src/pages/AgentDetail.tsx
import { useParams } from "react-router-dom";
import { analyticsApi, type AgentDetail as AgentDetailType } from "../api/client";
import { useApi } from "../hooks/useApi";

export default function AgentDetail() {
  const { id } = useParams<{ id: string }>();
  const { data, loading, error } = useApi<AgentDetailType>(
    () => analyticsApi.getAgentDetail(id!),
    [id]
  );

  if (loading) return <div>加载中...</div>;
  if (error) return <div style={{ color: "red" }}>加载失败: {error.message}</div>;
  if (!data) return <div>无数据</div>;

  const { agent, session_stats: ss, recent_events: events } = data;

  return (
    <div>
      <h1 style={{ fontSize: "1.5rem", marginBottom: "1rem" }}>
        Agent: {String(agent.name || agent.agent_id || "").slice(0, 30)}
      </h1>
      <div style={{ display: "grid", gridTemplateColumns: "repeat(auto-fill, minmax(180px, 1fr))", gap: "1rem", marginBottom: "1.5rem" }}>
        <StatCard title="状态" value={String(agent.status || "")} />
        <StatCard title="类型" value={String(agent.agent_type || "")} />
        <StatCard title="总会话" value={ss.total ?? 0} />
        <StatCard title="完成" value={ss.completed ?? 0} />
        <StatCard title="活跃中" value={ss.active ?? 0} />
        <StatCard title="纠纷" value={ss.disputed ?? 0} />
        <StatCard title="独立客户" value={ss.unique_clients ?? 0} />
        <StatCard title="总成交量" value={`${Number(ss.total_volume || 0).toFixed(2)} MSG`} />
      </div>

      <h2 style={{ fontSize: "1.1rem", margin: "1rem 0" }}>近期事件</h2>
      <pre style={{ background: "#f9fafb", padding: "1rem", borderRadius: 8, fontSize: "0.8rem", maxHeight: 400, overflow: "auto" }}>
        {JSON.stringify(events, null, 2)}
      </pre>
    </div>
  );
}
import StatCard from "../components/StatCard";
// frontend/src/pages/Transactions.tsx
import { useState, useEffect } from "react";
import { analyticsApi } from "../api/client";

interface TxSummary {
  hash: string;
  block_height: number;
  status: string;
  sender: string;
  contract: string;
  action: string;
  fee_amount: string;
  created_at: string;
}

export default function Transactions() {
  const [txs, setTxs] = useState<TxSummary[]>([]);
  const [loading, setLoading] = useState(true);

  useEffect(() => {
    setLoading(true);
    // 使用 custom_query 简化;生产环境应有专用端点
    analyticsApi.getTxVolume(1).then(() => {
      fetch("http://localhost:8000/api/v1/blocks?limit=5")
        .then((r) => r.json())
        .then((d) => {
          setTxs((d.data || []).flatMap((b: any) => b.transactions || []));
          setLoading(false);
        });
    });
  }, []);

  if (loading) return <div>加载中...</div>;

  return (
    <div>
      <h1 style={{ fontSize: "1.5rem", marginBottom: "1rem" }}>交易浏览</h1>
      <div style={{ background: "#fff", borderRadius: 12, overflow: "hidden", boxShadow: "0 1px 3px rgba(0,0,0,0.08)" }}>
        <table style={{ width: "100%", borderCollapse: "collapse" }}>
          <thead>
            <tr style={{ background: "#f9fafb", textAlign: "left" }}>
              <th style={thStyle}>Tx Hash</th>
              <th style={thStyle}>区块</th>
              <th style={thStyle}>发送者</th>
              <th style={thStyle}>合约</th>
              <th style={thStyle}>操作</th>
              <th style={thStyle}>状态</th>
              <th style={thStyle}>费用</th>
            </tr>
          </thead>
          <tbody>
            {txs.map((tx) => (
              <tr key={tx.hash} style={{ borderBottom: "1px solid #f3f4f6" }}>
                <td style={{ ...tdStyle, fontFamily: "monospace", fontSize: "0.85rem" }}>
                  {tx.hash.slice(0, 16)}...
                </td>
                <td style={tdStyle}>{tx.block_height}</td>
                <td style={{ ...tdStyle, fontFamily: "monospace", fontSize: "0.85rem" }}>
                  {tx.sender.slice(0, 12)}...
                </td>
                <td style={{ ...tdStyle, fontFamily: "monospace", fontSize: "0.85rem" }}>
                  {tx.contract?.slice(0, 12)}...
                </td>
                <td style={tdStyle}>{tx.action}</td>
                <td style={tdStyle}>
                  <span style={{ color: tx.status === "success" ? "#10b981" : "#ef4444" }}>
                    {tx.status}
                  </span>
                </td>
                <td style={tdStyle}>{Number(tx.fee_amount).toFixed(6)}</td>
              </tr>
            ))}
          </tbody>
        </table>
      </div>
    </div>
  );
}
// frontend/src/pages/Contracts.tsx
import { analyticsApi, type ContractRanking } from "../api/client";
import { useApi } from "../hooks/useApi";
import ContractTable from "../components/ContractTable";

export default function Contracts() {
  const { data, loading } = useApi<ContractRanking[]>(() => analyticsApi.getTopContracts(50, 30));
  return (
    <div>
      <h1 style={{ fontSize: "1.5rem", marginBottom: "1rem" }}>合约分析</h1>
      <ContractTable data={data ?? []} loading={loading} />
    </div>
  );
}
// frontend/src/pages/Gas.tsx
import { analyticsApi, type GasTrend } from "../api/client";
import { useApi } from "../hooks/useApi";
import GasPriceChart from "../charts/GasPriceChart";
import GasDistributionChart from "../charts/GasDistributionChart";

export default function Gas() {
  const trend = useApi<GasTrend[]>(() => analyticsApi.getGasTrend(30));
  const dist = useApi(() => analyticsApi.getGasDistribution());
  return (
    <div>
      <h1 style={{ fontSize: "1.5rem", marginBottom: "1rem" }}>Gas 分析</h1>
      <GasPriceChart data={trend.data ?? []} loading={trend.loading} />
      <GasDistributionChart data={dist.data ?? []} loading={dist.loading} />
    </div>
  );
}
// frontend/src/charts/GasDistributionChart.tsx
import {
  ResponsiveContainer, BarChart, Bar, XAxis, YAxis, CartesianGrid, Tooltip, Legend,
} from "recharts";
import ChartContainer from "../components/ChartContainer";

interface GasDistItem { tier: string; count: number; avg_price: string; }

interface Props { data: GasDistItem[]; loading: boolean; }

const TIER_LABELS: Record<string, string> = {
  free: "免费", very_low: "极低", low: "低", medium: "中等", high: "高", very_high: "极高",
};

export default function GasDistributionChart({ data, loading }: Props) {
  const mapped = data.map((d) => ({ ...d, tier: TIER_LABELS[d.tier] || d.tier }));
  return (
    <ChartContainer title="Gas 价格分布(7天)" loading={loading}>
      <ResponsiveContainer width="100%" height={300}>
        <BarChart data={mapped}>
          <CartesianGrid strokeDasharray="3 3" stroke="#eee" />
          <XAxis dataKey="tier" />
          <YAxis tick={{ fontSize: 12 }} />
          <Tooltip />
          <Legend />
          <Bar dataKey="count" fill="#2d6cdf" name="交易数" />
        </BarChart>
      </ResponsiveContainer>
    </ChartContainer>
  );
}

7. 部署与运维

7.1 Docker Compose 全栈部署

# docker-compose.yml
version: "3.9"

volumes:
  pgdata:
  redis_data:

services:
  postgres:
    image: postgres:16-alpine
    environment:
      POSTGRES_DB: msgchain_analytics
      POSTGRES_USER: msgchain
      POSTGRES_PASSWORD: msgchain_secret
    volumes:
      - pgdata:/var/lib/postgresql/data
      - ./sql/init.sql:/docker-entrypoint-initdb.d/01-init.sql
      - ./sql/migrations:/docker-entrypoint-initdb.d/migrations
    ports:
      - "5432:5432"
    healthcheck:
      test: ["CMD-SHELL", "pg_isready -U msgchain -d msgchain_analytics"]
      interval: 10s
      timeout: 5s
      retries: 5
    deploy:
      resources:
        limits:
          memory: 2G
        reservations:
          memory: 1G

  redis:
    image: redis:7-alpine
    ports:
      - "6379:6379"
    volumes:
      - redis_data:/data
    healthcheck:
      test: ["CMD", "redis-cli", "ping"]
      interval: 10s
      timeout: 5s
      retries: 3

  indexer:
    build:
      context: .
      dockerfile: Dockerfile.indexer
    depends_on:
      postgres:
        condition: service_healthy
      redis:
        condition: service_healthy
    environment:
      RPC_URL: https://rpc.msgchain.example.com
      WSS_URL: wss://rpc.msgchain.example.com/websocket
      DATABASE_URL: postgresql://msgchain:msgchain_secret@postgres:5432/msgchain_analytics
      REDIS_URL: redis://redis:6379/0
      INDEXER_START_HEIGHT: "0"
      INDEXER_POLL_INTERVAL: "2"
      INDEXER_BATCH_SIZE: "100"
      LOG_LEVEL: INFO
    restart: unless-stopped

  api:
    build:
      context: .
      dockerfile: Dockerfile.api
    depends_on:
      postgres:
        condition: service_healthy
    environment:
      DATABASE_URL: postgresql://msgchain:msgchain_secret@postgres:5432/msgchain_analytics
      CORS_ORIGINS: "*"
      LOG_LEVEL: INFO
    ports:
      - "8000:8000"
    restart: unless-stopped
    deploy:
      resources:
        limits:
          memory: 512M

  dashboard:
    build:
      context: ./frontend
      dockerfile: Dockerfile
    depends_on:
      - api
    environment:
      VITE_API_URL: http://localhost:8000
    ports:
      - "3000:80"
    restart: unless-stopped

  nginx:
    image: nginx:alpine
    volumes:
      - ./nginx.conf:/etc/nginx/nginx.conf:ro
    ports:
      - "80:80"
    depends_on:
      - api
      - dashboard
    restart: unless-stopped

  pgadmin:
    image: dpage/pgadmin4
    environment:
      PGADMIN_DEFAULT_EMAIL: admin@msgchain.local
      PGADMIN_DEFAULT_PASSWORD: admin_secret
    ports:
      - "5050:80"
    profiles:
      - dev
    restart: unless-stopped

  prometheus:
    image: prom/prometheus:v2.45
    volumes:
      - ./prometheus.yml:/etc/prometheus/prometheus.yml:ro
      - prometheus_data:/prometheus
    ports:
      - "9090:9090"
    profiles:
      - monitoring
    restart: unless-stopped

  grafana:
    image: grafana/grafana:10.0
    environment:
      GF_SECURITY_ADMIN_PASSWORD: grafana_secret
    volumes:
      - grafana_data:/var/lib/grafana
      - ./grafana/dashboards:/etc/grafana/provisioning/dashboards:ro
      - ./grafana/datasources:/etc/grafana/provisioning/datasources:ro
    ports:
      - "3001:3000"
    profiles:
      - monitoring
    restart: unless-stopped

volumes:
  prometheus_data:
  grafana_data:

7.2 Dockerfile 构建

# Dockerfile.indexer — 索引器
FROM python:3.12-slim AS builder
WORKDIR /app
COPY requirements.txt .
RUN pip install --no-cache-dir -r requirements.txt && \
    pip install --no-cache-dir aiohttp asyncpg redis websockets

FROM python:3.12-slim
WORKDIR /app
COPY --from=builder /usr/local/lib/python3.12/site-packages /usr/local/lib/python3.12/site-packages
COPY indexer/ ./indexer/
CMD ["python", "-m", "indexer.main"]
# Dockerfile.api — API 服务
FROM python:3.12-slim AS builder
WORKDIR /app
COPY api/requirements.txt .
RUN pip install --no-cache-dir -r requirements.txt
RUN pip install --no-cache-dir fastapi uvicorn strawberry-graphql asyncpg

FROM python:3.12-slim
WORKDIR /app
COPY --from=builder /usr/local/lib/python3.12/site-packages /usr/local/lib/python3.12/site-packages
COPY api/ ./api/
COPY indexer/ ./indexer/
CMD ["uvicorn", "api.main:app", "--host", "0.0.0.0", "--port", "8000", "--workers", "4"]
# frontend/Dockerfile — Vite 构建 + Nginx
FROM node:20-alpine AS builder
WORKDIR /app
COPY package.json .
RUN npm ci
COPY . .
RUN npm run build

FROM nginx:alpine
COPY --from=builder /app/dist /usr/share/nginx/html
COPY nginx-default.conf /etc/nginx/conf.d/default.conf
EXPOSE 80
CMD ["nginx", "-g", "daemon off;"]

7.3 Nginx 配置

# nginx.conf
events { worker_connections 1024; }

http {
  upstream api_upstream {
    server api:8000;
    keepalive 32;
  }

  upstream dashboard_upstream {
    server dashboard:80;
  }

  server {
    listen 80;
    server_name analytics.msgchain.example.com;

    # 静态资源缓存
    location /assets/ {
      proxy_pass http://dashboard_upstream;
      expires 1y;
      add_header Cache-Control "public, immutable";
    }

    # API 代理
    location /api/ {
      proxy_pass http://api_upstream;
      proxy_http_version 1.1;
      proxy_set_header Upgrade $http_upgrade;
      proxy_set_header Connection "upgrade";
      proxy_set_header Host $host;
      proxy_set_header X-Real-IP $remote_addr;
      proxy_read_timeout 60s;

      # 请求体限制(GraphQL 查询可能较大)
      client_max_body_size 1M;
    }

    # GraphQL 端点
    location /graphql {
      proxy_pass http://api_upstream;
      proxy_http_version 1.1;
      proxy_set_header Host $host;
      proxy_read_timeout 30s;
    }

    # 前端 SPA(所有其他路径)
    location / {
      proxy_pass http://dashboard_upstream;
      proxy_set_header Host $host;
      proxy_set_header X-Real-IP $remote_addr;
      expires -1;
    }

    # 压缩
    gzip on;
    gzip_types application/json text/plain text/css application/javascript;
    gzip_min_length 1000;
    gzip_proxied any;

    # 安全头
    add_header X-Content-Type-Options nosniff;
    add_header X-Frame-Options DENY;
    add_header X-XSS-Protection "1; mode=block";
  }
}

7.4 监控与告警

# prometheus.yml
global:
  scrape_interval: 15s
  evaluation_interval: 15s

scrape_configs:
  - job_name: "api"
    static_configs:
      - targets: ["api:8000"]
    metrics_path: /metrics

  - job_name: "indexer"
    static_configs:
      - targets: ["indexer:8080"]
    metrics_path: /metrics

  - job_name: "postgres"
    static_configs:
      - targets: ["postgres:9187"]  # postgres_exporter 侧车

  - job_name: "redis"
    static_configs:
      - targets: ["redis:9121"]     # redis_exporter 侧车

  - job_name: "node"
    static_configs:
      - targets: ["node-exporter:9100"]
# grafana/datasources/prometheus.yml
apiVersion: 1
datasources:
  - name: Prometheus
    type: prometheus
    url: http://prometheus:9090
    access: proxy
    isDefault: true
# grafana/dashboards/dashboard.yml
apiVersion: 1
providers:
  - name: "MSG Chain"
    orgId: 1
    folder: ""
    type: file
    disableDeletion: false
    editable: true
    options:
      path: /etc/grafana/provisioning/dashboards/json

7.5 性能调优清单

PostgreSQL 调优

# postgresql.custom.conf
# 专为分析型负载优化
shared_buffers = 512MB          # 内存的 25%
effective_cache_size = 1.5GB    # 内存的 50%
work_mem = 64MB                 # 排序/哈希内存
maintenance_work_mem = 256MB    # VACUUM/索引维护
random_page_cost = 1.1          # SSD 优化
effective_io_concurrency = 200  # SSD 并发
wal_buffers = 16MB
max_worker_processes = 8
max_parallel_workers_per_gather = 4
max_parallel_workers = 8
parallel_tuple_cost = 0.01
parallel_setup_cost = 25
jit = on
jit_above_cost = 100000         # 大查询启用 JIT

# 分析型查询调优
from_collapse_limit = 8
join_collapse_limit = 8

# 统计信息
default_statistics_target = 500

查询缓存策略

# 索引器内置缓存(+ Redis)
CACHE_DEFAULTS = {
    "overview": 5,             # 概览缓 5 秒
    "tx_volume": 60,           # 交易量缓 60 秒
    "dau": 300,                # DAU 缓 5 分钟
    "gas_trend": 60,
    "top_contracts": 120,
    "agent_rankings": 60,
    "agent_detail": 30,
    "event_distribution": 300,
}
// 前端 SWR 式缓存
// frontend/src/hooks/useApi.ts 已内置 loading/error 状态;
// 生产环境推荐使用 TanStack Query (react-query) 开启后台刷新和缓存 GC
// npm install @tanstack/react-query

// 示例集成:
// import { QueryClient, QueryClientProvider, useQuery } from "@tanstack/react-query";
// const queryClient = new QueryClient({ defaultOptions: { queries: { staleTime: 30_000 } } });
// function useOverview() {
//   return useQuery({ queryKey: ["overview"], queryFn: () => analyticsApi.getOverview(), refetchInterval: 10_000 });
// }

7.6 索引维护策略

-- 7.6.1 定期重建高写入表索引
REINDEX TABLE CONCURRENTLY transactions;
REINDEX TABLE CONCURRENTLY events;

-- 7.6.2 更新统计信息
ANALYZE transactions;
ANALYZE events;
ANALYZE daily_stats;
ANALYZE hourly_stats;

-- 7.6.3 清理过期数据(保留 90 天明细,永久保留聚合)
DELETE FROM events WHERE created_at < NOW() - INTERVAL '90 days';
DELETE FROM transactions WHERE created_at < NOW() - INTERVAL '90 days';
DELETE FROM blocks WHERE time < NOW() - INTERVAL '90 days';
DELETE FROM hourly_stats WHERE bucket < NOW() - INTERVAL '90 days';

7.7 Crontab 定期任务

# ── 数据保留 ──
0 3 * * * psql $DATABASE_URL -c "DELETE FROM events WHERE created_at < NOW() - INTERVAL '90 days'"
0 4 * * * psql $DATABASE_URL -c "DELETE FROM transactions WHERE created_at < NOW() - INTERVAL '90 days'"
0 5 * * * psql $DATABASE_URL -c "VACUUM ANALYZE transactions; VACUUM ANALYZE events;"

# ── 物化视图刷新 ──
*/30 * * * * psql $DATABASE_URL -c "REFRESH MATERIALIZED VIEW CONCURRENTLY mv_dau"
*/15 * * * * psql $DATABASE_URL -c "REFRESH MATERIALIZED VIEW CONCURRENTLY mv_contract_rankings"
*/10 * * * * psql $DATABASE_URL -c "REFRESH MATERIALIZED VIEW CONCURRENTLY mv_agent_economics"
0 * * * *   psql $DATABASE_URL -c "REFRESH MATERIALIZED VIEW CONCURRENTLY mv_gas_price_distribution"

# ── 数据一致性检查 ──
0 6 * * 0 psql $DATABASE_URL -f /sql/consistency_checks.sql > /var/log/consistency.log 2>&1

# ── 备份 ──
0 2 * * * pg_dump -U msgchain -h localhost msgchain_analytics -F c -f /backup/msgchain_$(date +\\%Y\\%m\\%d).dump

7.8 故障恢复流程

#!/bin/bash
# scripts/disaster_recovery.sh — MSG Chain 分析集群灾难恢复

set -euo pipefail
BACKUP_DIR="/backup"
LATEST=$(ls -t "$BACKUP_DIR"/msgchain_*.dump | head -1)

echo "[1/5] 停止服务..."
docker-compose stop api indexer dashboard

echo "[2/5] 恢复 PostgreSQL..."
docker-compose exec -T postgres pg_restore -U msgchain -d msgchain_analytics \
    --clean --if-exists "$LATEST"

echo "[3/5] 重建索引..."
docker-compose exec -T postgres psql -U msgchain -d msgchain_analytics -c "REINDEX DATABASE msgchain_analytics;"

echo "[4/5] 刷新物化视图..."
docker-compose exec -T postgres psql -U msgchain -d msgchain_analytics <<SQL
REFRESH MATERIALIZED VIEW CONCURRENTLY mv_dau;
REFRESH MATERIALIZED VIEW CONCURRENTLY mv_contract_rankings;
REFRESH MATERIALIZED VIEW CONCURRENTLY mv_agent_economics;
REFRESH MATERIALIZED VIEW CONCURRENTLY mv_gas_price_distribution;
SQL

echo "[5/5] 启动服务..."
docker-compose up -d
echo "恢复完成,最新备份:$LATEST"

7.9 安全配置

# 安全中间件(FastAPI)
from fastapi import Request
from starlette.middleware.base import BaseHTTPMiddleware
import time


class RateLimitMiddleware(BaseHTTPMiddleware):
    def __init__(self, app, rate_limit: int = 100, window: int = 60):
        super().__init__(app)
        self.rate_limit = rate_limit
        self.window = window
        self.requests: dict[str, list[float]] = {}

    async def dispatch(self, request: Request, call_next):
        # GraphQL 端点的请求速率限制
        if request.url.path in ("/graphql", "/api/v1/query"):
            client_ip = request.client.host if request.client else "unknown"
            now = time.time()
            window_start = now - self.window
            self.requests.setdefault(client_ip, [])
            self.requests[client_ip] = [t for t in self.requests[client_ip] if t > window_start]
            if len(self.requests[client_ip]) >= self.rate_limit:
                from fastapi.responses import JSONResponse
                return JSONResponse(
                    status_code=429,
                    content={"error": "Rate limit exceeded", "retry_after": self.window},
                )
            self.requests[client_ip].append(now)
        return await call_next(next)

7.10 可观测性

# api/metrics.py — Prometheus 指标暴露
from prometheus_client import Counter, Histogram, generate_latest, REGISTRY
from fastapi import Response

REQUEST_COUNT = Counter("api_requests_total", "Total API requests", ["method", "endpoint", "status"])
REQUEST_LATENCY = Histogram("api_request_duration_seconds", "API request latency", ["method", "endpoint"])


@app.get("/metrics")
async def metrics():
    return Response(content=generate_latest(REGISTRY), media_type="text/plain")
# 结构化日志(JSON 格式)
import structlog
import logging.config

logging.config.dictConfig({
    "version": 1,
    "formatters": {
        "json": {
            "()": structlog.stdlib.ProcessorFormatter,
            "processor": structlog.processors.JSONRenderer(),
        },
    },
    "handlers": {
        "console": {
            "class": "logging.StreamHandler",
            "formatter": "json",
        },
    },
    "loggers": {
        "indexer": {"level": "INFO", "handlers": ["console"]},
        "api": {"level": "INFO", "handlers": ["console"]},
    },
})

7.11 水平扩展方案

# docker-compose.prod.yml — 生产扩展
version: "3.9"
services:
  api:
    deploy:
      replicas: 4                # API 多副本
      resources:
        limits:
          memory: 1G
    environment:
      ENABLE_CACHE: "true"
      CACHE_TTL_OVERVIEW: "2"

  indexer:
    deploy:
      replicas: 3                # 分区索引(每个负责不同 block 范围)
    environment:
      INDEXER_SHARD: "${INDEXER_SHARD:-0}"
      INDEXER_SHARD_TOTAL: "${INDEXER_SHARD_TOTAL:-3}"

  postgres:
    deploy:
      placement:
        constraints:
          - node.role == manager
    # 生产环境建议使用托管 PostgreSQL(如 RDS、Cloud SQL)
    # 或使用 Patroni + pgpool-II 搭建高可用集群
# 分区索引实现
class ShardedIndexer(MsgChainIndexer):
    def __init__(self, shard: int, total_shards: int, *args, **kwargs):
        super().__init__(*args, **kwargs)
        self.shard = shard
        self.total_shards = total_shards

    def should_process(self, height: int) -> bool:
        return height % self.total_shards == self.shard

7.12 成本估算

资源 规格 月度成本 (USD) 说明
PostgreSQL 2 vCPU, 4GB RAM, 100GB SSD ~$50-80 托管数据库
Redis Cache 1GB 内存 ~$15-20 ElastiCache/Memorystore
API 服务器 2 vCPU, 4GB RAM × 2 ~$60-100 中等流量
索引器 1 vCPU, 2GB RAM × 2 ~$30-50 并行索引
对象存储 500GB (备份) ~$10-15 S3/GCS
监控 免费 (Prometheus + Grafana) $0 自托管
总计 ~$165-265/mo

7.13 上线检查清单

□ 数据库迁移已执行且验证
□ 初始回填已完成,区块高度一致
□ 实时索引器运行超 1 小时,无重连错误
□ API 健康检查返回 200
□ 前端仪表盘所有图表正常加载
□ CORS 配置限制为已知域名
□ Nginx 反向代理启用 HTTPS (Let's Encrypt)
□ Rate Limiting 已配置
□ Prometheus 目标均在 UP 状态
□ Grafana 仪表盘显示基线指标
□ 备份 cron 任务已设置并测试
□ 日志轮转已配置(Docker json-file max-size/max-file)
□ 故障恢复演练完成
□ 负载测试通过(建议 1000 QPS 持续 10 分钟)