Websocket을 활용한 event indexing

Natreeum's Blog·2026년 1월 9일

Blockchain

목록 보기
4/6

github

https://github.com/natreeum/HybridIndexer

개요

기존 인덱싱 방식은 http 방식(Polling)으로 진행. 이로인해 설정된 주기로만 요청하므로 최대 ‘설정된 주기까지의 delay’가 생김. 이로 인해 낙관적 업데이트와 같은 추가 작업을 필요로 합니다.

구상

  1. ws provider를 통한 event subscription으로 event 발생 시 트랜잭션(event) 수집
  2. 웹소켓으로 수집된 트랜잭션은 DB에 저장
    a. 최초 status 는 unconfirmed상태로 저장
  3. batch indexing을 통해 unconfirmed 상태의 트랜잭션에 대해 confirm
    a. removed : true 인 경우(reorg) 삭제
    b. 누락된 데이터 업데이트
  4. ws connection 이 유실되었을 경우 재연결 시도

0. Contract

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

1. Websocket으로 Emit되는 event 확인

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 발생 시 서버측에서 이벤트 확인

2. DB Table

확장성을 고려해서 공통 column / eventArgs column을 구분

공통 column

  • transaction_hash
  • block_number
  • log_index
  • transaction_index
  • confirmed

Event Args Column

  • ...

3. Indexer Constructor

  • wsUrl : WebSocket RPC URL
  • httpUrl : http RPC URL
  • contractCfg
    - address : 컨트랙트 주소
    • abi : 컨트랙트 abi
  • eventName : 받아올 event의 이름
  • tableName : 받아온 event data 를 저장하는 table 이름
  • argKeys : args 이름 배열 [”eventArgName1”, “eventArgName2”]
  • confirmations : block confirmation 수
  • chunkSize : polling 최대 블록 수
  • pollIntervalMs : polling 인터벌 ms
  • pingIntervalMs : ws connection status 체크를 위한 ping 인터벌 ms

구현

Hybrid Indexer

const { 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";
  }
};

서버 실행

Event 발생 시 즉시 DB에 row 생성

새로운 row 확인

새로운 row가 insert된 내용의 log

Polling 으로 Tx Confirm

confirmed status 가 0->1 로 변경됨

server - confirmed event log

Websocket Ping 으로 health check

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

profile
BlockChain DEV

0개의 댓글