86 lines
3.3 KiB
Python
86 lines
3.3 KiB
Python
import asyncio
|
|
import logging
|
|
import signal
|
|
import sys
|
|
from typing import Optional
|
|
import sys
|
|
import os
|
|
sys.path.insert(0, os.path.abspath(os.path.join(os.path.dirname(__file__), "..")))
|
|
|
|
from src.config import indexer_settings
|
|
from src.ingestion.fetcher import EVMFetcher
|
|
from src.decoding.evm import EVMDecoder
|
|
from src.projections.store import ProjectionStore
|
|
from src.core.checkpoint import CheckpointState
|
|
|
|
logging.basicConfig(level=logging.INFO, format="%(asctime)s [%(levelname)s] %(name)s: %(message)s")
|
|
logger = logging.getLogger("WTFX-Indexer")
|
|
|
|
class IndexerService:
|
|
"""WTFX 确定性索引服务核心执行引擎"""
|
|
|
|
def __init__(self):
|
|
self.fetcher = EVMFetcher(indexer_settings.RPC_URL, indexer_settings.CHAIN_ID)
|
|
self.decoder = EVMDecoder(indexer_settings.CHAIN_ID)
|
|
self.projection_store = ProjectionStore()
|
|
self.checkpoint = CheckpointState(
|
|
chain_id=indexer_settings.CHAIN_ID,
|
|
last_processed_block=indexer_settings.START_BLOCK
|
|
)
|
|
self.is_running = False
|
|
|
|
async def run(self):
|
|
self.is_running = True
|
|
logger.info(f"Starting WTFX Indexer on Chain {indexer_settings.CHAIN_ID} (RPC: {indexer_settings.RPC_URL})")
|
|
logger.info(f"Monitoring Controller: {indexer_settings.CONTROLLER_ADDRESS}")
|
|
|
|
while self.is_running:
|
|
try:
|
|
latest_block = await self.fetcher.get_latest_block_number()
|
|
current_block = self.checkpoint.last_processed_block
|
|
|
|
if current_block < latest_block:
|
|
to_block = min(current_block + indexer_settings.BATCH_SIZE, latest_block)
|
|
logger.info(f"Scanning blocks {current_block + 1} -> {to_block} (Latest: {latest_block})...")
|
|
|
|
# 1. Fetch raw logs
|
|
logs = await self.fetcher.fetch_logs(
|
|
from_block=current_block + 1,
|
|
to_block=to_block
|
|
)
|
|
|
|
# 2. Decode & Normalize
|
|
for log in logs:
|
|
normalized = self.decoder.decode_log(log)
|
|
if normalized:
|
|
logger.info(f"Decoded Event: {normalized.event_type} (Tx: {normalized.tx_hash[:10]}...)")
|
|
# 3. Apply to Projection
|
|
self.projection_store.apply_event(normalized)
|
|
|
|
# 4. Advance Checkpoint
|
|
self.checkpoint.last_processed_block = to_block
|
|
self.checkpoint.updated_at = int(asyncio.get_event_loop().time())
|
|
else:
|
|
await asyncio.sleep(indexer_settings.POLL_INTERVAL_SECONDS)
|
|
|
|
except Exception as e:
|
|
logger.error(f"Error in indexer polling loop: {e}", exc_info=True)
|
|
await asyncio.sleep(indexer_settings.POLL_INTERVAL_SECONDS * 2)
|
|
|
|
def stop(self):
|
|
logger.info("Stopping Indexer service...")
|
|
self.is_running = False
|
|
|
|
async def main():
|
|
service = IndexerService()
|
|
loop = asyncio.get_running_loop()
|
|
for sig in (signal.SIGINT, signal.SIGTERM):
|
|
try:
|
|
loop.add_signal_handler(sig, service.stop)
|
|
except NotImplementedError:
|
|
pass # Windows compatibility
|
|
await service.run()
|
|
|
|
if __name__ == "__main__":
|
|
asyncio.run(main())
|