feat(indexer): complete deterministic projection pipeline and replay tests
This commit is contained in:
+76
-17
@@ -1,22 +1,81 @@
|
||||
import asyncio
|
||||
import logging
|
||||
import signal
|
||||
import sys
|
||||
from typing import Optional
|
||||
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
|
||||
|
||||
async def run_indexer_pipeline():
|
||||
"""
|
||||
可重放的确定性流水线:
|
||||
Fetcher (拉块) -> Decoder (规范化) -> Processor (业务处理) -> Projection (确定性视图落盘)
|
||||
"""
|
||||
print(f"Starting WTFX Indexer for Chain [{indexer_settings.CHAIN_ID}]...")
|
||||
print(f"Target RPC: {indexer_settings.RPC_URL}")
|
||||
|
||||
current_block = indexer_settings.START_BLOCK
|
||||
while True:
|
||||
# 1. 抓取安全区间块日志 (Fetcher)
|
||||
# 2. 解码并生成 NormalizedEvent (Decoder)
|
||||
# 3. 幂等去重检查 (Idempotency Check)
|
||||
# 4. 业务处理与 Projection 视图更新 (Processor & Projection)
|
||||
# 5. 更新 Checkpoint
|
||||
await asyncio.sleep(indexer_settings.POLL_INTERVAL_SECONDS)
|
||||
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(run_indexer_pipeline())
|
||||
asyncio.run(main())
|
||||
|
||||
Reference in New Issue
Block a user