MSG Chain 事件索引与 GraphQL API 完全指南
数据来源:MSG Chain 代码库核实
主网状态: No-Go — 当前 MSGChain 主网裁决为 No-Go,以下内容反映代码实际状态,不代表生产可用。
目录
- 概述
- Tendermint WebSocket 事件
- CosmWasm 合约事件
- Indexer 数据管道
- PostgreSQL 数据模型
- GraphQL API 服务
- Indexer 事件索引
- Agent API 事件网关
- 实时订阅前端
- 搜索服务
- Docker 部署
- 性能与扩展
1. 概述
1.1 为什么需要事件索引
MSG Chain 是一条基于 Cosmos SDK + CosmWasm 的应用链。链上合约(Agent Payment、Agent Registry、DID Registry、AI Constitution、Micropayment Session)会持续产生大量事件。原始链上事件通过 Tendermint WebSocket 暴露,但存在以下问题:
- 无持久化:WebSocket 事件是瞬时的,断连后历史事件丢失
- 无查询能力:原始 RPC 只支持按 tx hash 或 block height 检索,不支持事件类型过滤
- 无结构化数据:事件以扁平
{key, value}数组形式存在,需要二次解析 - 无聚合视图:无法直接查询"某个 Agent 的所有支付会话"或"某个 DID 的所有变更历史"
事件索引系统(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 内置的完整重连策略:
- 指数退避:初始 1s,乘以 2,最大 30s
- 最大重试次数:10 次,之后触发
reconnect_failed - 自动恢复订阅:重连后重新订阅所有之前的查询
- Ping/Pong:每 30s 发送 ping 维持连接
- 断开检测: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— 仅限本地预检,不可用于生产环境。
