https://github.com/natreeum/HybridIndexer
기존 인덱싱 방식은 http 방식(Polling)으로 진행. 이로인해 설정된 주기로만 요청하므로 최대 ‘설정된 주기까지의 delay’가 생김. 이로 인해 낙관적 업데이트와 같은 추가 작업을 필요로 합니다.
Event 를 emit 하는 단순한 컨트랙트 배포
code : Test.sol
// SPDX-License-Identifier: MIT
pragma solidity ^0.8.0;
contract Test {
event LogMessage(string message);
event TestTransfer(address indexed from, address indexed to, uint256 value);
function emitLogMessage(string memory message) public {
emit LogMessage(message);
}
function transfer(address to, uint256 value) public {
emit TestTransfer(msg.sender, to, value);
}
}
code : indexer.js
const { ethers } = require("ethers");
module.exports = class HybridIndexer {
constructor({
wsUrl,
httpUrl,
contractCfg, // {address, abi}
events, // ["Transfer", ...]
confirmations = 10,
chunkSize = 2000,
poolIntervalMs = 5000,
}) {
this.wsUrl = wsUrl;
this.httpUrl = httpUrl;
this.contractCfg = contractCfg;
this.events = events;
this.confirmations = confirmations;
this.chunkSize = chunkSize;
this.poolIntervalMs = poolIntervalMs;
this.wsProvider = null;
this.httpProvider = null;
this.wsContract = null;
this.httpContract = null;
this.running = false;
}
async start() {
this.running = true;
this.wsProvider = new ethers.WebSocketProvider(this.wsUrl);
this.httpProvider = new ethers.JsonRpcProvider(this.httpUrl);
this.wsContract = new ethers.Contract(
this.contractCfg.address,
this.contractCfg.abi,
this.wsProvider
);
console.log(`WebSocket Provider connected to ${this.wsUrl}`);
this.httpContract = new ethers.Contract(
this.contractCfg.address,
this.contractCfg.abi,
this.httpProvider
);
console.log(`HTTP Provider connected to ${this.httpUrl}`);
this._subscribeEvents();
console.log("HybridIndexer started");
}
_subscribeEvents() {
this.events.forEach((eventName) => {
this.wsContract.on(eventName, async (...args) => {
console.log(args);
});
});
}
};
code : server.js
const express = require("express");
const Indexer = require("./indexer/hybridIndexer");
const app = express();
// Fuji Network configuration
const wsUrl = "wss://api.avax-test.network/ext/bc/C/ws";
const httpUrl = "https://api.avax-test.network/ext/bc/C/rpc";
const abi = require("./abi/ABI.json"); // 예시 ABI 파일
const contractAddress = "0xa85761417ed8ae28378612a63282b4A576949377";
const contractCfg = {address : contractAddress, abi}
const events = ["LogMessage","TestTransfer"]
app.get("/", (req, res) => {
return res.send("Hello, World!");
});
const indexer = new Indexer({
wsUrl,
httpUrl,
contractCfg, // ethers.Contract Instance
events, // ["Transfer", ...]
});
indexer.start()
app.listen(3000, () => {
console.log("Server is running on port 3000");
});
서버를 실행하면 indexer가 실행되고, events에 대해 subscription 실행

Test.sol 의 LogMessage Event 발생 시 서버측에서 이벤트 확인

확장성을 고려해서 공통 column / eventArgs column을 구분
wsUrl : WebSocket RPC URLhttpUrl : http RPC URLcontractCfgaddress : 컨트랙트 주소abi : 컨트랙트 abieventName : 받아올 event의 이름tableName : 받아온 event data 를 저장하는 table 이름argKeys : args 이름 배열 [”eventArgName1”, “eventArgName2”]confirmations : block confirmation 수chunkSize : polling 최대 블록 수pollIntervalMs : polling 인터벌 mspingIntervalMs : ws connection status 체크를 위한 ping 인터벌 msconst { ethers } = require("ethers");
const log = require("../libraries/timestampLog");
const insertEvent = require("../database/insert");
const updateEvent = require("../database/update");
const deleteEvent = require("../database/delete");
const updateLastQueryBlockNumber = require("../database/updateLastQueryBlock");
const getLastBlockNumber = require("../database/getLastBlockNumber");
const getEvent = require("../database/getEvent");
module.exports = class HybridIndexer {
constructor({
wsUrl,
httpUrl,
contractCfg, // {address, abi}
eventName,
tableName,
argKeys, // ["eventArgName1", "eventArgName2"]
confirmations = 10,
chunkSize = 1000,
pollIntervalMs = 30000, // Default 30 seconds
pingIntervalMs = 300000, // Default 5 minutes
}) {
this.wsUrl = wsUrl;
this.httpUrl = httpUrl;
this.contractCfg = contractCfg;
this.eventName = eventName;
this.confirmations = confirmations;
this.chunkSize = chunkSize;
this.pollIntervalMs = pollIntervalMs;
this.intervalPoll = null;
this.mutexLock = false;
this.lastConfirmedBlock = 0;
this.wsProvider = null;
this.httpProvider = null;
this.wsContract = null;
this.httpContract = null;
this.tableName = tableName;
this.argKeys = argKeys;
this.intervalPing = null;
this.pingIntervalMs = pingIntervalMs;
this.running = false;
}
async start() {
this.running = true;
this.wsProvider = new ethers.WebSocketProvider(this.wsUrl);
this.httpProvider = new ethers.JsonRpcProvider(this.httpUrl);
this.wsContract = new ethers.Contract(
this.contractCfg.address,
this.contractCfg.abi,
this.wsProvider
);
log(`WebSocket Provider is set to ${this.wsUrl}`);
this.wsProvider.websocket.on("open", () => {
log(`WebSocket connected to ${this.wsUrl}`);
});
this.httpContract = new ethers.Contract(
this.contractCfg.address,
this.contractCfg.abi,
this.httpProvider
);
log(`HTTP Provider is set to ${this.httpUrl}`);
this._subscribeEvents();
log("HybridIndexer started");
// reconnect on error
this.wsProvider.websocket.on("error", () => {
log(`WebSocket Error Occurred`);
this.wsProvider?.removeAllListeners?.();
log(`Removing all listeners from WebSocket Provider`);
this.wsProvider?.destroy?.();
log(`WebSocket Provider destroyed`);
if (this.running) {
log("Attempting to reconnect...");
this.start();
}
});
// reconnect on close
this.wsProvider.websocket.on("close", () => {
log(`WebSocket Connection Closed`);
this.wsProvider?.removeAllListeners?.();
log(`Removing all listeners from WebSocket Provider`);
this.wsProvider?.destroy?.();
log(`WebSocket Provider destroyed`);
// Reconnect if still running
if (this.running) {
log("Attempting to reconnect...");
this.start();
}
});
// To keep the WebSocket connection alive, send periodic pings
this.intervalPing = setInterval(() => {
if (this.wsProvider.websocket.readyState === 1) {
this.wsProvider.websocket.ping();
log("WebSocket ping sent");
} else {
log("WebSocket is not open, skipping ping");
}
}, this.pingIntervalMs);
this._confirmEvents();
}
stop() {
this.running = false;
try {
this.wsProvider?.removeAllListeners?.();
this.wsProvider?.destroy?.();
this.wsProvider = null;
} catch (error) {
error("Error stopping WebSocket Provider:", error);
}
// Clear the ping interval if it exists
if (this.intervalPing) {
clearInterval(this.intervalPing);
this.intervalPing = null;
log("Ping interval cleared");
}
log("HybridIndexer stopped");
}
_subscribeEvents() {
this.wsContract.on(this.eventName, async (...args) => {
const event = args[args.length - 1]; // ethers v6: 마지막 인자가 Event
await insertEvent({
tableName: this.tableName,
transactionHash: event.log.transactionHash,
blockNumber: event.log.blockNumber,
logIndex: event.log.index,
transactionIndex: event.log.transactionIndex,
argKeys: this.argKeys,
argValues: event.log.args,
});
log(
`Event has been inserted into the database : ${event.log.blockNumber}-${event.log.index}-${event.log.transactionIndex}`
);
});
log(`Subscribed to event: ${this.eventName}`);
}
_confirmEvents() {
this.intervalPoll = setInterval(async () => {
if (this.mutexLock) {
console.log(`Mutex is locked`);
return;
}
this.mutexLock = true;
try {
const network = await this.httpProvider.getNetwork();
const lastBlockNumber = await getLastBlockNumber({
eventName: this.eventName,
contractAddress: this.contractCfg.address,
networkId: network.chainId.toString(),
});
console.log("last block number : ", lastBlockNumber);
const currentBlockNumber = await this.httpProvider.getBlockNumber();
console.log("Current block number : ", currentBlockNumber);
const filter = this.httpContract.filters[this.eventName]();
const fromBlockNumber = lastBlockNumber + 1;
const toBlockNumber = Math.min(
currentBlockNumber - this.confirmations,
fromBlockNumber + this.chunkSize
);
console.log(fromBlockNumber);
console.log(toBlockNumber);
// invalid block range
if (fromBlockNumber >= toBlockNumber) {
console.log(
`Block range is invalid: from:${fromBlockNumber} >= to:${toBlockNumber}`
);
throw new Error(`Invalid block range`);
}
const events = await this.httpContract.queryFilter(
filter,
fromBlockNumber,
toBlockNumber
);
for (const e of events) {
const {
transactionHash,
blockNumber,
removed,
index,
transactionIndex,
args,
} = e;
// Check if this event exist on db
const exists = await getEvent({
tableName: this.tableName,
transactionHash,
blockNumber,
index,
transactionIndex,
});
if (exists.length > 0) {
if (!removed) {
// if the event is not removed, update "confirmed" as 1
await updateEvent({
tableName: this.tableName,
transactionHash,
blockNumber,
index,
transactionIndex,
});
} else {
// if the event is removed, delete the row
await deleteEvent({
tableName: this.tableName,
transactionHash,
blockNumber,
index,
transactionIndex,
});
}
} else {
// if the event is not exist on db, insert it
await insertEvent({
tableName: this.tableName,
transactionHash,
blockNumber,
logIndex: index,
transactionIndex,
argKeys: this.argKeys,
argValues: args,
});
}
}
console.log(`${events.length} events processed`);
await updateLastQueryBlockNumber({
tableName: "event_query_info",
eventName: this.eventName,
contractAddress: this.contractCfg.address,
networkId: network.chainId.toString(),
lastQueryBlockNumber: toBlockNumber,
});
} catch (e) {
console.log(e);
} finally {
this.mutexLock = false;
}
}, this.pollIntervalMs);
}
websocketStatus() {
return this.wsProvider?.websocket?.readyState === 1
? "Connected"
: "Disconnected";
}
};

새로운 row 확인

새로운 row가 insert된 내용의 log

confirmed status 가 0->1 로 변경됨

server - confirmed event log

Websocket Ping 으로 health check

Polling 할 block range 를 기록하기 위한 테이블
