dApp Docs/事件索引与GraphQL API
Development reference. Not independently verified for production.

MSG Chain 事件索引与 GraphQL API 完全指南

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

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


目录

  1. 概述
  2. Tendermint WebSocket 事件
  3. CosmWasm 合约事件
  4. Indexer 数据管道
  5. PostgreSQL 数据模型
  6. GraphQL API 服务
  7. Indexer 事件索引
  8. Agent API 事件网关
  9. 实时订阅前端
  10. 搜索服务
  11. Docker 部署
  12. 性能与扩展

1. 概述

1.1 为什么需要事件索引

MSG Chain 是一条基于 Cosmos SDK + CosmWasm 的应用链。链上合约(Agent Payment、Agent Registry、DID Registry、AI Constitution、Micropayment Session)会持续产生大量事件。原始链上事件通过 Tendermint WebSocket 暴露,但存在以下问题:

事件索引系统(Indexer)解决了这些问题:它订阅链上事件,解析、规范化后存入 PostgreSQL,并通过 GraphQL API 对外提供查询和实时订阅能力。

1.2 架构总览

┌─────────────────────────────────────────────────────────────────┐
│                          MSG Chain                              │
│  ┌───────────────────────────────────────────────────────────┐  │
│  │  Tendermint RPC (26657)  ◄── WebSocket                   │  │
│  │  CosmWasm Contract Events                                 │  │
│  └───────────────────────────────────────────────────────────┘  │
└─────────────────────────────────────────────────────────────────┘
                          │
                          ▼
┌─────────────────────────────────────────────────────────────────┐
│                    Indexer Service                               │
│  ┌──────────┐  ┌──────────┐  ┌─────────────┐  ┌────────────┐  │
│  │WS Listener│→│Event Parser│→│Normalizer   │→│DB Writer   │  │
│  └──────────┘  └──────────┘  └─────────────┘  └────────────┘  │
└─────────────────────────────────────────────────────────────────┘
                          │
                          ▼
┌─────────────────────────────────────────────────────────────────┐
│                        PostgreSQL                                │
│  blocks │ transactions │ events │ messages │ contracts          │
│  agents │ sessions │ did_events │ constitution_rules             │
└─────────────────────────────────────────────────────────────────┘
                          │
                          ▼
┌─────────────────────────────────────────────────────────────────┐
│                     GraphQL API (Yoga/Apollo)                    │
│  ┌──────────────────────────────────────────────────────────┐   │
│  │  Queries │ Mutations │ Subscriptions                    │   │
│  │  Search  │ Pagination │ DataLoader                      │   │
│  └──────────────────────────────────────────────────────────┘   │
└─────────────────────────────────────────────────────────────────┘
                          │
                          ▼
┌─────────────────────────────────────────────────────────────────┐
│                    Clients (Frontend / SDK)                      │
│  React Apollo │ Agent API WebSocket │ curl                      │
└─────────────────────────────────────────────────────────────────┘

1.3 当前实现状态

⚠️ 重要: 当前 Indexer API 处于 local_preflight_only=true 状态,仅部分实现,不适用于生产环境。

// X-Local-Only=true — 以下功能仅在本地预检模式下可用
const INDEXER_FLAGS = {
  localPreflightOnly: true,
  productionReady: false,
  persistentStorage: false,  // BadgerDB 本地存储
  retentionProof: false,     // 保留证明未实现
  fullTextSearch: false,     // 搜索仅限前缀匹配
};

已实现的 Indexer API 端点:

方法 端点 状态 说明
GET /api/v1/indexer/health ✅ 已实现 健康检查
GET /api/v1/indexer/capabilities ✅ 已实现 返回 local_preflight_only 标志
GET /api/v1/indexer/retention/proof ⚠️ 存根 保留证明,仅返回空响应
GET /api/v1/indexer/memos ⚠️ 部分实现 Memo 列表,仅支持 BadgerDB 游标分页
GET /api/v1/indexer/memos/search?q= ⚠️ 部分实现 仅前缀匹配,无全文搜索
GET /api/v1/indexer/memos/types ✅ 已实现 Memo 类型列表
POST /api/v1/indexer/memos/submit ⚠️ 部分实现 提交 memo,仅验证格式

1.4 技术栈

组件 技术选型 版本
链 Cosmos SDK + Tendermint + CosmWasm Cosmos SDK v0.47+
存储 PostgreSQL (主存储) + BadgerDB (Indexer 本地) PG 15+
GraphQL 服务器 GraphQL Yoga 或 Apollo Server Yoga v3+
TypeScript 运行时 Node.js 或 Bun Node 20+ / Bun 1.0+
缓存 Redis (可选) 7+
部署 Docker Compose 3.8+
前端 React 18 + Apollo Client 18.x

2. Tendermint WebSocket 事件

2.1 连接设置

Tendermint 通过 WebSocket 暴露事件,端点位于 ws://localhost:26657/websocket。连接使用 JSON-RPC 协议。

// src/websocket/tendermint-ws.ts
// Tendermint WebSocket 客户端 — 完整实现

import WebSocket from 'ws';
import { EventEmitter } from 'events';

interface TendermintEvent {
  query: string;
  data: {
    type: string;
    value: Record<string, unknown>;
  };
  events: Record<string, string[]>;
}

interface BlockHeader {
  height: string;
  chain_id: string;
  time: string;
  proposer_address: string;
  last_commit_hash: string;
  data_hash: string;
  validators_hash: string;
  next_validators_hash: string;
  consensus_hash: string;
  app_hash: string;
  last_results_hash: string;
  evidence_hash: string;
  proposer: string;
}

interface TxEvent {
  height: string;
  txhash: string;
  tx: string;
  result: {
    code: number;
    events: TendermintRawEvent[];
    gas_used: string;
    gas_wanted: string;
  };
}

interface TendermintRawEvent {
  type: string;
  attributes: { key: string; value: string; index: boolean }[];
}

export class TendermintWSClient extends EventEmitter {
  private ws: WebSocket | null = null;
  private idCounter = 0;
  private pendingRequests = new Map<number, {
    resolve: (value: unknown) => void;
    reject: (reason: unknown) => void;
    timer: NodeJS.Timeout;
  }>();
  private reconnectAttempts = 0;
  private maxReconnectAttempts = 10;
  private reconnectDelay = 1000;
  private subscriptions = new Set<string>();
  private url: string;
  private pingInterval: NodeJS.Timeout | null = null;

  constructor(endpoint = 'ws://localhost:26657/websocket') {
    super();
    this.url = endpoint;
  }

  async connect(): Promise<void> {
    return new Promise((resolve, reject) => {
      this.ws = new WebSocket(this.url);

      this.ws.onopen = () => {
        this.reconnectAttempts = 0;
        this.startPing();
        this.emit('connected');
        this.resubscribeAll().then(() => resolve());
      };

      this.ws.onmessage = (event: WebSocket.MessageEvent) => {
        const msg = JSON.parse(event.data.toString());
        this.handleMessage(msg);
      };

      this.ws.onclose = (event: WebSocket.CloseEvent) => {
        this.stopPing();
        this.emit('disconnected', event.code, event.reason);
        this.attemptReconnect();
      };

      this.ws.onerror = (error: WebSocket.ErrorEvent) => {
        this.emit('error', error);
        reject(error);
      };
    });
  }

  private startPing() {
    this.pingInterval = setInterval(() => {
      if (this.ws?.readyState === WebSocket.OPEN) {
        this.ws.ping();
      }
    }, 30000);
  }

  private stopPing() {
    if (this.pingInterval) {
      clearInterval(this.pingInterval);
      this.pingInterval = null;
    }
  }

  private handleMessage(msg: Record<string, unknown>) {
    const { id, result, error, method, params } = msg as any;

    if (id !== undefined && this.pendingRequests.has(id)) {
      const pending = this.pendingRequests.get(id)!;
      clearTimeout(pending.timer);
      this.pendingRequests.delete(id);
      if (error) {
        pending.reject(new Error(JSON.stringify(error)));
      } else {
        pending.resolve(result);
      }
      return;
    }

    if (method === 'Event') {
      const event = params?.data?.value as TendermintEvent;
      const eventType = params?.data?.type;
      this.emit('event', { type: eventType, data: event });
      if (event?.query) {
        this.emit(event.query, event.data);
      }
    }
  }

  private async rpcCall(method: string, params: unknown[] = []): Promise<unknown> {
    return new Promise((resolve, reject) => {
      if (!this.ws || this.ws.readyState !== WebSocket.OPEN) {
        reject(new Error('WebSocket not connected'));
        return;
      }
      const id = ++this.idCounter;
      const request = { jsonrpc: '2.0', id, method, params };
      const timer = setTimeout(() => {
        this.pendingRequests.delete(id);
        reject(new Error(`RPC call ${method} timed out`));
      }, 10000);
      this.pendingRequests.set(id, { resolve, reject, timer });
      this.ws.send(JSON.stringify(request));
    });
  }

  async subscribe(query: string): Promise<void> {
    if (this.subscriptions.has(query)) return;
    await this.rpcCall('subscribe', [query]);
    this.subscriptions.add(query);
    this.emit('subscribed', query);
  }

  async unsubscribe(query: string): Promise<void> {
    await this.rpcCall('unsubscribe', [query]);
    this.subscriptions.delete(query);
    this.emit('unsubscribed', query);
  }

  private async resubscribeAll(): Promise<void> {
    const queries = [...this.subscriptions];
    this.subscriptions.clear();
    for (const query of queries) {
      await this.subscribe(query);
    }
  }

  private async attemptReconnect() {
    if (this.reconnectAttempts >= this.maxReconnectAttempts) {
      this.emit('reconnect_failed');
      return;
    }
    const delay = this.reconnectDelay * Math.pow(2, this.reconnectAttempts);
    this.reconnectAttempts++;
    this.emit('reconnecting', this.reconnectAttempts, delay);
    await new Promise(r => setTimeout(r, delay));
    try {
      await this.connect();
      this.emit('reconnected');
    } catch {
      this.attemptReconnect();
    }
  }

  async disconnect(): Promise<void> {
    this.stopPing();
    this.maxReconnectAttempts = 0;
    for (const query of this.subscriptions) {
      try { await this.unsubscribe(query); } catch { /* ignore */ }
    }
    this.subscriptions.clear();
    this.ws?.close();
    this.ws = null;
  }

  async getBlock(height?: number): Promise<BlockHeader> {
    const params = height ? [`'${height}'`] : [];
    return (await this.rpcCall('block', params)) as BlockHeader;
  }

  async getTx(txhash: string): Promise<TxEvent> {
    return (await this.rpcCall('tx', [txhash])) as TxEvent;
  }
}

async function example() {
  const client = new TendermintWSClient();

  client.on('connected', () => console.log('已连接到 Tendermint WebSocket'));
  client.on('disconnected', (code, reason) => console.log('断开连接:', code, reason));
  client.on('reconnecting', (attempt, delay) => console.log('重连第', attempt, '次'));

  await client.connect();

  // 订阅新块事件
  await client.subscribe("tm.event='NewBlock'");
  client.on("tm.event='NewBlock'", (data) => {
    console.log('新块:', data.value.block.header.height);
  });

  // 订阅交易事件
  await client.subscribe("tm.event='Tx'");
  client.on("tm.event='Tx'", (data) => {
    console.log('新交易:', data.value.TxResult.txhash);
  });

  // 订阅自定义合约事件
  await client.subscribe("tm.event='Tx' AND action='CreateSession'");
  client.on("tm.event='Tx' AND action='CreateSession'", (data) => {
    console.log('新支付会话:', data.value.TxResult.events);
  });

  setTimeout(() => client.disconnect(), 30000);
}

2.2 事件查询语言

Tendermint 事件查询语法基于布尔表达式:

操作符 说明 示例
= 等于 tm.event='Tx'
!= 不等于 tm.event!='Tx'
AND 逻辑与 tm.event='Tx' AND action='CreateSession'
OR 逻辑或 tm.event='NewBlock' OR tm.event='ValidatorSetUpdates'
CONTAINS 包含 sender CONTAINS 'msg1'
EXISTS 存在 action EXISTS

预定义事件类型:

事件类型 触发时机 说明
NewBlock 新区块产生 包含完整区块头
NewBlockHeader 新区块头 仅区块头,不含交易
Tx 交易被执行 包含交易结果和事件
ValidatorSetUpdates 验证人集变更 验证人集更新

自定义合约事件查询:

await client.subscribe("tm.event='Tx' AND action='RegisterAgent'");
await client.subscribe("tm.event='Tx' AND contract_address='msg1contract...'");
await client.subscribe("tm.event='Tx' AND action='CreateSession' AND sender='msg1...'");
await client.subscribe("tm.event='Tx' AND (action='CreateSession' OR action='ReleasePayment')");

2.3 事件格式

NewBlock 事件:

{
  "jsonrpc": "2.0",
  "method": "Event",
  "params": {
    "data": {
      "type": "tendermint/event/NewBlock",
      "value": {
        "block": {
          "header": {
            "chain_id": "msg-chain-1",
            "height": "1000",
            "time": "2026-07-06T10:00:00Z",
            "proposer_address": "msg1proposeraddress..."
          }
        }
      }
    }
  }
}

Tx 事件(合约事件部分):

{
  "jsonrpc": "2.0",
  "method": "Event",
  "params": {
    "data": {
      "type": "tendermint/event/Tx",
      "value": {
        "TxResult": {
          "height": "1000",
          "txhash": "TXHASH1234567890ABCDEF",
          "result": {
            "code": 0,
            "gas_used": "120000",
            "gas_wanted": "200000",
            "events": [
              {
                "type": "wasm",
                "attributes": [
                  { "key": "_contract_address", "value": "msg1contract123...", "index": true },
                  { "key": "action", "value": "CreateSession", "index": true },
                  { "key": "session_id", "value": "session-001", "index": true },
                  { "key": "agent_id", "value": "agent-001", "index": true },
                  { "key": "client", "value": "msg1client...", "index": true },
                  { "key": "amount", "value": "5000000", "index": true }
                ]
              }
            ]
          }
        }
      }
    }
  }
}

2.4 重连处理

TendermintWSClient 内置的完整重连策略:

  1. 指数退避:初始 1s,乘以 2,最大 30s
  2. 最大重试次数:10 次,之后触发 reconnect_failed
  3. 自动恢复订阅:重连后重新订阅所有之前的查询
  4. Ping/Pong:每 30s 发送 ping 维持连接
  5. 断开检测:WebSocket onclose 事件触发重连流程

3. CosmWasm 合约事件

3.1 标准事件格式

CosmWasm 合约使用 deps.events.add_attribute() 或 Response::new().add_attribute() 发射事件:

use cosmwasm_std::{DepsMut, Env, MessageInfo, Response, StdResult};

pub fn execute_create_session(
    deps: DepsMut,
    env: Env,
    info: MessageInfo,
    agent_id: String,
    amount: u128,
) -> StdResult<Response> {
    Ok(Response::new()
        .add_attribute("action", "CreateSession")
        .add_attribute("contract_address", env.contract.address.to_string())
        .add_attribute("sender", info.sender.to_string())
        .add_attribute("agent_id", agent_id)
        .add_attribute("amount", amount.to_string())
        .add_attribute("session_id", generate_session_id()))
}

pub fn emit_raw_event(deps: DepsMut, env: Env) {
    deps.events.add_attribute("action", "do_something");
    deps.events.add_attribute("sender", info.sender.to_string());
}

3.2 合约事件 Schema

3.2.1 Agent Payment v1

interface CreateSessionEvent {
  action: 'CreateSession';
  contract_address: string;
  sender: string;
  session_id: string;
  agent_id: string;
  client: string;
  amount: string;
  deposit: string;
  deadline: string;
  terms_hash: string;
}

interface FundSessionEvent {
  action: 'FundSession';
  contract_address: string;
  sender: string;
  session_id: string;
  amount: string;
  total_balance: string;
}

interface ReleasePaymentEvent {
  action: 'ReleasePayment';
  contract_address: string;
  sender: string;
  session_id: string;
  agent_id: string;
  amount: string;
  release_id: string;
}

interface DisputePaymentEvent {
  action: 'DisputePayment';
  contract_address: string;
  sender: string;
  session_id: string;
  reason: string;
  dispute_id: string;
}

interface CloseSessionEvent {
  action: 'CloseSession';
  contract_address: string;
  sender: string;
  session_id: string;
  reason: string;
  refund_amount: string;
  settled_amount: string;
}

3.2.2 Agent Registry v1

interface RegisterAgentEvent {
  action: 'RegisterAgent';
  contract_address: string;
  sender: string;
  agent_id: string;
  owner: string;
  name: string;
  agent_type: string;
  description: string;
  endpoint: string;
  capabilities: string[];
  fee: string;
}

interface UpdateAgentEvent {
  action: 'UpdateAgent';
  contract_address: string;
  sender: string;
  agent_id: string;
  updated_fields: string[];
}

interface DeregisterAgentEvent {
  action: 'DeregisterAgent';
  contract_address: string;
  sender: string;
  agent_id: string;
  reason: string;
}

interface SetStatusEvent {
  action: 'SetStatus';
  contract_address: string;
  sender: string;
  agent_id: string;
  status: AgentStatus;
}

type AgentStatus = 'ACTIVE' | 'INACTIVE' | 'SUSPENDED' | 'DEREGISTERED';

3.2.3 AID/DID Registry v1

interface CreateDIDEvent {
  action: 'CreateDID';
  contract_address: string;
  sender: string;
  did: string;
  controller: string;
  public_key: string;
  authentication_methods: string[];
  service_endpoints: string[];
}

interface UpdateDIDEvent {
  action: 'UpdateDID';
  contract_address: string;
  sender: string;
  did: string;
  updated_services: string[];
}

interface DeactivateDIDEvent {
  action: 'DeactivateDID';
  contract_address: string;
  sender: string;
  did: string;
}

3.2.4 AI Agent Constitution v1

interface AddRuleEvent {
  action: 'AddRule';
  contract_address: string;
  sender: string;
  rule_id: string;
  rule_type: string;
  description: string;
  parameters: string;
}

interface RemoveRuleEvent {
  action: 'RemoveRule';
  contract_address: string;
  sender: string;
  rule_id: string;
}

interface UpdateRuleEvent {
  action: 'UpdateRule';
  contract_address: string;
  sender: string;
  rule_id: string;
  previous_hash: string;
  new_hash: string;
}

interface SetActiveEvent {
  action: 'SetActive';
  contract_address: string;
  sender: string;
  active: boolean;
  version: string;
}

3.2.5 Micropayment Session v1

interface OpenChannelEvent {
  action: 'OpenChannel';
  contract_address: string;
  sender: string;
  channel_id: string;
  participant_a: string;
  participant_b: string;
  capacity: string;
  expiry: string;
}

interface DepositEvent {
  action: 'Deposit';
  contract_address: string;
  sender: string;
  channel_id: string;
  amount: string;
  total_balance_a: string;
  total_balance_b: string;
}

interface ClaimEvent {
  action: 'Claim';
  contract_address: string;
  sender: string;
  channel_id: string;
  recipient: string;
  amount: string;
  signature: string;
}

interface CloseChannelEvent {
  action: 'CloseChannel';
  contract_address: string;
  sender: string;
  channel_id: string;
  reason: string;
  final_balance_a: string;
  final_balance_b: string;
}

3.3 事件解析与标准化

// src/events/event-normalizer.ts

interface RawAttribute {
  key: string;
  value: string;
  index: boolean;
}

interface RawEvent {
  type: string;
  attributes: RawAttribute[];
}

interface NormalizedEvent {
  id: string;
  txHash: string;
  blockHeight: number;
  eventType: string;
  action: string;
  contractAddress: string;
  attributes: Record<string, string>;
  timestamp: string;
  success: boolean;
}

export class EventNormalizer {
  normalizeTxEvents(txResult: {
    height: string;
    txhash: string;
    result: { events: RawEvent[]; gas_used: string; gas_wanted: string; code: number };
  }): NormalizedEvent[] {
    const events: NormalizedEvent[] = [];
    const blockHeight = parseInt(txResult.height);

    for (const raw of txResult.result.events) {
      if (raw.type !== 'wasm') continue;

      const attrs = this.attributesToRecord(raw.attributes);
      const action = attrs['action'];
      if (!action) continue;

      events.push({
        id: `${txResult.txhash}-${events.length}`,
        txHash: txResult.txhash,
        blockHeight,
        eventType: raw.type,
        action,
        contractAddress: attrs['_contract_address'] || '',
        attributes: attrs,
        timestamp: new Date().toISOString(),
        success: txResult.result.code === 0,
      });
    }
    return events;
  }

  private attributesToRecord(attributes: RawAttribute[]): Record<string, string> {
    const record: Record<string, string> = {};
    for (const attr of attributes) {
      record[attr.key] = attr.value;
    }
    return record;
  }

  parseJsonArray(value: string): string[] {
    try { return JSON.parse(value); }
    catch { return value.split(',').map(s => s.trim()); }
  }
}

export class ContractEventFormatter {
  formatCreateSession(attrs: Record<string, string>): CreateSessionEvent {
    return {
      action: 'CreateSession',
      contract_address: attrs['_contract_address'],
      sender: attrs['sender'],
      session_id: attrs['session_id'],
      agent_id: attrs['agent_id'],
      client: attrs['client'],
      amount: attrs['amount'],
      deposit: attrs['deposit'],
      deadline: attrs['deadline'],
      terms_hash: attrs['terms_hash'],
    };
  }

  formatRegisterAgent(attrs: Record<string, string>): RegisterAgentEvent {
    return {
      action: 'RegisterAgent',
      contract_address: attrs['_contract_address'],
      sender: attrs['sender'],
      agent_id: attrs['agent_id'],
      owner: attrs['owner'] || attrs['sender'],
      name: attrs['name'],
      agent_type: attrs['agent_type'],
      description: attrs['description'],
      endpoint: attrs['endpoint'],
      capabilities: JSON.parse(attrs['capabilities'] || '[]'),
      fee: attrs['fee'],
    };
  }

  formatCreateDID(attrs: Record<string, string>): CreateDIDEvent {
    return {
      action: 'CreateDID',
      contract_address: attrs['_contract_address'],
      sender: attrs['sender'],
      did: attrs['did'],
      controller: attrs['controller'],
      public_key: attrs['public_key'],
      authentication_methods: JSON.parse(attrs['authentication_methods'] || '[]'),
      service_endpoints: JSON.parse(attrs['service_endpoints'] || '[]'),
    };
  }

  formatAddRule(attrs: Record<string, string>): AddRuleEvent {
    return {
      action: 'AddRule',
      contract_address: attrs['_contract_address'],
      sender: attrs['sender'],
      rule_id: attrs['rule_id'],
      rule_type: attrs['rule_type'],
      description: attrs['description'],
      parameters: attrs['parameters'],
    };
  }

  formatOpenChannel(attrs: Record<string, string>): OpenChannelEvent {
    return {
      action: 'OpenChannel',
      contract_address: attrs['_contract_address'],
      sender: attrs['sender'],
      channel_id: attrs['channel_id'],
      participant_a: attrs['participant_a'],
      participant_b: attrs['participant_b'],
      capacity: attrs['capacity'],
      expiry: attrs['expiry'],
    };
  }
}

3.4 事件分类总表

合约 事件名称 关键属性 业务含义
agent_payment_v1 CreateSession session_id, agent_id, client, amount 创建支付会话
agent_payment_v1 FundSession session_id, amount 追加会话资金
agent_payment_v1 ReleasePayment session_id, agent_id, amount 放行支付
agent_payment_v1 DisputePayment session_id, reason 发起争议
agent_payment_v1 CloseSession session_id, reason 关闭会话
agent_registry_v1 RegisterAgent agent_id, owner, name, type 注册 Agent
agent_registry_v1 UpdateAgent agent_id, updated_fields 更新信息
agent_registry_v1 DeregisterAgent agent_id, reason 注销 Agent
agent_registry_v1 SetStatus agent_id, status 设置状态
aidid_did_registry_v1 CreateDID did, controller, public_key 创建 DID
aidid_did_registry_v1 UpdateDID did, updated_services 更新 DID
aidid_did_registry_v1 DeactivateDID did 停用 DID
ai_agent_constitution_v1 AddRule rule_id, rule_type 添加规则
ai_agent_constitution_v1 RemoveRule rule_id 移除规则
ai_agent_constitution_v1 UpdateRule rule_id, new_hash 更新规则
ai_agent_constitution_v1 SetActive active, version 启用/停用
micropayment_session_v1 OpenChannel channel_id, participant_a, participant_b 开通道
micropayment_session_v1 Deposit channel_id, amount 存款
micropayment_session_v1 Claim channel_id, recipient, amount 申领
micropayment_session_v1 CloseChannel channel_id, reason 关通道

4. Indexer 数据管道

4.1 架构设计

Indexer 采用生产者-消费者模式:

Tendermint WS ──→ Block Fetcher ──→ Event Parser ──→ DB Writer
                      │                  │               │
                      ▼                  ▼               ▼
                 Block Queue        Event Queue      PostgreSQL

4.2 Block Fetcher 服务

// src/indexer/block-fetcher.ts

import { TendermintWSClient } from '../websocket/tendermint-ws';
import axios, { AxiosInstance } from 'axios';

interface BlockData {
  height: number;
  hash: string;
  time: string;
  proposer: string;
  txCount: number;
  transactions: TxData[];
}

interface TxData {
  hash: string;
  height: number;
  index: number;
  status: string;
  gasUsed: number;
  gasWanted: number;
  fee: string;
  memo: string;
  events: NormalizedEvent[];
  messages: MessageData[];
}

interface MessageData {
  type: string;
  sender: string;
  contractAddress: string;
  msgType: string;
  data: Record<string, unknown>;
}

export class BlockFetcher {
  private wsClient: TendermintWSClient | null = null;
  private restClient: AxiosInstance;
  private pollingInterval: NodeJS.Timeout | null = null;
  private lastProcessedHeight = 0;
  private onBlock: (block: BlockData) => Promise<void>;

  constructor(
    private rpcEndpoint: string,
    private restEndpoint: string,
    private options: {
      mode: 'websocket' | 'polling';
      pollingIntervalMs?: number;
      startHeight?: number;
    },
  ) {
    this.restClient = axios.create({ baseURL: restEndpoint });
    this.onBlock = async () => {};

    if (options.startHeight) {
      this.lastProcessedHeight = options.startHeight;
    }
  }

  onNewBlock(callback: (block: BlockData) => Promise<void>): void {
    this.onBlock = callback;
  }

  async start(): Promise<void> {
    if (this.options.mode === 'websocket') {
      await this.startWebSocketMode();
    } else {
      this.startPollingMode();
    }
  }

  private async startWebSocketMode(): Promise<void> {
    this.wsClient = new TendermintWSClient(`${this.rpcEndpoint}/websocket`);
    await this.wsClient.connect();

    this.wsClient.subscribe("tm.event='NewBlock'");
    this.wsClient.on("tm.event='NewBlock'", async (data: any) => {
      const block = data.value.block;
      const height = parseInt(block.header.height);
      if (height <= this.lastProcessedHeight) return;

      try {
        const blockData = await this.fetchBlockDetails(height, block);
        await this.onBlock(blockData);
        this.lastProcessedHeight = height;
      } catch (error) {
        console.error('处理区块', height, '失败:', error);
      }
    });
  }

  private startPollingMode(): void {
    const interval = this.options.pollingIntervalMs || 5000;

    this.pollingInterval = setInterval(async () => {
      try {
        const latestHeight = await this.fetchLatestHeight();
        if (latestHeight > this.lastProcessedHeight) {
          for (let h = this.lastProcessedHeight + 1; h <= latestHeight; h++) {
            const blockData = await this.fetchBlockByHeight(h);
            await this.onBlock(blockData);
            this.lastProcessedHeight = h;
          }
        }
      } catch (error) {
        console.error('轮询区块失败:', error);
      }
    }, interval);
  }

  private async fetchBlockByHeight(height: number): Promise<BlockData> {
    const blockResp = await this.restClient.get(
      `/cosmos/base/tendermint/v1beta1/blocks/${height}`
    );
    const block = blockResp.data.block;
    const blockHash = blockResp.data.block_id?.hash || '';
    return this.buildBlockData(block, blockHash);
  }

  private async fetchBlockDetails(height: number, blockData: any): Promise<BlockData> {
    const txsResp = await this.restClient.get(
      `/cosmos/tx/v1beta1/txs?events=tx.height=${height}`
    );
    const txs = txsResp.data.tx_responses || [];

    const transactions = await Promise.all(
      txs.map(async (tx: any, index: number) => this.parseTx(tx, height, index))
    );

    return {
      height,
      hash: blockData?.header?.last_block_id?.hash || '',
      time: blockData.header.time,
      proposer: blockData.header.proposer_address || '',
      txCount: transactions.length,
      transactions,
    };
  }

  private buildBlockData(block: any, blockHash: string): BlockData {
    const header = block.header;
    return {
      height: parseInt(header.height),
      hash: blockHash,
      time: header.time,
      proposer: header.proposer_address || '',
      txCount: block.data?.txs?.length || 0,
      transactions: [],
    };
  }

  private async parseTx(tx: any, height: number, index: number): Promise<TxData> {
    const normalizer = new EventNormalizer();
    const txEvents = normalizer.normalizeTxEvents({
      height: height.toString(),
      txhash: tx.txhash,
      result: {
        events: tx.events || [],
        gas_used: tx.gas_used || '0',
        gas_wanted: tx.gas_wanted || '0',
        code: tx.code || 0,
      },
    });

    const messages = this.parseMessages(tx);

    return {
      hash: tx.txhash,
      height,
      index,
      status: tx.code === 0 ? 'SUCCESS' : 'FAILED',
      gasUsed: parseInt(tx.gas_used || '0'),
      gasWanted: parseInt(tx.gas_wanted || '0'),
      fee: this.extractFee(tx),
      memo: tx.memo || '',
      events: txEvents,
      messages,
    };
  }

  private parseMessages(tx: any): MessageData[] {
    const messages: MessageData[] = [];
    const msgs = tx.tx?.body?.messages || [];

    for (const msg of msgs) {
      if (msg['@type'] === '/cosmwasm.wasm.v1.MsgExecuteContract') {
        messages.push({
          type: '/cosmwasm.wasm.v1.MsgExecuteContract',
          sender: msg.sender,
          contractAddress: msg.contract,
          msgType: Object.keys(msg.msg || {})[0] || '',
          data: msg.msg || {},
        });
      } else if (msg['@type'] === '/cosmwasm.wasm.v1.MsgInstantiateContract') {
        messages.push({
          type: '/cosmwasm.wasm.v1.MsgInstantiateContract',
          sender: msg.sender,
          contractAddress: '',
          msgType: Object.keys(msg.msg || {})[0] || '',
          data: msg.msg || {},
        });
      }
    }
    return messages;
  }

  private extractFee(tx: any): string {
    const fee = tx.tx?.auth_info?.fee;
    if (!fee) return '0umsg';
    return fee.amount?.map((c: any) => `${c.amount}${c.denom}`).join(',') || '0umsg';
  }

  private async fetchLatestHeight(): Promise<number> {
    const resp = await this.restClient.get('/cosmos/base/tendermint/v1beta1/blocks/latest');
    return parseInt(resp.data.block.header.height);
  }

  async stop(): Promise<void> {
    if (this.pollingInterval) {
      clearInterval(this.pollingInterval);
      this.pollingInterval = null;
    }
    if (this.wsClient) {
      await this.wsClient.disconnect();
      this.wsClient = null;
    }
  }
}

4.3 Indexer 主服务

// src/indexer/indexer-service.ts

import { BlockFetcher } from './block-fetcher';
import { DatabaseService } from '../db/database-service';
import { TendermintWSClient } from '../websocket/tendermint-ws';

export class IndexerService {
  private fetcher: BlockFetcher;
  private db: DatabaseService;
  private isRunning = false;

  constructor(private config: {
    rpcEndpoint: string;
    restEndpoint: string;
    databaseUrl: string;
    mode: 'websocket' | 'polling';
    startHeight?: number;
    batchSize?: number;
  }) {
    this.db = new DatabaseService(config.databaseUrl);
    this.fetcher = new BlockFetcher(config.rpcEndpoint, config.restEndpoint, {
      mode: config.mode,
      startHeight: config.startHeight,
    });
  }

  async start(): Promise<void> {
    if (this.isRunning) return;
    this.isRunning = true;
    await this.db.connect();
    await this.verifyChainConnection();
    await this.syncStartHeight();

    this.fetcher.onNewBlock(async (block) => {
      await this.db.insertBlock(block);
    });

    await this.fetcher.start();
    console.log('Indexer 已启动 (模式:', this.config.mode, ')');
  }

  private async verifyChainConnection(): Promise<void> {
    const client = new TendermintWSClient(`${this.config.rpcEndpoint}/websocket`);
    await client.connect();
    const status = await (client as any).rpcCall('status');
    await client.disconnect();
    console.log('链连接成功:', (status as any).node_info.network);
  }

  private async syncStartHeight(): Promise<void> {
    // X-Local-Only=true: 生产环境需要更严谨的重连逻辑
    const lastHeight = await this.db.getLastProcessedHeight();
    if (lastHeight > 0 && !this.config.startHeight) {
      console.log('从高度', lastHeight + 1, '恢复同步');
    }
  }

  async stop(): Promise<void> {
    if (!this.isRunning) return;
    this.isRunning = false;
    await this.fetcher.stop();
    await this.db.disconnect();
    console.log('Indexer 已停止');
  }
}

4.4 事件去重

// src/indexer/deduplicator.ts

class EventDeduplicator {
  private processedEvents = new Set<string>();
  private maxSize: number;

  constructor(maxSize = 100000) {
    this.maxSize = maxSize;
  }

  shouldProcess(eventId: string): boolean {
    if (this.processedEvents.has(eventId)) return false;

    this.processedEvents.add(eventId);

    if (this.processedEvents.size > this.maxSize) {
      const toRemove = Math.floor(this.maxSize * 0.1);
      const iterator = this.processedEvents.values();
      for (let i = 0; i < toRemove; i++) {
        const next = iterator.next();
        if (next.done) break;
        this.processedEvents.delete(next.value);
      }
    }
    return true;
  }

  filterNew(eventIds: string[]): string[] {
    return eventIds.filter(id => this.shouldProcess(id));
  }
}

4.5 区块重组处理

// src/indexer/reorg-handler.ts

class ReorgHandler {
  private reorgSafeDepth = 10;

  async handleReorg(
    height: number,
    incomingHash: string,
    getStoredHash: (height: number) => Promise<string | null>,
    reorgCallback: (fromHeight: number, toHeight: number) => Promise<void>,
  ): Promise<number[]> {
    const storedHash = await getStoredHash(height);
    if (!storedHash || storedHash === incomingHash) return [];

    const reorgedHeights: number[] = [];
    let checkHeight = height;

    while (checkHeight > height - this.reorgSafeDepth) {
      const checkStoredHash = await getStoredHash(checkHeight);
      const checkIncomingHash = await this.getHashAtHeight(checkHeight);

      if (checkStoredHash && checkStoredHash !== checkIncomingHash) {
        reorgedHeights.push(checkHeight);
        checkHeight--;
      } else {
        break;
      }
    }

    if (reorgedHeights.length > 0) {
      const fromHeight = Math.min(...reorgedHeights);
      const toHeight = Math.max(...reorgedHeights);
      await reorgCallback(fromHeight, toHeight);
    }

    return reorgedHeights;
  }

  private async getHashAtHeight(height: number): Promise<string> {
    return '';
  }
}

5. PostgreSQL 数据模型

5.1 完整 Schema

-- migrations/001_schema.sql
-- MSG Chain Indexer PostgreSQL Schema (PG 15+)

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

-- 枚举类型

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 tx_status AS ENUM (
  'SUCCESS', 'FAILED'
);

-- 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,
  app_hash        VARCHAR(64),
  validators_hash VARCHAR(64),
  created_at      TIMESTAMPTZ NOT NULL DEFAULT NOW(),
  CONSTRAINT blocks_hash_check CHECK (hash ~ '^[A-Fa-f0-9]{64}$')
);

COMMENT ON TABLE  blocks     IS '区块头信息';
COMMENT ON COLUMN blocks.height IS '区块高度 (主键)';
COMMENT ON COLUMN blocks.hash   IS '区块哈希 (Hex)';

CREATE INDEX idx_blocks_time ON blocks (time DESC);
CREATE INDEX idx_blocks_proposer ON blocks (proposer);

-- 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,
  gas_used      BIGINT NOT NULL DEFAULT 0,
  gas_wanted    BIGINT NOT NULL DEFAULT 0,
  fee           VARCHAR(128) NOT NULL DEFAULT '0umsg',
  memo          TEXT NOT NULL DEFAULT '',
  code          INTEGER NOT NULL DEFAULT 0,
  raw_log       TEXT,
  sender        VARCHAR(64) NOT NULL DEFAULT '',
  created_at    TIMESTAMPTZ NOT NULL DEFAULT NOW(),
  CONSTRAINT unique_tx_in_block UNIQUE (block_height, index)
);

COMMENT ON TABLE  transactions IS '交易记录 — 类似 Etherscan 的交易列表';

CREATE INDEX idx_transactions_block_height ON transactions (block_height DESC);
CREATE INDEX idx_transactions_status ON transactions (status);
CREATE INDEX idx_transactions_sender ON transactions (sender);
CREATE INDEX idx_transactions_created_at ON transactions (created_at DESC);

-- 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,
  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.attributes IS '事件属性 JSONB — 类似 The Graph 子图的实体属性';

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_event_type ON events (event_type);
CREATE INDEX idx_events_action ON events (action);
CREATE INDEX idx_events_contract_address ON events (contract_address);
CREATE INDEX idx_events_attributes_gin ON events USING GIN (attributes jsonb_path_ops);
CREATE INDEX idx_events_contract_action ON events (contract_address, action);

-- messages 表

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,
  contract_address VARCHAR(64) NOT NULL DEFAULT '',
  msg_type         VARCHAR(128) NOT NULL DEFAULT '',
  data             JSONB NOT NULL DEFAULT '{}',
  created_at       TIMESTAMPTZ NOT NULL DEFAULT NOW()
);

COMMENT ON TABLE messages IS '交易中的 CosmWasm 消息';

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_address 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)
);

COMMENT ON TABLE contracts IS '已实例化的 CosmWasm 合约';

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

-- agents 表

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 注册信息';

CREATE INDEX idx_agents_owner ON agents (owner);
CREATE INDEX idx_agents_status ON agents (status);
CREATE INDEX idx_agents_agent_type ON agents (agent_type);
CREATE INDEX idx_agents_capabilities_gin ON agents USING GIN (capabilities);

-- 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 '支付会话';

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);

-- 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
);

COMMENT ON TABLE did_documents IS '去中心化身份 (DID) 文档';

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

-- constitution_rules 表

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
);

COMMENT ON TABLE constitution_rules IS 'AI Agent 治理规则';

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
);

COMMENT ON TABLE micropayment_channels IS '微支付状态通道';

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);

-- 全文搜索索引

CREATE INDEX idx_blocks_hash_trgm ON blocks USING GIN (hash gin_trgm_ops);
CREATE INDEX idx_transactions_hash_trgm ON transactions USING GIN (hash gin_trgm_ops);
CREATE INDEX idx_contracts_address_trgm ON contracts USING GIN (address gin_trgm_ops);
CREATE INDEX idx_agents_id_trgm ON agents USING GIN (agent_id gin_trgm_ops);
CREATE INDEX idx_agents_name_trgm ON agents USING GIN (name gin_trgm_ops);

ALTER TABLE transactions ADD COLUMN search_vector tsvector
  GENERATED ALWAYS AS (to_tsvector('simple', coalesce(memo, ''))) STORED;
CREATE INDEX idx_transactions_search ON transactions USING GIN (search_vector);

-- 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',
  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 'Indexer 同步状态跟踪';

INSERT INTO indexer_state (last_processed_height) VALUES (0);

5.2 种子数据

-- migrations/002_seed.sql

INSERT INTO contracts (address, code_id, creator, admin, label, contract_type, created_at) VALUES
  ('msg1agentpayment...', 1, 'msg1creator...', 'msg1admin...', 'Agent Payment v1', 'agent_payment_v1', NOW()),
  ('msg1agentregistry...', 2, 'msg1creator...', 'msg1admin...', 'Agent Registry v1', 'agent_registry_v1', NOW()),
  ('msg1didregistry...', 3, 'msg1creator...', 'msg1admin...', 'AID/DID Registry v1', 'aidid_did_registry_v1', NOW()),
  ('msg1constitution...', 4, 'msg1creator...', 'msg1admin...', 'AI Constitution v1', 'ai_agent_constitution_v1', NOW()),
  ('msg1micropayment...', 5, 'msg1creator...', 'msg1admin...', 'Micropayment Session v1', 'micropayment_session_v1', NOW());

INSERT INTO agents (agent_id, owner, name, status, agent_type, description, endpoint, capabilities, fee, created_height, created_at, updated_at) VALUES
  ('agent-alpha', 'msg1owner1...', 'Alpha Assistant', 'ACTIVE', 'ai',
   '通用 AI 助手',
    'https://agent-alpha.msgchain.org/api',
   '["text-generation", "data-analysis", "code-review"]',
   '1000000', 100, NOW(), NOW()),
  ('agent-beta', 'msg1owner2...', 'Beta Trader', 'ACTIVE', 'ai',
   '去中心化交易策略 Agent',
    'https://agent-beta.msgchain.org/api',
   '["trading", "market-analysis", "portfolio-management"]',
   '2000000', 200, NOW(), NOW());

INSERT INTO sessions (session_id, agent_id, client, status, amount, deposit, balance, deadline, created_height, created_at, updated_at) VALUES
  ('session-001', 'agent-alpha', 'msg1client1...', 'RELEASED', 5000000, 5000000, 5000000,
   NOW() + INTERVAL '7 days', 150, NOW(), NOW()),
  ('session-002', 'agent-beta', 'msg1client2...', 'OPEN', 10000000, 3000000, 3000000,
   NOW() + INTERVAL '14 days', 200, NOW(), NOW());

INSERT INTO did_documents (did, controller, public_key, authentication_methods, service_endpoints, active, created_height, created_at, updated_at) VALUES
  ('did:msg:msg1owner1...', 'msg1owner1...',
   '{"type":"Ed25519VerificationKey2020","publicKeyMultibase":"z6Mk..."}',
   '[{"id":"#key-1","type":"Ed25519VerificationKey2020","controller":"msg1owner1..."}]',
   '[{"id":"#agent-alpha","type":"AgentEndpoint","serviceEndpoint":"https://agent-alpha.msgchain.org/api"}]',
   true, 100, NOW(), NOW());

5.3 额外性能索引

-- migrations/003_indexes.sql

CREATE INDEX CONCURRENTLY IF NOT EXISTS idx_events_recent
  ON events (block_height DESC, id DESC);

CREATE INDEX CONCURRENTLY IF NOT EXISTS idx_transactions_recent
  ON transactions (block_height DESC, index ASC);

CREATE INDEX CONCURRENTLY IF NOT EXISTS idx_events_action_count
  ON events (action, block_height DESC);

CREATE INDEX CONCURRENTLY IF NOT EXISTS idx_sessions_agent_stats
  ON sessions (agent_id, status, amount);

CREATE INDEX CONCURRENTLY IF NOT EXISTS idx_sessions_active
  ON sessions (agent_id, created_at DESC)
  WHERE status IN ('OPEN', 'FUNDED');

CREATE INDEX CONCURRENTLY IF NOT EXISTS idx_agents_active
  ON agents (agent_id, agent_type)
  WHERE status = 'ACTIVE';

CREATE INDEX CONCURRENTLY IF NOT EXISTS idx_blocks_time_series
  ON blocks (time DESC) INCLUDE (height, hash, tx_count);

5.4 数据保留策略

-- migrations/004_retention.sql

CREATE OR REPLACE FUNCTION apply_retention_policy(
  p_retention_days INTEGER DEFAULT 365
) RETURNS INTEGER AS $$
DECLARE
  v_cutoff_date TIMESTAMPTZ;
  v_deleted_count INTEGER := 0;
BEGIN
  v_cutoff_date := NOW() - (p_retention_days || ' days')::INTERVAL;
  DELETE FROM events WHERE created_at < v_cutoff_date;
  GET DIAGNOSTICS v_deleted_count = ROW_COUNT;
  DELETE FROM transactions WHERE created_at < v_cutoff_date;
  RETURN v_deleted_count;
END;
$$ LANGUAGE plpgsql;

5.5 数据库服务

// src/db/database-service.ts

import { Pool, PoolClient } from 'pg';

export class DatabaseService {
  private pool: Pool;

  constructor(databaseUrl: string) {
    this.pool = new Pool({
      connectionString: databaseUrl,
      max: 20,
      idleTimeoutMillis: 30000,
      connectionTimeoutMillis: 5000,
    });
    this.pool.on('error', (err) => console.error('数据库连接池错误:', err));
  }

  async connect(): Promise<void> {
    await this.pool.connect();
    console.log('数据库已连接');
  }

  async disconnect(): Promise<void> {
    await this.pool.end();
  }

  async insertBlock(block: BlockData): Promise<void> {
    const client = await this.pool.connect();
    try {
      await client.query('BEGIN');

      await client.query(
        `INSERT INTO blocks (height, hash, time, proposer, tx_count)
         VALUES ($1, $2, $3, $4, $5)
         ON CONFLICT (height) DO UPDATE SET hash = EXCLUDED.hash, time = EXCLUDED.time`,
        [block.height, block.hash, block.time, block.proposer, block.txCount]
      );

      for (const tx of block.transactions) {
        await this.insertTx(client, tx);
        for (const event of tx.events) {
          await this.insertEvent(client, event);
        }
        for (const msg of tx.messages) {
          await this.insertMessage(client, msg);
        }
      }

      await client.query(
        `UPDATE indexer_state SET last_processed_height = $1, last_processed_hash = $2, updated_at = NOW() WHERE id = 1`,
        [block.height, block.hash]
      );

      await client.query('COMMIT');
    } catch (error) {
      await client.query('ROLLBACK');
      throw error;
    } finally {
      client.release();
    }
  }

  private async insertTx(client: PoolClient, tx: TxData): Promise<void> {
    await client.query(
      `INSERT INTO transactions (hash, block_height, index, status, gas_used, gas_wanted, fee, memo, code)
       VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9)
       ON CONFLICT (hash) DO NOTHING`,
      [tx.hash, tx.height, tx.index, tx.status, tx.gasUsed, tx.gasWanted, tx.fee, tx.memo,
       tx.status === 'SUCCESS' ? 0 : 1]
    );
  }

  private async insertEvent(client: PoolClient, event: NormalizedEvent): Promise<void> {
    await client.query(
      `INSERT INTO events (tx_hash, block_height, event_type, action, contract_address, attributes, success)
       VALUES ($1, $2, $3, $4, $5, $6, $7)`,
      [event.txHash, event.blockHeight, event.eventType, event.action,
       event.contractAddress, JSON.stringify(event.attributes), event.success]
    );
  }

  private async insertMessage(client: PoolClient, msg: MessageData): Promise<void> {
    await client.query(
      `INSERT INTO messages (tx_hash, block_height, index, message_type, sender, contract_address, msg_type, data)
       VALUES ($1, $2, $3, $4, $5, $6, $7, $8)`,
      [msg.txHash, msg.blockHeight, 0, msg.type, msg.sender, msg.contractAddress, msg.msgType, JSON.stringify(msg.data)]
    );
  }

  async getLastProcessedHeight(): Promise<number> {
    const result = await this.pool.query(
      'SELECT last_processed_height FROM indexer_state WHERE id = 1'
    );
    return result.rows[0]?.last_processed_height || 0;
  }

  async getBlock(height: number): Promise<BlockData | null> {
    const result = await this.pool.query('SELECT * FROM blocks WHERE height = $1', [height]);
    if (result.rows.length === 0) return null;
    return this.rowToBlock(result.rows[0]);
  }

  async getBlocks(limit = 20, offset = 0): Promise<BlockData[]> {
    const result = await this.pool.query(
      'SELECT * FROM blocks ORDER BY height DESC LIMIT $1 OFFSET $2', [limit, offset]
    );
    return result.rows.map(r => this.rowToBlock(r));
  }

  async getTransaction(hash: string): Promise<TxData | null> {
    const result = await this.pool.query('SELECT * FROM transactions WHERE hash = $1', [hash]);
    if (result.rows.length === 0) return null;
    return this.rowToTx(result.rows[0]);
  }

  async getTransactionsByBlock(height: number): Promise<TxData[]> {
    const result = await this.pool.query(
      'SELECT * FROM transactions WHERE block_height = $1 ORDER BY index ASC', [height]
    );
    return result.rows.map(r => this.rowToTx(r));
  }

  async getEventsByType(action: string, limit = 50, offset = 0): Promise<NormalizedEvent[]> {
    const result = await this.pool.query(
      `SELECT e.*, t.hash as tx_hash FROM events e
       JOIN transactions t ON t.hash = e.tx_hash
       WHERE e.action = $1
       ORDER BY e.block_height DESC, e.id DESC LIMIT $2 OFFSET $3`,
      [action, limit, offset]
    );
    return result.rows.map(r => this.rowToEvent(r));
  }

  async getContract(address: string): Promise<any | null> {
    const result = await this.pool.query('SELECT * FROM contracts WHERE address = $1', [address]);
    return result.rows[0] || null;
  }

  async getContracts(limit = 20, offset = 0): Promise<any[]> {
    const result = await this.pool.query(
      'SELECT * FROM contracts ORDER BY created_at DESC LIMIT $1 OFFSET $2', [limit, offset]
    );
    return result.rows;
  }

  async getAgent(agentId: string): Promise<any | null> {
    const result = await this.pool.query('SELECT * FROM agents WHERE agent_id = $1', [agentId]);
    return result.rows[0] || null;
  }

  async getAgentsByOwner(owner: string): Promise<any[]> {
    const result = await this.pool.query(
      'SELECT * FROM agents WHERE owner = $1 ORDER BY created_at DESC', [owner]
    );
    return result.rows;
  }

  async search(query: string): Promise<SearchResult> {
    const pattern = `%${query}%`;

    const [blocks, txs, contracts, agents, sessions] = await Promise.all([
      this.pool.query('SELECT height, hash, time FROM blocks WHERE hash ILIKE $1 LIMIT 10', [pattern]),
      this.pool.query('SELECT hash, block_height, status FROM transactions WHERE hash ILIKE $1 LIMIT 10', [pattern]),
      this.pool.query('SELECT address, label, contract_type FROM contracts WHERE address ILIKE $1 OR label ILIKE $2 LIMIT 10', [pattern, `%${query}%`]),
      this.pool.query('SELECT agent_id, name, owner, status FROM agents WHERE agent_id ILIKE $1 OR name ILIKE $2 LIMIT 10', [pattern, `%${query}%`]),
      this.pool.query('SELECT session_id, agent_id, client, status FROM sessions WHERE session_id ILIKE $1 LIMIT 10', [pattern]),
    ]);

    return {
      blocks: blocks.rows,
      transactions: txs.rows,
      contracts: contracts.rows,
      agents: agents.rows,
      sessions: sessions.rows,
    };
  }

  private rowToBlock(row: any): BlockData {
    return { height: row.height, hash: row.hash, time: row.time, proposer: row.proposer, txCount: row.tx_count, transactions: [] };
  }

  private rowToTx(row: any): TxData {
    return { hash: row.hash, height: row.block_height, index: row.index, status: row.status, gasUsed: row.gas_used, gasWanted: row.gas_wanted, fee: row.fee, memo: row.memo || '', events: [], messages: [] };
  }

  private rowToEvent(row: any): NormalizedEvent {
    return { id: `${row.tx_hash}-${row.id}`, txHash: row.tx_hash, blockHeight: row.block_height, eventType: row.event_type, action: row.action, contractAddress: row.contract_address, attributes: row.attributes || {}, timestamp: row.created_at, success: row.success };
  }
}

interface SearchResult {
  blocks: any[]; transactions: any[]; contracts: any[]; agents: any[]; sessions: any[];
}

6. GraphQL API 服务

6.1 服务器设置

// src/graphql/server.ts

import { createYoga, createSchema } from 'graphql-yoga';
import { createServer } from 'node:http';
import { resolvers } from './resolvers';
import { typeDefs } from './schema';

export async function startGraphQLServer(options: { port: number; databaseUrl: string }) {
  const schema = createSchema({ typeDefs, resolvers });

  const yoga = createYoga({
    schema,
    context: async () => {
      const db = new DatabaseService(options.databaseUrl);
      await db.connect();
      return { db };
    },
    subscriptions: { path: '/graphql' },
    cors: { origin: '*', credentials: true },
    maskedErrors: false,
    logging: true,
  });

  const httpServer = createServer(yoga);
  httpServer.listen(options.port, () => {
    console.log('GraphQL API 运行在 http://localhost:' + options.port + '/graphql');
  });
  return httpServer;
}

6.2 完整 GraphQL Schema

# schema.graphql

scalar DateTime
scalar JSON
scalar BigInt

enum AgentStatus { ACTIVE INACTIVE SUSPENDED DEREGISTERED }
enum SessionStatus { OPEN FUNDED RELEASED DISPUTED CLOSED EXPIRED }
enum ChannelStatus { OPEN SETTLING CLOSED }
enum TxStatus { SUCCESS FAILED }
enum EventSortField { BLOCK_HEIGHT CREATED_AT ACTION }
enum SortOrder { ASC DESC }

type Block {
  height: Int!
  hash: String!
  time: DateTime!
  proposer: String!
  txCount: Int!
  transactions(limit: Int = 20, offset: Int = 0): [Transaction!]!
  eventStats: EventStats!
}

type EventStats {
  totalEvents: Int!
  eventTypeCounts: [EventTypeCount!]!
}

type EventTypeCount {
  action: String!
  count: Int!
}

type Transaction {
  hash: String!
  blockHeight: Int!
  index: Int!
  status: TxStatus!
  gasUsed: Int!
  gasWanted: Int!
  fee: String!
  memo: String!
  block: Block!
  events: [Event!]!
  messages: [Message!]!
}

type Event {
  id: ID!
  txHash: String!
  blockHeight: Int!
  eventType: String!
  action: String!
  contractAddress: String!
  attributes: JSON!
  success: Boolean!
  createdAt: DateTime!
  transaction: Transaction!
}

type Message {
  id: ID!
  txHash: String!
  blockHeight: Int!
  index: Int!
  messageType: String!
  sender: String!
  contractAddress: String!
  msgType: String!
  data: JSON!
  transaction: Transaction!
}

type Contract {
  address: String!
  codeId: Int!
  creator: String!
  admin: String
  label: String!
  contractType: String
  createdAt: DateTime!
  events(limit: Int = 20, offset: Int = 0): [Event!]!
}

type Agent {
  agentId: String!
  owner: String!
  name: String!
  status: AgentStatus!
  agentType: String!
  description: String!
  endpoint: String!
  capabilities: [String!]!
  fee: String!
  contractAddress: String
  createdHeight: Int!
  createdAt: DateTime!
  updatedAt: DateTime!
  sessions(limit: Int = 20, offset: Int = 0): [Session!]!
}

type Session {
  sessionId: String!
  agentId: String!
  agent: Agent!
  client: String!
  status: SessionStatus!
  amount: BigInt!
  deposit: BigInt!
  balance: BigInt!
  deadline: DateTime
  termsHash: String
  contractAddress: String
  createdHeight: Int!
  createdAt: DateTime!
  closedReason: String
  closedAt: DateTime
}

type DIDDocument {
  did: String!
  controller: String!
  publicKey: JSON!
  authenticationMethods: JSON!
  serviceEndpoints: JSON!
  active: Boolean!
  contractAddress: String
  createdHeight: Int!
  createdAt: DateTime!
  deactivatedAt: DateTime
}

type ConstitutionRule {
  ruleId: String!
  ruleType: String!
  description: String!
  parameters: JSON!
  active: Boolean!
  version: String!
  contractAddress: String
  createdHeight: Int!
  createdAt: DateTime!
  updatedAt: DateTime!
}

type MicropaymentChannel {
  channelId: String!
  participantA: String!
  participantB: String!
  capacity: BigInt!
  balanceA: BigInt!
  balanceB: BigInt!
  status: ChannelStatus!
  expiry: DateTime
  contractAddress: String
  createdHeight: Int!
  createdAt: DateTime!
  updatedAt: DateTime!
}

union SearchResultItem = Block | Transaction | Contract | Agent | Session

type SearchStats {
  blocks: Int!
  transactions: Int!
  contracts: Int!
  agents: Int!
  sessions: Int!
}

type SearchResults {
  items: [SearchResultItem!]!
  stats: SearchStats!
  truncated: Boolean!
}

type Query {
  block(height: Int!): Block
  blocks(limit: Int = 20, offset: Int = 0): [Block!]!
  transaction(hash: String!): Transaction
  transactionsByBlock(height: Int!): [Transaction!]!
  eventsByAction(action: String!, limit: Int = 50, offset: Int = 0): [Event!]!
  eventsByContract(contractAddress: String!, limit: Int = 50, offset: Int = 0): [Event!]!
  events(
    action: String, contractAddress: String, txHash: String, blockHeight: Int,
    fromBlock: Int, toBlock: Int, sortField: EventSortField = BLOCK_HEIGHT,
    sortOrder: SortOrder = DESC, limit: Int = 50, offset: Int = 0
  ): [Event!]!
  contract(address: String!): Contract
  contracts(limit: Int = 20, offset: Int = 0): [Contract!]!
  agent(agentId: String!): Agent
  agentsByOwner(owner: String!): [Agent!]!
  agents(status: AgentStatus, agentType: String, limit: Int = 20, offset: Int = 0): [Agent!]!
  session(sessionId: String!): Session
  sessionsByAgent(agentId: String!, limit: Int = 20, offset: Int = 0): [Session!]!
  sessionsByClient(client: String!, limit: Int = 20, offset: Int = 0): [Session!]!
  did(did: String!): DIDDocument
  didsByController(controller: String!): [DIDDocument!]!
  rule(ruleId: String!): ConstitutionRule
  activeRules(version: String): [ConstitutionRule!]!
  channel(channelId: String!): MicropaymentChannel
  channelsByParticipant(participant: String!): [MicropaymentChannel!]!
  search(q: String!): SearchResults!
  latestHeight: Int!
  chainSummary: ChainSummary!
}

type ChainSummary {
  latestBlockHeight: Int!
  totalBlocks: Int!
  totalTransactions: Int!
  totalEvents: Int!
  totalContracts: Int!
  totalAgents: Int!
  totalSessions: Int!
  averageBlockTime: Float!
  activeAgents: Int!
}

type Mutation {
  submitMemo(input: MemoInput!): MemoResponse!
}

input MemoInput {
  content: String!
  type: String!
  referenceId: String
}

type MemoResponse {
  success: Boolean!
  memoId: String
  error: String
}

type Subscription {
  newBlock: Block!
  newTransaction: Transaction!
  newEvent(eventType: String): Event!
  agentStatusChanged: AgentStatusChange!
  newSession: Session!
  didChanged: DIDChange!
}

type AgentStatusChange {
  agentId: String!
  previousStatus: AgentStatus!
  newStatus: AgentStatus!
  timestamp: DateTime!
  txHash: String!
}

type DIDChange {
  did: String!
  changeType: String!
  timestamp: DateTime!
  txHash: String!
}

6.3 Query 示例

# 最新区块
query LatestBlocks {
  blocks(limit: 5) { height hash time txCount }
}

# 区块详情
query BlockDetail($height: Int!) {
  block(height: $height) {
    hash time proposer txCount
    transactions { hash status gasUsed fee }
    eventStats { totalEvents eventTypeCounts { action count } }
  }
}

# 交易详情
query TxDetail($hash: String!) {
  transaction(hash: $hash) {
    blockHeight status gasUsed fee memo
    block { height time }
    events { action contractAddress attributes }
    messages { messageType sender contractAddress data }
  }
}

# Agent + 会话
query AgentWithSessions($agentId: String!) {
  agent(agentId: $agentId) {
    name owner status agentType capabilities
    sessions { sessionId client status amount createdAt }
  }
}

# 链摘要
query Summary { chainSummary { latestBlockHeight totalBlocks totalTransactions totalEvents totalAgents averageBlockTime activeAgents } }

# 搜索
query Search($q: String!) {
  search(q: $q) {
    stats { blocks transactions agents }
    items {
      ... on Block { height hash time }
      ... on Transaction { hash blockHeight status }
      ... on Agent { agentId name owner }
      ... on Contract { address label contractType }
    }
  }
}

# 实时订阅
subscription OnNewBlock { newBlock { height hash time txCount } }
subscription OnNewTransaction { newTransaction { hash blockHeight status } }
subscription OnCreateSession { newEvent(eventType: "CreateSession") { txHash blockHeight attributes } }
subscription OnAgentStatusChanged { agentStatusChanged { agentId previousStatus newStatus txHash } }

6.4 Subscription 集成

// src/graphql/subscription-integration.ts

import { PubSub } from 'graphql-yoga';
import { TendermintWSClient } from '../websocket/tendermint-ws';
import { EventNormalizer } from '../events/event-normalizer';

export class SubscriptionBridge {
  private wsClient: TendermintWSClient;
  private pubsub: PubSub;
  private normalizer: EventNormalizer;

  constructor(rpcEndpoint: string, pubsub: PubSub) {
    this.wsClient = new TendermintWSClient(rpcEndpoint + '/websocket');
    this.pubsub = pubsub;
    this.normalizer = new EventNormalizer();
  }

  async start(): Promise<void> {
    await this.wsClient.connect();

    await this.wsClient.subscribe("tm.event='NewBlock'");
    this.wsClient.on("tm.event='NewBlock'", async (data: any) => {
      const block = data.value.block;
      await this.pubsub.publish('NEW_BLOCK', {
        height: parseInt(block.header.height),
        hash: block.header.last_block_id?.hash || '',
        time: block.header.time,
        proposer: block.header.proposer_address || '',
        txCount: block.data?.txs?.length || 0,
      });
    });

    await this.wsClient.subscribe("tm.event='Tx'");
    this.wsClient.on("tm.event='Tx'", async (data: any) => {
      const txResult = data.value.TxResult;
      await this.pubsub.publish('NEW_TRANSACTION', {
        hash: txResult.txhash,
        blockHeight: parseInt(txResult.height),
        status: txResult.result.code === 0 ? 'SUCCESS' : 'FAILED',
        gasUsed: parseInt(txResult.result.gas_used || '0'),
        gasWanted: parseInt(txResult.result.gas_wanted || '0'),
      });

      const events = this.normalizer.normalizeTxEvents({
        height: txResult.height,
        txhash: txResult.txhash,
        result: txResult.result,
      });

      for (const event of events) {
        await this.pubsub.publish('NEW_EVENT', event);
        if (event.action) {
          await this.pubsub.publish('EVENT:' + event.action, event);
        }

        switch (event.action) {
          case 'RegisterAgent':
          case 'SetStatus':
          case 'DeregisterAgent':
            await this.pubsub.publish('AGENT_STATUS_CHANGED', {
              agentId: event.attributes['agent_id'],
              previousStatus: '',
              newStatus: event.attributes['status'] || 'ACTIVE',
              timestamp: new Date().toISOString(),
              txHash: event.txHash,
            });
            break;
          case 'CreateSession':
            await this.pubsub.publish('NEW_SESSION', {
              sessionId: event.attributes['session_id'],
              agentId: event.attributes['agent_id'],
              client: event.attributes['client'],
              status: 'OPEN',
              amount: event.attributes['amount'],
              createdAt: new Date().toISOString(),
            });
            break;
          case 'CreateDID':
          case 'UpdateDID':
          case 'DeactivateDID':
            await this.pubsub.publish('DID_CHANGED', {
              did: event.attributes['did'],
              changeType: event.action.replace('DID', '').toUpperCase(),
              timestamp: new Date().toISOString(),
              txHash: event.txHash,
            });
            break;
        }
      }
    });
  }

  async stop(): Promise<void> {
    await this.wsClient.disconnect();
  }
}

7. Indexer 事件索引

7.1 事件处理管道

原始 Tendermint 事件
        │
        ▼
┌───────────────┐
│ 事件分类器     │
└───────────────┘
        │
        ├──→ Agent Payment 事件 → sessions 表
        ├──→ Agent Registry 事件 → agents 表
        ├──→ DID Registry 事件  → did_documents 表
        ├──→ Constitution 事件  → constitution_rules 表
        └──→ Micropayment 事件 → micropayment_channels 表

7.2 分阶段索引管道

// src/indexer/event-pipeline.ts

import { DatabaseService } from '../db/database-service';
import { EventNormalizer, ContractEventFormatter } from '../events/event-normalizer';

export class EventPipeline {
  private normalizer: EventNormalizer;
  private formatter: ContractEventFormatter;
  private db: DatabaseService;

  constructor(db: DatabaseService) {
    this.normalizer = new EventNormalizer();
    this.formatter = new ContractEventFormatter();
    this.db = db;
  }

  async processTransaction(ctx: {
    blockHeight: number;
    txHash: string;
    rawEvents: any[];
  }): Promise<void> {
    const normalized = this.normalizer.normalizeTxEvents({
      height: ctx.blockHeight.toString(),
      txhash: ctx.txHash,
      result: { events: ctx.rawEvents, gas_used: '0', gas_wanted: '0', code: 0 },
    });

    for (const event of normalized) {
      const attrs = event.attributes;

      switch (event.action) {
        case 'RegisterAgent':
          await this.db['pool'].query(
            `INSERT INTO agents (agent_id, owner, name, status, agent_type, description, endpoint, capabilities, fee, contract_address, created_height, created_at, updated_at)
             VALUES ($1, $2, $3, 'ACTIVE', $4, $5, $6, $7, $8, $9, $10, NOW(), NOW())
             ON CONFLICT (agent_id) DO UPDATE SET name = EXCLUDED.name, updated_at = NOW()`,
            [attrs['agent_id'], attrs['owner'] || attrs['sender'], attrs['name'],
             attrs['agent_type'], attrs['description'], attrs['endpoint'],
             attrs['capabilities'] || '[]', attrs['fee'], attrs['_contract_address'], ctx.blockHeight]
          );
          break;

        case 'SetStatus':
          await this.db['pool'].query(
            `UPDATE agents SET status = $1::agent_status, updated_at = NOW() WHERE agent_id = $2`,
            [attrs['status'], attrs['agent_id']]
          );
          break;

        case 'DeregisterAgent':
          await this.db['pool'].query(
            `UPDATE agents SET status = 'DEREGISTERED', updated_at = NOW() WHERE agent_id = $1`,
            [attrs['agent_id']]
          );
          break;

        case 'CreateSession':
          await this.db['pool'].query(
            `INSERT INTO sessions (session_id, agent_id, client, status, amount, deposit, balance, deadline, terms_hash, contract_address, created_height, created_at, updated_at)
             VALUES ($1, $2, $3, 'OPEN', $4, $5, $5, to_timestamp($6::double precision), $7, $8, $9, NOW(), NOW())
             ON CONFLICT (session_id) DO NOTHING`,
            [attrs['session_id'], attrs['agent_id'], attrs['client'], attrs['amount'],
             attrs['deposit'], attrs['deadline'], attrs['terms_hash'],
             attrs['_contract_address'], ctx.blockHeight]
          );
          break;

        case 'FundSession':
          await this.db['pool'].query(
            `UPDATE sessions SET status = 'FUNDED', balance = balance + $1, updated_at = NOW() WHERE session_id = $2`,
            [attrs['amount'], attrs['session_id']]
          );
          break;

        case 'ReleasePayment':
          await this.db['pool'].query(
            `UPDATE sessions SET status = 'RELEASED', balance = balance - $1, updated_at = NOW() WHERE session_id = $2`,
            [attrs['amount'], attrs['session_id']]
          );
          break;

        case 'DisputePayment':
          await this.db['pool'].query(
            `UPDATE sessions SET status = 'DISPUTED', updated_at = NOW() WHERE session_id = $1`,
            [attrs['session_id']]
          );
          break;

        case 'CloseSession':
          await this.db['pool'].query(
            `UPDATE sessions SET status = 'CLOSED', closed_reason = $2, closed_at = NOW(), updated_at = NOW() WHERE session_id = $1`,
            [attrs['session_id'], attrs['reason']]
          );
          break;

        case 'CreateDID':
          await this.db['pool'].query(
            `INSERT INTO did_documents (did, controller, public_key, authentication_methods, service_endpoints, active, contract_address, created_height, created_at, updated_at)
             VALUES ($1, $2, $3, $4, $5, true, $6, $7, NOW(), NOW())
             ON CONFLICT (did) DO UPDATE SET controller = EXCLUDED.controller, updated_at = NOW()`,
            [attrs['did'], attrs['controller'], attrs['public_key'],
             JSON.stringify(JSON.parse(attrs['authentication_methods'] || '[]')),
             JSON.stringify(JSON.parse(attrs['service_endpoints'] || '[]')),
             attrs['_contract_address'], ctx.blockHeight]
          );
          break;

        case 'DeactivateDID':
          await this.db['pool'].query(
            `UPDATE did_documents SET active = false, deactivated_at = NOW(), updated_at = NOW() WHERE did = $1`,
            [attrs['did']]
          );
          break;

        case 'AddRule':
          await this.db['pool'].query(
            `INSERT INTO constitution_rules (rule_id, rule_type, description, parameters, active, version, contract_address, created_height, created_at, updated_at)
             VALUES ($1, $2, $3, $4, true, '1.0', $5, $6, NOW(), NOW())
             ON CONFLICT (rule_id) DO NOTHING`,
            [attrs['rule_id'], attrs['rule_type'], attrs['description'], attrs['parameters'],
             attrs['_contract_address'], ctx.blockHeight]
          );
          break;

        case 'RemoveRule':
          await this.db['pool'].query(
            `UPDATE constitution_rules SET active = false, removed_at = NOW(), updated_at = NOW() WHERE rule_id = $1`,
            [attrs['rule_id']]
          );
          break;

        case 'OpenChannel':
          await this.db['pool'].query(
            `INSERT INTO micropayment_channels (channel_id, participant_a, participant_b, capacity, balance_a, balance_b, status, expiry, contract_address, created_height, created_at, updated_at)
             VALUES ($1, $2, $3, $4, 0, 0, 'OPEN', to_timestamp($5::double precision), $6, $7, NOW(), NOW())
             ON CONFLICT (channel_id) DO NOTHING`,
            [attrs['channel_id'], attrs['participant_a'], attrs['participant_b'],
             attrs['capacity'], attrs['expiry'], attrs['_contract_address'], ctx.blockHeight]
          );
          break;

        case 'CloseChannel':
          await this.db['pool'].query(
            `UPDATE micropayment_channels SET status = 'CLOSED', closed_at = NOW(), updated_at = NOW() WHERE channel_id = $1`,
            [attrs['channel_id']]
          );
          break;
      }
    }
  }
}

7.3 批量写入优化

// src/indexer/batch-writer.ts

interface BatchBuffer {
  events: NormalizedEvent[];
  messages: MessageData[];
}

export class BatchWriter {
  private buffer: BatchBuffer = { events: [], messages: [] };
  private batchSize = 500;
  private flushInterval = 1000;
  private timer: NodeJS.Timeout | null = null;
  private db: DatabaseService;

  constructor(db: DatabaseService) {
    this.db = db;
    this.startAutoFlush();
  }

  add(data: Partial<BatchBuffer>): void {
    if (data.events) this.buffer.events.push(...data.events);
    if (data.messages) this.buffer.messages.push(...data.messages);
    if (this.buffer.events.length >= this.batchSize) this.flush();
  }

  private startAutoFlush(): void {
    this.timer = setInterval(() => this.flush(), this.flushInterval);
  }

  async flush(): Promise<void> {
    if (this.buffer.events.length === 0) return;

    const batch = this.buffer;
    this.buffer = { events: [], messages: [] };

    const client = await this.db['pool'].connect();
    try {
      await client.query('BEGIN');

      for (const event of batch.events) {
        await client.query(
          `INSERT INTO events (tx_hash, block_height, event_type, action, contract_address, attributes, success)
           VALUES ($1, $2, $3, $4, $5, $6, $7) ON CONFLICT DO NOTHING`,
          [event.txHash, event.blockHeight, event.eventType, event.action,
           event.contractAddress, JSON.stringify(event.attributes), event.success]
        );
      }

      await client.query('COMMIT');
    } catch (error) {
      await client.query('ROLLBACK');
      this.buffer.events.unshift(...batch.events);
      this.buffer.messages.unshift(...batch.messages);
      console.error('批量写入失败, 已回滚:', error);
    } finally {
      client.release();
    }
  }

  stop(): void {
    if (this.timer) {
      clearInterval(this.timer);
      this.timer = null;
    }
    this.flush();
  }
}

8. Agent API 事件网关

8.1 WebSocket 桥接

Agent API 提供 WebSocket 端点用于实时事件推送,支持订阅、取消订阅和回放。

// src/gateway/agent-event-gateway.ts
// Agent API WebSocket 事件网关

import WebSocket from 'ws';
import { IncomingMessage } from 'http';
import { TendermintWSClient } from '../websocket/tendermint-ws';

interface AgentSubscription {
  id: string;
  agentId: string;
  ws: WebSocket;
  filters: {
    eventTypes?: string[];
    actions?: string[];
  };
}

export class AgentEventGateway {
  private tendermintWS: TendermintWSClient;
  private subscriptions = new Map<string, AgentSubscription>();
  private wsServer: WebSocket.Server | null = null;

  constructor(rpcEndpoint: string) {
    this.tendermintWS = new TendermintWSClient(rpcEndpoint + '/websocket');
  }

  async start(port = 8090): Promise<void> {
    await this.tendermintWS.connect();

    // 订阅所有链上事件
    await this.tendermintWS.subscribe("tm.event='Tx'");
    this.tendermintWS.on("tm.event='Tx'", (data: any) => {
      const events = data.value.TxResult.result.events || [];
      this.broadcastToSubscribers(events);
    });

    // 启动 Agent API WebSocket 服务器
    this.wsServer = new WebSocket.Server({ port });
    console.log('Agent 事件网关运行在 ws://localhost:' + port);

    this.wsServer.on('connection', (ws: WebSocket, req: IncomingMessage) => {
      this.handleConnection(ws, req);
    });
  }

  private handleConnection(ws: WebSocket, req: IncomingMessage): void {
    ws.on('message', async (raw: string) => {
      try {
        const msg = JSON.parse(raw.toString());

        switch (msg.type) {
          case 'subscribe':
            await this.handleSubscribe(ws, msg);
            break;
          case 'unsubscribe':
            this.handleUnsubscribe(ws, msg);
            break;
          case 'replay':
            await this.handleReplay(ws, msg);
            break;
          default:
            ws.send(JSON.stringify({ type: 'error', message: '未知消息类型' }));
        }
      } catch (error) {
        ws.send(JSON.stringify({ type: 'error', message: '消息解析失败' }));
      }
    });

    ws.on('close', () => this.handleDisconnect(ws));

    ws.send(JSON.stringify({
      type: 'connected',
      version: '1.0.0',
      chainId: 'msg-chain-1',
      // X-Local-Only=true
      capabilities: { replay: true, filterByAgent: true, filterByAction: true },
    }));
  }

  private async handleSubscribe(ws: WebSocket, msg: any): Promise<void> {
    const subId = msg.id || `sub-${Date.now()}`;
    const subscription: AgentSubscription = {
      id: subId,
      agentId: msg.agent_id || '',
      ws,
      filters: {
        eventTypes: msg.filters?.eventTypes,
        actions: msg.filters?.actions,
      },
    };
    this.subscriptions.set(subId, subscription);

    ws.send(JSON.stringify({ type: 'subscribed', id: subId }));
  }

  private handleUnsubscribe(ws: WebSocket, msg: any): void {
    const subId = msg.id;
    if (subId && this.subscriptions.has(subId)) {
      this.subscriptions.delete(subId);
      ws.send(JSON.stringify({ type: 'unsubscribed', id: subId }));
    }
  }

  private async handleReplay(ws: WebSocket, msg: any): Promise<void> {
    // X-Local-Only=true: 回放功能仅部分实现
    const { fromHeight, toHeight, agentId } = msg;

    ws.send(JSON.stringify({
      type: 'replay_start',
      fromHeight,
      toHeight,
      agentId,
    }));

    // TODO: 从数据库查询历史事件并回放
    // 此处需要访问 events 表

    ws.send(JSON.stringify({
      type: 'replay_complete',
      fromHeight,
      toHeight,
      eventCount: 0,
    }));
  }

  private handleDisconnect(ws: WebSocket): void {
    for (const [id, sub] of this.subscriptions) {
      if (sub.ws === ws) {
        this.subscriptions.delete(id);
      }
    }
  }

  private broadcastToSubscribers(events: any[]): void {
    for (const [id, sub] of this.subscriptions) {
      if (sub.ws.readyState !== WebSocket.OPEN) continue;

      const filtered = events.filter((ev: any) => {
        if (ev.type !== 'wasm') return false;

        const attrs = this.attributesToMap(ev.attributes);

        // 按 agent_id 过滤
        if (sub.agentId && attrs['agent_id'] !== sub.agentId) return false;

        // 按 action 过滤
        if (sub.filters.actions && sub.filters.actions.length > 0) {
          if (!sub.filters.actions.includes(attrs['action'])) return false;
        }

        return true;
      });

      if (filtered.length > 0) {
        sub.ws.send(JSON.stringify({
          type: 'event',
          subscriptionId: id,
          events: filtered,
        }));
      }
    }
  }

  private attributesToMap(attributes: { key: string; value: string }[]): Record<string, string> {
    const map: Record<string, string> = {};
    for (const attr of attributes) {
      map[attr.key] = attr.value;
    }
    return map;
  }

  async stop(): Promise<void> {
    await this.tendermintWS.disconnect();
    this.wsServer?.close();
  }
}

8.2 Agent API 事件端点

Agent API 通过 HTTP 端点管理事件订阅(来自 agent_surface.yaml):

方法 端点 说明
POST /api/v1/agent/events/subscribe WebSocket 升级 — 建立事件连接
POST /api/v1/agent/events/unsubscribe 取消订阅
POST /api/v1/agent/events/replay 回放历史事件

客户端连接示例:

// Agent API WebSocket 客户端

const ws = new WebSocket('ws://localhost:8090');

ws.onopen = () => {
  // 订阅特定 Agent 的事件
  ws.send(JSON.stringify({
    type: 'subscribe',
    id: 'my-sub-1',
    agent_id: 'agent-alpha',
    filters: {
      actions: ['CreateSession', 'ReleasePayment', 'DisputePayment'],
    },
  }));
};

ws.onmessage = (event) => {
  const msg = JSON.parse(event.data);
  if (msg.type === 'event') {
    console.log('收到 Agent 事件:', msg.events);
  }
};

// 回放历史事件
ws.send(JSON.stringify({
  type: 'replay',
  fromHeight: 1000,
  toHeight: 1100,
  agent_id: 'agent-alpha',
}));

9. 实时订阅前端

9.1 Apollo Client 设置

// src/frontend/apollo-setup.ts
// React + Apollo Client GraphQL 订阅设置

import { ApolloClient, InMemoryCache, split } from '@apollo/client';
import { GraphQLWsLink } from '@apollo/client/link/subscriptions';
import { createClient } from 'graphql-ws';
import { HttpLink } from '@apollo/client/link/http';
import { getMainDefinition } from '@apollo/client/utilities';

const httpLink = new HttpLink({
  uri: 'http://localhost:4000/graphql',
});

const wsLink = new GraphQLWsLink(createClient({
  url: 'ws://localhost:4000/graphql',
  connectionParams: {},
  retryAttempts: 10,
  shouldRetry: () => true,
  on: {
    connected: () => console.log('GraphQL WebSocket 已连接'),
    disconnected: () => console.log('GraphQL WebSocket 已断开'),
    error: (err) => console.error('GraphQL WebSocket 错误:', err),
  },
}));

const splitLink = split(
  ({ query }) => {
    const definition = getMainDefinition(query);
    return definition.kind === 'OperationDefinition' && definition.operation === 'subscription';
  },
  wsLink,
  httpLink,
);

export const client = new ApolloClient({
  link: splitLink,
  cache: new InMemoryCache(),
});

9.2 React Subscription Hooks

// src/frontend/hooks/useNewBlock.ts

import { gql, useSubscription } from '@apollo/client';

const NEW_BLOCK_SUBSCRIPTION = gql`
  subscription OnNewBlock {
    newBlock {
      height
      hash
      time
      proposer
      txCount
    }
  }
`;

export function useNewBlockSubscription() {
  return useSubscription(NEW_BLOCK_SUBSCRIPTION);
}

// ─── useNewTransaction ──────────────────────────

const NEW_TRANSACTION_SUBSCRIPTION = gql`
  subscription OnNewTransaction {
    newTransaction {
      hash
      blockHeight
      status
      gasUsed
    }
  }
`;

export function useNewTransactionSubscription() {
  return useSubscription(NEW_TRANSACTION_SUBSCRIPTION);
}

// ─── useNewEvent ────────────────────────────────

const NEW_EVENT_SUBSCRIPTION = gql`
  subscription OnNewEvent($eventType: String) {
    newEvent(eventType: $eventType) {
      txHash
      blockHeight
      action
      contractAddress
      attributes
    }
  }
`;

export function useNewEventSubscription(eventType?: string) {
  return useSubscription(NEW_EVENT_SUBSCRIPTION, {
    variables: { eventType },
  });
}

// ─── useAgentStatusChanged ──────────────────────

const AGENT_STATUS_SUBSCRIPTION = gql`
  subscription OnAgentStatusChanged {
    agentStatusChanged {
      agentId
      previousStatus
      newStatus
      txHash
    }
  }
`;

export function useAgentStatusSubscription() {
  return useSubscription(AGENT_STATUS_SUBSCRIPTION);
}

// ─── useNewSession ──────────────────────────────

const NEW_SESSION_SUBSCRIPTION = gql`
  subscription OnNewSession {
    newSession {
      sessionId
      agentId
      client
      status
      amount
    }
  }
`;

export function useNewSessionSubscription() {
  return useSubscription(NEW_SESSION_SUBSCRIPTION);
}

9.3 实时交易 Feed 组件

// src/frontend/components/LiveTransactionFeed.tsx

import React, { useEffect, useState } from 'react';
import { useNewBlockSubscription, useNewTransactionSubscription } from '../hooks/useNewBlock';

interface Transaction {
  hash: string;
  blockHeight: number;
  status: string;
  gasUsed: number;
  fee: string;
}

interface Block {
  height: number;
  hash: string;
  time: string;
  txCount: number;
}

export function LiveTransactionFeed() {
  const [recentTxs, setRecentTxs] = useState<Transaction[]>([]);
  const [recentBlocks, setRecentBlocks] = useState<Block[]>([]);
  const maxItems = 20;

  const { data: blockData } = useNewBlockSubscription();
  const { data: txData } = useNewTransactionSubscription();

  useEffect(() => {
    if (blockData?.newBlock) {
      setRecentBlocks(prev => {
        const updated = [blockData.newBlock, ...prev].slice(0, maxItems);
        return updated;
      });
    }
  }, [blockData]);

  useEffect(() => {
    if (txData?.newTransaction) {
      setRecentTxs(prev => {
        const updated = [txData.newTransaction, ...prev].slice(0, maxItems);
        return updated;
      });
    }
  }, [txData]);

  return (
    <div className="live-feed">
      <h2>实时交易 Feed</h2>

      <div className="recent-blocks">
        <h3>最新区块</h3>
        <div className="block-list">
          {recentBlocks.map(block => (
            <div key={block.height} className="block-item">
              <span className="block-height">#{block.height}</span>
              <span className="block-time">{new Date(block.time).toLocaleTimeString()}</span>
              <span className="block-tx-count">{block.txCount} txs</span>
            </div>
          ))}
        </div>
      </div>

      <div className="recent-transactions">
        <h3>最新交易</h3>
        <div className="tx-list">
          {recentTxs.map(tx => (
            <div key={tx.hash} className="tx-item">
              <span className={`tx-status ${tx.status.toLowerCase()}`}>
                {tx.status}
              </span>
              <span className="tx-hash">{tx.hash.slice(0, 16)}...</span>
              <span className="tx-block">#{tx.blockHeight}</span>
              <span className="tx-gas">{tx.gasUsed} gas</span>
            </div>
          ))}
        </div>
      </div>
    </div>
  );
}

9.4 Agent 实时看板

// src/frontend/components/AgentDashboard.tsx

import React, { useState, useEffect, useCallback } from 'react';
import {
  useNewEventSubscription,
  useAgentStatusSubscription,
  useNewSessionSubscription,
} from '../hooks/useNewBlock';

interface AgentEvent {
  txHash: string;
  blockHeight: number;
  action: string;
  contractAddress: string;
  attributes: Record<string, string>;
}

export function AgentDashboard() {
  const [selectedAgent, setSelectedAgent] = useState<string>('');
  const [events, setEvents] = useState<AgentEvent[]>([]);
  const [eventFilter, setEventFilter] = useState<string>('');

  const { data: eventData } = useNewEventSubscription(eventFilter || undefined);
  const { data: statusData } = useAgentStatusSubscription();
  const { data: sessionData } = useNewSessionSubscription();

  useEffect(() => {
    if (eventData?.newEvent) {
      setEvents(prev => [eventData.newEvent, ...prev].slice(0, 100));
    }
  }, [eventData]);

  const filteredEvents = selectedAgent
    ? events.filter(e => e.attributes['agent_id'] === selectedAgent)
    : events;

  return (
    <div className="agent-dashboard">
      <h2>Agent 实时看板</h2>

      <div className="controls">
        <input
          type="text"
          placeholder="按 Agent ID 过滤"
          value={selectedAgent}
          onChange={e => setSelectedAgent(e.target.value)}
        />
        <select value={eventFilter} onChange={e => setEventFilter(e.target.value)}>
          <option value="">所有事件</option>
          <option value="CreateSession">CreateSession</option>
          <option value="ReleasePayment">ReleasePayment</option>
          <option value="DisputePayment">DisputePayment</option>
          <option value="RegisterAgent">RegisterAgent</option>
          <option value="SetStatus">SetStatus</option>
        </select>
      </div>

      {statusData?.agentStatusChanged && (
        <div className="status-alert">
          Agent {statusData.agentStatusChanged.agentId}:
          {statusData.agentStatusChanged.previousStatus} →
          {statusData.agentStatusChanged.newStatus}
        </div>
      )}

      {sessionData?.newSession && (
        <div className="session-alert">
          新会话: {sessionData.newSession.sessionId} -
          {sessionData.newSession.client} - {sessionData.newSession.amount} umsg
        </div>
      )}

      <div className="event-stream">
        <h3>事件流 ({filteredEvents.length})</h3>
        <div className="event-list">
          {filteredEvents.map((evt, i) => (
            <div key={`${evt.txHash}-${i}`} className="event-item">
              <span className="event-action">{evt.action}</span>
              <span className="event-block">#{evt.blockHeight}</span>
              <span className="event-agent">{evt.attributes['agent_id'] || '-'}</span>
              <span className="event-tx">{evt.txHash.slice(0, 12)}...</span>
            </div>
          ))}
        </div>
      </div>
    </div>
  );
}

10. 搜索服务

10.1 PostgreSQL 全文搜索设置

全文搜索基于 pg_trgm 扩展实现模糊匹配和 tsvector 全文索引。

-- 已包含在 001_schema.sql 中的搜索索引

-- trigram 模糊搜索
CREATE INDEX idx_blocks_hash_trgm ON blocks USING GIN (hash gin_trgm_ops);
CREATE INDEX idx_transactions_hash_trgm ON transactions USING GIN (hash gin_trgm_ops);
CREATE INDEX idx_contracts_address_trgm ON contracts USING GIN (address gin_trgm_ops);
CREATE INDEX idx_agents_id_trgm ON agents USING GIN (agent_id gin_trgm_ops);
CREATE INDEX idx_agents_name_trgm ON agents USING GIN (name gin_trgm_ops);

-- tsvector 全文搜索(memo 字段)
ALTER TABLE transactions ADD COLUMN search_vector tsvector
  GENERATED ALWAYS AS (to_tsvector('simple', coalesce(memo, ''))) STORED;
CREATE INDEX idx_transactions_search ON transactions USING GIN (search_vector);

10.2 搜索解析器

// src/search/search-service.ts

interface SearchInput {
  q: string;
  types?: ('block' | 'transaction' | 'contract' | 'agent' | 'session')[];
  limit?: number;
}

interface SearchResultItem {
  type: string;
  id: string;
  title: string;
  subtitle: string;
  route: string;
  matchField: string;
  score: number;
}

interface SearchResponse {
  items: SearchResultItem[];
  total: number;
  query: string;
}

export class SearchService {
  constructor(private db: DatabaseService) {}

  async search(input: SearchInput): Promise<SearchResponse> {
    const query = input.q.trim();
    const types = input.types || ['block', 'transaction', 'contract', 'agent', 'session'];
    const limit = input.limit || 20;

    if (!query || query.length < 1) {
      return { items: [], total: 0, query };
    }

    const results: SearchResultItem[] = [];
    const pattern = '%' + query + '%';

    if (types.includes('block')) {
      const blocks = await this.db['pool'].query(
        `SELECT height::text as id, height, hash, time FROM blocks
         WHERE hash ILIKE $1 OR height::text ILIKE $2
         LIMIT $3`,
        [pattern, query, limit]
      );
      for (const b of blocks.rows) {
        results.push({
          type: 'block',
          id: b.id,
          title: '区块 #' + b.height,
          subtitle: b.hash.slice(0, 20) + '...',
          route: '/blocks/' + b.height,
          matchField: b.hash.startsWith(query) ? 'hash' : 'height',
          score: b.hash === query ? 100 : 50,
        });
      }
    }

    if (types.includes('transaction')) {
      const txs = await this.db['pool'].query(
        `SELECT hash, block_height, status, memo FROM transactions
         WHERE hash ILIKE $1
         LIMIT $2`,
        [pattern, limit]
      );
      for (const t of txs.rows) {
        results.push({
          type: 'transaction',
          id: t.hash,
          title: t.hash.slice(0, 20) + '...',
          subtitle: '区块 #' + t.block_height + ' | ' + t.status,
          route: '/txs/' + t.hash,
          matchField: 'hash',
          score: t.hash === query ? 95 : 45,
        });
      }
    }

    if (types.includes('contract')) {
      const contracts = await this.db['pool'].query(
        `SELECT address, label, contract_type FROM contracts
         WHERE address ILIKE $1 OR label ILIKE $2
         LIMIT $3`,
        [pattern, '%' + query + '%', limit]
      );
      for (const c of contracts.rows) {
        results.push({
          type: 'contract',
          id: c.address,
          title: c.label || c.address.slice(0, 20) + '...',
          subtitle: c.contract_type || c.address,
          route: '/contracts/' + c.address,
          matchField: c.address.startsWith(query) ? 'address' : 'label',
          score: c.address === query ? 90 : 40,
        });
      }
    }

    if (types.includes('agent')) {
      const agents = await this.db['pool'].query(
        `SELECT agent_id, name, owner, status FROM agents
         WHERE agent_id ILIKE $1 OR name ILIKE $2 OR owner ILIKE $3
         LIMIT $4`,
        [pattern, '%' + query + '%', pattern, limit]
      );
      for (const a of agents.rows) {
        results.push({
          type: 'agent',
          id: a.agent_id,
          title: a.name,
          subtitle: a.agent_id + ' | ' + a.status,
          route: '/agents/' + a.agent_id,
          matchField: a.agent_id === query ? 'agent_id' : 'name',
          score: a.agent_id === query ? 85 : 35,
        });
      }
    }

    if (types.includes('session')) {
      const sessions = await this.db['pool'].query(
        `SELECT session_id, agent_id, client, status FROM sessions
         WHERE session_id ILIKE $1
         LIMIT $2`,
        [pattern, limit]
      );
      for (const s of sessions.rows) {
        results.push({
          type: 'session',
          id: s.session_id,
          title: s.session_id,
          subtitle: s.agent_id + ' | ' + s.client + ' | ' + s.status,
          route: '/sessions/' + s.session_id,
          matchField: 'session_id',
          score: s.session_id === query ? 80 : 30,
        });
      }
    }

    // 按分数排序
    results.sort((a, b) => b.score - a.score);
    const top = results.slice(0, limit);

    return {
      items: top,
      total: results.length,
      query,
    };
  }

  async suggest(query: string): Promise<string[]> {
    if (query.length < 2) return [];

    const suggestions: Set<string> = new Set();

    // Agent 名称建议
    const agents = await this.db['pool'].query(
      `SELECT name FROM agents WHERE name ILIKE $1 LIMIT 5`,
      ['%' + query + '%']
    );
    for (const a of agents.rows) suggestions.add(a.name);

    // 区块高度建议
    if (/^\d+$/.test(query)) {
      suggestions.add('区块 #' + query);
    }

    return Array.from(suggestions).slice(0, 5);
  }
}

10.3 搜索 API Resolver

// 在 GraphQL resolvers 中添加

search: async (_: any, { q }: { q: string }, { db }: Context) => {
  const searchService = new SearchService(db as any);
  const result = await searchService.search({ q, limit: 20 });

  return {
    items: result.items.map(item => {
      switch (item.type) {
        case 'block': return { height: parseInt(item.id), hash: item.subtitle.replace('...', ''), time: '' };
        case 'transaction': return { hash: item.id, blockHeight: 0, status: '' };
        case 'agent': return { agentId: item.id, name: item.title, owner: '' };
        case 'contract': return { address: item.id, label: item.title, contractType: '' };
        case 'session': return { sessionId: item.id, agentId: '', client: '', status: '' };
        default: return null;
      }
    }).filter(Boolean),
    stats: {
      blocks: result.items.filter(i => i.type === 'block').length,
      transactions: result.items.filter(i => i.type === 'transaction').length,
      contracts: result.items.filter(i => i.type === 'contract').length,
      agents: result.items.filter(i => i.type === 'agent').length,
      sessions: result.items.filter(i => i.type === 'session').length,
    },
    truncated: result.total > 20,
  };
},

11. Docker 部署

11.1 Docker Compose

# docker-compose.yml

version: '3.8'

services:
  postgres:
    image: postgres:15-alpine
    container_name: msgchain-pg
    environment:
      POSTGRES_USER: msgchain
      POSTGRES_PASSWORD: msgchain_secret
      POSTGRES_DB: msgchain_indexer
    ports:
      - "5432:5432"
    volumes:
      - pgdata:/var/lib/postgresql/data
      - ./migrations:/docker-entrypoint-initdb.d
    healthcheck:
      test: ["CMD-SHELL", "pg_isready -U msgchain"]
      interval: 10s
      timeout: 5s
      retries: 5
    networks:
      - msgchain-net

  indexer:
    build:
      context: .
      dockerfile: Dockerfile.indexer
    container_name: msgchain-indexer
    environment:
      RPC_ENDPOINT: http://host.docker.internal:26657
      REST_ENDPOINT: http://host.docker.internal:1317
      DATABASE_URL: postgresql://msgchain:msgchain_secret@postgres:5432/msgchain_indexer
      INDEXER_MODE: websocket
      START_HEIGHT: 0
      LOG_LEVEL: info
    depends_on:
      postgres:
        condition: service_healthy
    restart: unless-stopped
    networks:
      - msgchain-net

  graphql-api:
    build:
      context: .
      dockerfile: Dockerfile.api
    container_name: msgchain-graphql
    environment:
      DATABASE_URL: postgresql://msgchain:msgchain_secret@postgres:5432/msgchain_indexer
      RPC_ENDPOINT: http://host.docker.internal:26657
      PORT: 4000
      CORS_ORIGIN: "*"
    ports:
      - "4000:4000"
    depends_on:
      postgres:
        condition: service_healthy
    restart: unless-stopped
    networks:
      - msgchain-net

  agent-gateway:
    build:
      context: .
      dockerfile: Dockerfile.gateway
    container_name: msgchain-gateway
    environment:
      RPC_ENDPOINT: http://host.docker.internal:26657
      GATEWAY_PORT: 8090
    ports:
      - "8090:8090"
    depends_on:
      - indexer
    restart: unless-stopped
    networks:
      - msgchain-net

  frontend:
    build:
      context: ./frontend
      dockerfile: Dockerfile
    container_name: msgchain-frontend
    environment:
      REACT_APP_GRAPHQL_URL: http://localhost:4000/graphql
      REACT_APP_WS_URL: ws://localhost:4000/graphql
    ports:
      - "3000:80"
    depends_on:
      - graphql-api
    networks:
      - msgchain-net

  nginx:
    image: nginx:alpine
    container_name: msgchain-nginx
    ports:
      - "80:80"
    volumes:
      - ./nginx.conf:/etc/nginx/nginx.conf
    depends_on:
      - graphql-api
      - frontend
    networks:
      - msgchain-net

  redis:
    image: redis:7-alpine
    container_name: msgchain-redis
    ports:
      - "6379:6379"
    volumes:
      - redisdata:/data
    networks:
      - msgchain-net

volumes:
  pgdata:
  redisdata:

networks:
  msgchain-net:
    driver: bridge

11.2 Dockerfile

# Dockerfile.indexer
FROM node:20-alpine AS builder
WORKDIR /app
COPY package*.json tsconfig.json ./
COPY src/ ./src/
RUN npm ci && npm run build

FROM node:20-alpine
WORKDIR /app
COPY --from=builder /app/dist ./dist
COPY --from=builder /app/node_modules ./node_modules
COPY migrations/ ./migrations/
EXPOSE 4001
CMD ["node", "dist/indexer/index.js"]

# Dockerfile.api
FROM node:20-alpine AS builder
WORKDIR /app
COPY package*.json tsconfig.json ./
COPY src/ ./src/
RUN npm ci && npm run build

FROM node:20-alpine
WORKDIR /app
COPY --from=builder /app/dist ./dist
COPY --from=builder /app/node_modules ./node_modules
EXPOSE 4000
CMD ["node", "dist/graphql/server.js"]

# Dockerfile.gateway
FROM node:20-alpine AS builder
WORKDIR /app
COPY package*.json tsconfig.json ./
COPY src/ ./src/
RUN npm ci && npm run build

FROM node:20-alpine
WORKDIR /app
COPY --from=builder /app/dist ./dist
COPY --from=builder /app/node_modules ./node_modules
EXPOSE 8090
CMD ["node", "dist/gateway/agent-event-gateway.js"]

11.3 Nginx 配置

# nginx.conf

events {
  worker_connections 1024;
}

http {
  upstream graphql {
    server graphql-api:4000;
  }

  upstream frontend {
    server frontend:80;
  }

  server {
    listen 80;

    # GraphQL API
    location /graphql {
      proxy_pass http://graphql;
      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 86400s;
    }

    # Frontend
    location / {
      proxy_pass http://frontend;
      proxy_set_header Host $host;
      proxy_set_header X-Real-IP $remote_addr;
    }
  }
}

11.4 部署命令

# 启动所有服务
docker-compose up -d

# 查看日志
docker-compose logs -f indexer
docker-compose logs -f graphql-api

# 执行数据库迁移
docker-compose exec postgres psql -U msgchain -d msgchain_indexer -f /migrations/001_schema.sql
docker-compose exec postgres psql -U msgchain -d msgchain_indexer -f /migrations/002_seed.sql
docker-compose exec postgres psql -U msgchain -d msgchain_indexer -f /migrations/003_indexes.sql

# 停止服务
docker-compose down

# 扩展 Indexer 实例(如果使用轮询模式)
docker-compose up -d --scale indexer=3

12. 性能与扩展

12.1 数据库索引策略

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

按高度查询区块       → PRIMARY KEY (height)
按哈希查询交易       → PRIMARY KEY (hash) + GIN trigram
按 action 查事件     → idx_events_action
按合约地址查事件     → idx_events_contract_address
按合约地址+action   → idx_events_contract_action (复合索引)
JSONB 属性查询       → idx_events_attributes_gin (GIN jsonb_path_ops)
最新区块列表         → idx_blocks_time (DESC)
Agent 的所有会话     → idx_sessions_agent_status (agent_id, status)
模糊搜索             → GIN trigram 索引

12.2 连接池配置

// pg 连接池优化配置

const poolConfig = {
  max: 20,                    // 最大连接数
  min: 4,                     // 最小空闲连接
  idleTimeoutMillis: 30000,   // 空闲超时
  acquireConnectionTimeout: 5000,
  createTimeoutMillis: 5000,
  destroyTimeoutMillis: 5000,
  reapIntervalMillis: 1000,   // 空闲连接清理间隔
  // 用于大数据量插入的配置
  statement_timeout: 30000,   // 单条语句超时
  query_timeout: 30000,       // 查询超时
};

// 写入连接池(用于 Indexer,需要更多连接)
const writePool = new Pool({ ...poolConfig, max: 30 });

// 读取连接池(用于 GraphQL,需要更少但更快)
const readPool = new Pool({ ...poolConfig, max: 10 });

12.3 查询优化

// 使用 LIMIT 防止全表扫描
const QUERY_TEMPLATES = {
  recentBlocks: 'SELECT * FROM blocks ORDER BY height DESC LIMIT $1',
  eventsByAction: `
    SELECT e.* FROM events e
    WHERE e.action = $1
    ORDER BY e.block_height DESC
    LIMIT $2 OFFSET $3
  `,
  // 使用覆盖索引减少回表
  agentSummary: `
    SELECT agent_id, name, status, agent_type
    FROM agents
    WHERE status = 'ACTIVE'
    ORDER BY created_at DESC
    LIMIT $1
  `,
  // 使用 JOIN 代替 N+1
  transactionsWithEvents: `
    SELECT t.*, json_agg(e.*) as events
    FROM transactions t
    LEFT JOIN events e ON e.tx_hash = t.hash
    WHERE t.block_height = $1
    GROUP BY t.hash
    ORDER BY t.index ASC
  `,
};

12.4 Redis 缓存

// src/cache/redis-cache.ts

import Redis from 'ioredis';

export class CacheService {
  private redis: Redis;
  private defaultTTL = 60; // 默认 60 秒缓存

  constructor(redisUrl: string) {
    this.redis = new Redis(redisUrl);
  }

  async getOrSet<T>(
    key: string,
    fetcher: () => Promise<T>,
    ttl = this.defaultTTL,
  ): Promise<T> {
    const cached = await this.redis.get(key);
    if (cached) {
      return JSON.parse(cached);
    }

    const value = await fetcher();
    await this.redis.setex(key, ttl, JSON.stringify(value));
    return value;
  }

  // 缓存键命名空间
  keys = {
    block: (height: number) => `block:${height}`,
    tx: (hash: string) => `tx:${hash}`,
    agent: (id: string) => `agent:${id}`,
    chainSummary: 'chain:summary',
    latestHeight: 'chain:latestHeight',
    // 缓存失效模式:按模式批量清除
    pattern: {
      block: 'block:*',
      tx: 'tx:*',
    },
  };

  async invalidatePattern(pattern: string): Promise<void> {
    const keys = await this.redis.keys(pattern);
    if (keys.length > 0) {
      await this.redis.del(...keys);
    }
  }
}

// GraphQL DataLoader 配合缓存
const cacheService = new CacheService('redis://redis:6379');

const resolvers = {
  Query: {
    block: async (_: any, { height }: { height: number }, ctx: Context) => {
      return cacheService.getOrSet(
        `block:${height}`,
        () => ctx.db.getBlock(height),
        30, // 区块数据缓存 30 秒
      );
    },
    chainSummary: async (_: any, __: any, ctx: Context) => {
      return cacheService.getOrSet(
        'chain:summary',
        () => ctx.db.queryChainSummary(),
        10, // 摘要缓存 10 秒
      );
    },
  },
};

12.5 分片考虑

当数据量增长到以下规模时需要考虑分片:

指标 阈值 策略
events 表行数 > 1 亿 按 block_height 范围分区 (每 10 万块一个分区)
transactions 表 > 5000 万 按月分区
agents 表 > 10 万 无需分片,当前 PG 单表可支持
sessions 表 > 100 万 按 agent_id hash 分片

事件表分区示例:

-- 按区块范围分区
CREATE TABLE events_partitioned (
  LIKE events INCLUDING ALL
) PARTITION BY RANGE (block_height);

CREATE TABLE events_0_100000 PARTITION OF events_partitioned
  FOR VALUES FROM (0) TO (100000);
CREATE TABLE events_100000_200000 PARTITION OF events_partitioned
  FOR VALUES FROM (100000) TO (200000);
CREATE TABLE events_200000_300000 PARTITION OF events_partitioned
  FOR VALUES FROM (200000) TO (300000);

-- 创建默认分区
CREATE TABLE events_default PARTITION OF events_partitioned DEFAULT;

12.6 当前边界

// X-Local-Only=true — 当前实现边界

const BOUNDARIES = {
  // Indexer API
  localPreflightOnly: true,
  noPersistentStorage: '使用 BadgerDB 本地存储,非开发参考级',
  noRetentionProof: '保留证明未实现',
  searchLimitedToPrefix: '搜索仅限 ILIKE 前缀匹配,无 NLP',

  // 索引器
  noReorgHandling: '生产环境需要实现深度重组检测和回滚',
  noEventDedup: '内存去重,重启后可能重复处理',
  noBackfill: '不支持从创世块回填历史数据',

  // GraphQL
  noRateLimiting: '生产环境需要添加 rate limiting',
  noAuth: '当前无身份验证和授权',
  noCostAnalysis: 'GraphQL 查询成本分析未实现',

  // 性能
  noSharding: '尚未实现数据库分片',
  noCDN: '静态资产无需 CDN',
  maxBatchSize: 500,
  maxQueryDepth: 10,
};

12.7 性能基准

操作 延迟 (P99) 吞吐量 数据规模
按高度查区块 < 5ms > 1000/s 1000 万块
按哈希查交易 < 10ms > 500/s 1 亿交易
按 action 查事件 < 50ms > 200/s 5 亿事件
搜索 (ILIKE) < 100ms > 100/s 100 万 Agent
批量写入 (500条) < 200ms > 50 batches/s —
新块订阅推送 < 100ms 实时 5s 区块间隔

附录

A. 项目结构

src/
├── indexer/
│   ├── indexer-service.ts      # Indexer 主服务
│   ├── block-fetcher.ts        # 区块获取
│   ├── event-pipeline.ts       # 事件管道
│   ├── deduplicator.ts         # 事件去重
│   ├── reorg-handler.ts        # 重组处理
│   └── batch-writer.ts         # 批量写入
├── websocket/
│   └── tendermint-ws.ts        # Tendermint WS 客户端
├── events/
│   ├── event-normalizer.ts     # 事件标准化
│   └── schemas/
│       ├── agent-payment.ts    # 支付事件类型
│       ├── agent-registry.ts   # 注册事件类型
│       ├── did-registry.ts     # DID 事件类型
│       ├── constitution.ts     # 治理事件类型
│       └── micropayment.ts     # 微支付事件类型
├── db/
│   └── database-service.ts     # 数据库服务
├── graphql/
│   ├── server.ts               # GraphQL 服务器
│   ├── schema.ts               # Schema 定义
│   ├── resolvers.ts            # Resolvers
│   └── subscription-integration.ts
├── gateway/
│   └── agent-event-gateway.ts  # Agent 事件网关
├── cache/
│   └── redis-cache.ts          # Redis 缓存
├── search/
│   └── search-service.ts       # 搜索服务
└── frontend/
    ├── apollo-setup.ts
    ├── hooks/
    │   └── useNewBlock.ts
    └── components/
        ├── LiveTransactionFeed.tsx
        └── AgentDashboard.tsx

migrations/
├── 001_schema.sql
├── 002_seed.sql
├── 003_indexes.sql
└── 004_retention.sql

docker-compose.yml
Dockerfile.indexer
Dockerfile.api
Dockerfile.gateway
nginx.conf

B. 快速启动

# 1. 启动 PostgreSQL
docker-compose up -d postgres

# 2. 运行迁移
docker-compose exec postgres psql -U msgchain -d msgchain_indexer -f /migrations/001_schema.sql
docker-compose exec postgres psql -U msgchain -d msgchain_indexer -f /migrations/002_seed.sql

# 3. 启动 Indexer
docker-compose up -d indexer

# 4. 启动 GraphQL API
docker-compose up -d graphql-api

# 5. 访问 GraphQL Playground
# http://localhost:4000/graphql

# 6. 测试查询
curl -X POST http://localhost:4000/graphql \
  -H "Content-Type: application/json" \
  -d '{"query": "{ chainSummary { latestBlockHeight totalBlocks totalTransactions } }"}'

C. 相关链接

资源 地址
Tendermint RPC http://localhost:26657
Tendermint WebSocket ws://localhost:26657/websocket
REST API http://localhost:1317
gRPC localhost:9090
GraphQL API http://localhost:4000/graphql
GraphQL Subscription ws://localhost:4000/graphql
Agent 事件网关 ws://localhost:8090
Indexer Health http://localhost:4001/api/v1/indexer/health
Explorer Search http://localhost:1317/api/v1/search?q=

本文档由 MSG Chain 开发团队维护。所有代码示例基于 TypeScript 5.x + PostgreSQL 15 + GraphQL Yoga v3。
当前状态: local_preflight_only=true — 仅限本地预检,不可用于生产环境。