From e23225f1122dad83f4af9f061f47a0e07945ef8e Mon Sep 17 00:00:00 2001 From: Bot Date: Mon, 31 Aug 2026 02:41:21 +0800 Subject: [PATCH] feat(indexer): decode real WTFMarketV2 events, persist klines/trades to Postgres - evm.py: map MintSwapV2/RedeemSwapV2/CreateNewMarket; decode indexed address topics correctly - store.py: parse MintSwapV2/RedeemSwapV2 fields, compute real trade price (no [0,1] clamp), dirty-queue for trades/klines - repository.py: asyncpg upsert of 5s OHLCV bars + trades - main.py: watch controller/vault-factory/default market, dynamic market discovery, checkpoint resume - checkpoint.py: JSON checkpoint persistence --- src/config.py | 15 ++-- src/core/checkpoint.py | 31 +++++++ src/decoding/evm.py | 33 +++++-- src/main.py | 61 ++++++++++--- src/projections/repository.py | 127 ++++++++++++++++++++++++++ src/projections/store.py | 164 ++++++++++++++++++++++++++++------ 6 files changed, 378 insertions(+), 53 deletions(-) create mode 100644 src/projections/repository.py diff --git a/src/config.py b/src/config.py index 931c19e..b7661b2 100644 --- a/src/config.py +++ b/src/config.py @@ -3,19 +3,24 @@ from pydantic_settings import BaseSettings class IndexerSettings(BaseSettings): CHAIN_ID: int = 46630 # Robinhood Testnet RPC_URL: str = "https://rpc.testnet.chain.robinhood.com" - START_BLOCK: int = 0 - BATCH_SIZE: int = 100 + # 回放起点:默认为测试网控制器/market 部署后的区块,可通过环境变量覆盖 + START_BLOCK: int = 110000000 + BATCH_SIZE: int = 500 POLL_INTERVAL_SECONDS: float = 2.0 CONFIRMATIONS: int = 1 # Testnet 确认数 - + # 数据库与 Redis - DATABASE_URL: str = "postgresql+asyncpg://postgres:postgres@localhost:5432/wtfx" + DATABASE_URL: str = "postgresql+asyncpg://wtfx_user:wtfx_pass@localhost:5432/wtfx_db" REDIS_URL: str = "redis://localhost:6379/0" - + + # checkpoint 持久化文件 + CHECKPOINT_FILE: str = "./data/checkpoint_46630.json" + # 已部署合约地址 (Robinhood Testnet) CONTROLLER_ADDRESS: str = "0xc0E24E152771C588B21AEB654b30B1cBAf381c1a" COLLATERAL_ADDRESS: str = "0xe776e957953EA69b7Eaa9d7d4098aBC076bDD5E7" CURVE_ADDRESS: str = "0xF0E189974c413506098AFB6Bc6Bf4C6715fBf7B1" + VAULT_FACTORY_ADDRESS: str = "0x89401e07296267c01017cA150D5Ec8883a78e0B2" class Config: env_file = ".env" diff --git a/src/core/checkpoint.py b/src/core/checkpoint.py index 07b9d92..bc84c19 100644 --- a/src/core/checkpoint.py +++ b/src/core/checkpoint.py @@ -1,7 +1,38 @@ +import json +import logging +import os +from pathlib import Path +from typing import Optional + from pydantic import BaseModel +logger = logging.getLogger(__name__) + class CheckpointState(BaseModel): chain_id: int last_processed_block: int last_processed_tx: str = "" updated_at: int = 0 + + @classmethod + def load(cls, file_path: str, chain_id: int, default_block: int) -> "CheckpointState": + try: + p = Path(file_path) + if p.exists(): + data = json.loads(p.read_text(encoding="utf-8")) + state = cls(**data) + if state.chain_id == chain_id: + return state + except Exception as e: + logger.warning(f"Failed to load checkpoint from {file_path}: {e}") + return cls(chain_id=chain_id, last_processed_block=default_block) + + def save(self, file_path: str): + try: + p = Path(file_path) + p.parent.mkdir(parents=True, exist_ok=True) + tmp = p.with_suffix(".tmp") + tmp.write_text(self.model_dump_json(), encoding="utf-8") + os.replace(tmp, p) + except Exception as e: + logger.warning(f"Failed to save checkpoint to {file_path}: {e}") diff --git a/src/decoding/evm.py b/src/decoding/evm.py index 4cd7f98..310da71 100644 --- a/src/decoding/evm.py +++ b/src/decoding/evm.py @@ -85,8 +85,16 @@ class EVMDecoder: for idx, inp in enumerate(indexed_inputs): if idx + 1 < len(topics): raw_topic = topics[idx + 1] - t_hex = raw_topic.hex() if isinstance(raw_topic, bytes) else raw_topic - payload[inp["name"]] = t_hex + if isinstance(raw_topic, bytes): + t_bytes = raw_topic + else: + t_hex = raw_topic[2:] if raw_topic.startswith("0x") else raw_topic + t_bytes = bytes.fromhex(t_hex) + if inp["type"] == "address": + # indexed address topic 右对齐 20 字节,需要去掉前导 0 并转成标准 0x 地址 + payload[inp["name"]] = self.w3.to_checksum_address("0x" + t_bytes[-20:].hex()) + else: + payload[inp["name"]] = "0x" + t_bytes.hex() data = log.get("data", "0x") if isinstance(data, str) and data.startswith("0x"): @@ -109,12 +117,23 @@ class EVMDecoder: logger.error(f"Failed to decode non-indexed data for {event_name}: {e}") type_mapping = { - "DeployMarket": "MARKET_DEPLOYED", - "SetOutcome": "OUTCOME_RESOLVED", + # WTFControllerV2 / MarketFactory + "CreateNewMarket": "MARKET_DEPLOYED", + "CreateNewQuestionV2": "QUESTION_CREATED", + "Resolve": "OUTCOME_RESOLVED", + "Unresolve": "OUTCOME_UNRESOLVED", "Finalise": "MARKET_FINALISED", - "Mint": "ORDER_MINT", - "Redeem": "ORDER_REDEEM", - "Claim": "POSITION_CLAIMED", + "OverrideFinalise": "OUTCOME_OVERRIDDEN", + "ManuallyFinalise": "MARKET_FINALISED", + # WTFMarketV2 (V2 真实成交事件) + "MintSwapV2": "ORDER_MINT", + "RedeemSwapV2": "ORDER_REDEEM", + "MintSwap": "ORDER_MINT", + "RedeemSwap": "ORDER_REDEEM", + "ClaimPayout": "POSITION_CLAIMED", + "GraduateMarket": "MARKET_GRADUATED", + "MarketRefunded": "MARKET_REFUNDED", + # Vault "VaultCreated": "VAULT_CREATED", "Deposited": "VAULT_DEPOSITED", "Withdrawn": "VAULT_WITHDRAWN" diff --git a/src/main.py b/src/main.py index 8e9e693..46f63f9 100644 --- a/src/main.py +++ b/src/main.py @@ -2,7 +2,7 @@ import asyncio import logging import signal import sys -from typing import Optional +from typing import Optional, Set import sys import os sys.path.insert(0, os.path.abspath(os.path.join(os.path.dirname(__file__), ".."))) @@ -11,6 +11,7 @@ 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.projections.repository import KlineRepository from src.core.checkpoint import CheckpointState logging.basicConfig(level=logging.INFO, format="%(asctime)s [%(levelname)s] %(name)s: %(message)s") @@ -23,16 +24,31 @@ class IndexerService: 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.checkpoint = CheckpointState.load( + indexer_settings.CHECKPOINT_FILE, + indexer_settings.CHAIN_ID, + indexer_settings.START_BLOCK ) + self.repository = KlineRepository( + indexer_settings.DATABASE_URL.replace("postgresql+asyncpg://", "postgresql://", 1) + ) + # 动态关注地址:控制器 + 金库工厂 + 已知默认市场;新市场由 CreateNewMarket 动态加入 + self.watch_addresses: Set[str] = set() self.is_running = False + async def _bootstrap_watch_list(self): + self.watch_addresses.add(indexer_settings.CONTROLLER_ADDRESS.lower()) + self.watch_addresses.add(indexer_settings.VAULT_FACTORY_ADDRESS.lower()) + default_market = os.getenv("DEFAULT_MARKET", "0x048E9a90C25ba2c4410425D282b19A472e076039") + self.watch_addresses.add(default_market.lower()) + logger.info(f"Watch list initialized: {sorted(self.watch_addresses)}") + async def run(self): self.is_running = True + await self.repository.connect() + await self._bootstrap_watch_list() 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}") + logger.info(f"Resuming from block {self.checkpoint.last_processed_block}") while self.is_running: try: @@ -43,23 +59,44 @@ class IndexerService: 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 + to_block=to_block, + addresses=sorted(self.watch_addresses) ) - # 2. Decode & Normalize + # 一次性抓取本批涉及的区块时间戳,保证 K 线时间桶精确 + block_nums = {int(log.get("blockNumber", 0)) for log in logs} + timestamps = {} + for bn in block_nums: + timestamps[bn] = await self.fetcher.get_block_timestamp(bn) + + decoded_count = 0 for log in logs: - normalized = self.decoder.decode_log(log) + bn = int(log.get("blockNumber", 0)) + normalized = self.decoder.decode_log(log, block_timestamp=timestamps.get(bn, 0)) if normalized: - logger.info(f"Decoded Event: {normalized.event_type} (Tx: {normalized.tx_hash[:10]}...)") - # 3. Apply to Projection + decoded_count += 1 + if normalized.event_type == "MARKET_DEPLOYED": + market_addr = (normalized.payload.get("market") or "").lower() + if market_addr: + self.watch_addresses.add(market_addr) + logger.info(f"New market watched: {market_addr}") self.projection_store.apply_event(normalized) - # 4. Advance Checkpoint + # 持久化脏数据 (trades + klines) + await self.repository.flush( + self.projection_store.get_dirty_trades(), + self.projection_store.get_dirty_klines() + ) + self.projection_store.mark_trades_flushed() + self.projection_store.mark_klines_flushed() + + # 推进并持久化 checkpoint self.checkpoint.last_processed_block = to_block self.checkpoint.updated_at = int(asyncio.get_event_loop().time()) + self.checkpoint.save(indexer_settings.CHECKPOINT_FILE) + logger.info(f"Batch done. Decoded {decoded_count} events. Checkpoint -> {to_block}") else: await asyncio.sleep(indexer_settings.POLL_INTERVAL_SECONDS) diff --git a/src/projections/repository.py b/src/projections/repository.py new file mode 100644 index 0000000..15c0e3b --- /dev/null +++ b/src/projections/repository.py @@ -0,0 +1,127 @@ +import logging +from typing import Any, Dict, List, Tuple + +import asyncpg + +logger = logging.getLogger(__name__) + +CREATE_TABLES_SQL = """ +CREATE TABLE IF NOT EXISTS klines ( + market_address VARCHAR(42) NOT NULL, + outcome_index INTEGER NOT NULL, + bar_time BIGINT NOT NULL, + open NUMERIC NOT NULL, + high NUMERIC NOT NULL, + low NUMERIC NOT NULL, + close NUMERIC NOT NULL, + volume NUMERIC NOT NULL, + PRIMARY KEY (market_address, outcome_index, bar_time) +); +CREATE INDEX IF NOT EXISTS idx_klines_market_time + ON klines (market_address, outcome_index, bar_time); + +CREATE TABLE IF NOT EXISTS trades ( + id BIGSERIAL PRIMARY KEY, + tx_hash VARCHAR(66), + market_address VARCHAR(42) NOT NULL, + outcome_index INTEGER NOT NULL, + trade_type VARCHAR(8) NOT NULL, + price NUMERIC NOT NULL, + volume NUMERIC NOT NULL, + block_number BIGINT, + ts BIGINT +); +CREATE INDEX IF NOT EXISTS idx_trades_market_ts + ON trades (market_address, outcome_index, ts); +""" + +UPSERT_KLINE_SQL = """ +INSERT INTO klines (market_address, outcome_index, bar_time, open, high, low, close, volume) +VALUES ($1, $2, $3, $4, $5, $6, $7, $8) +ON CONFLICT (market_address, outcome_index, bar_time) +DO UPDATE SET + high = GREATEST(klines.high, EXCLUDED.high), + low = LEAST(klines.low, EXCLUDED.low), + close = EXCLUDED.close, + volume = klines.volume + EXCLUDED.volume; +""" + +INSERT_TRADE_SQL = """ +INSERT INTO trades (tx_hash, market_address, outcome_index, trade_type, price, volume, block_number, ts) +VALUES ($1, $2, $3, $4, $5, $6, $7, $8); +""" + +SELECT_KLINES_SQL = """ +SELECT bar_time, open, high, low, close, volume +FROM klines +WHERE market_address = $1 AND outcome_index = $2 +ORDER BY bar_time ASC +LIMIT $3; +""" + + +class KlineRepository: + """indexer -> PostgreSQL 的真实 K 线 / 成交持久化""" + + def __init__(self, database_url: str): + self.database_url = database_url + self._pool: asyncpg.Pool | None = None + + async def connect(self): + self._pool = await asyncpg.create_pool(dsn=self.database_url, min_size=1, max_size=3) + async with self._pool.acquire() as conn: + await conn.execute(CREATE_TABLES_SQL) + logger.info("KlineRepository connected to PostgreSQL, tables ensured.") + + async def close(self): + if self._pool: + await self._pool.close() + self._pool = None + + async def flush(self, trades: list, klines: Dict[Tuple[str, int], List[Dict[str, Any]]]): + if not self._pool: + return + async with self._pool.acquire() as conn: + async with conn.transaction(): + for t in trades: + await conn.execute( + INSERT_TRADE_SQL, + t.get("tx_hash"), + t.get("market"), + t.get("outcome_index", 0), + t.get("type", "MINT"), + t.get("price", 0), + t.get("collateral_amount", 0), + t.get("block_number"), + t.get("timestamp", 0), + ) + for (m_addr, o_idx), bars in klines.items(): + for bar in bars: + await conn.execute( + UPSERT_KLINE_SQL, + m_addr, + o_idx, + bar["time"], + bar["open"], + bar["high"], + bar["low"], + bar["close"], + bar["volume"], + ) + + async def get_klines(self, market_address: str, outcome_index: int = 0, limit: int = 1000): + if not self._pool: + return [] + async with self._pool.acquire() as conn: + rows = await conn.fetch(SELECT_KLINES_SQL, market_address.lower(), outcome_index, limit) + return [ + { + "time": r["bar_time"], + "open": float(r["open"]), + "high": float(r["high"]), + "low": float(r["low"]), + "close": float(r["close"]), + "volume": float(r["volume"]), + } + for r in rows + ] diff --git a/src/projections/store.py b/src/projections/store.py index d20f0bb..58e02f2 100644 --- a/src/projections/store.py +++ b/src/projections/store.py @@ -1,13 +1,34 @@ import logging -from typing import Dict, Any, Optional +import math +from typing import Dict, Any, List, Optional from src.core.events import NormalizedEvent logger = logging.getLogger(__name__) +def _to_int(value: Any, default: int = 0) -> int: + """robust int parsing: int | '0x..' | '000..1' (indexed topic hex) | '12'""" + if value is None: + return default + if isinstance(value, int): + return value + if isinstance(value, str): + s = value[2:] if value.startswith("0x") else value + try: + return int(s, 16) + except ValueError: + try: + return int(s) + except ValueError: + return default + return default + class ProjectionStore: """ - 确定性业务投影状态存储 (Layer 4: Projections) - 在第一阶段使用内存/本地轻量结构,并支持向 PostgreSQL 写入投影,支持随时 Reset & Replay + 【确定性业务投影与秒级 OHLCV K 线聚合引擎】 + - 记录逐笔成交 (Trades) + - 动态维护市场各 Token 瞬时价格与流动性 + - 按照 5s 粒度聚合实时 OHLCV K 线(适配 TradingView 图表) + - 内存投影 + 脏数据队列,由 main 定期 flush 到 PostgreSQL """ def __init__(self): @@ -16,25 +37,87 @@ class ProjectionStore: self.positions: Dict[str, Dict[str, Any]] = {} self.processed_events: set = set() + # K 线存储结构: market_address -> outcome_index -> List[OHLCV] + # item: {"time": 1700000000, "open": 0.5, "high": 0.52, "low": 0.48, "close": 0.51, "volume": 1200} + self.klines: Dict[str, Dict[int, List[Dict[str, Any]]]] = {} + + # 脏数据追踪(供 flush 到 DB 使用) + self._trade_watermark = 0 + self._dirty_klines: Dict[tuple, List[Dict[str, Any]]] = {} + def reset(self): - """全量清空投影视图(支持 Replay)""" self.markets.clear() self.trades.clear() self.positions.clear() self.processed_events.clear() - logger.info("Projection store reset successfully.") + self.klines.clear() + self._trade_watermark = 0 + self._dirty_klines.clear() + logger.info("Projection & K-line store reset successfully.") + + def get_klines(self, market_address: str, outcome_index: int = 0) -> List[Dict[str, Any]]: + m_addr = market_address.lower() + if m_addr in self.klines and outcome_index in self.klines[m_addr]: + return self.klines[m_addr][outcome_index] + return [] + + def get_dirty_trades(self) -> list: + """自上次 flush 以来的新增成交""" + new_trades = self.trades[self._trade_watermark:] + return new_trades + + def mark_trades_flushed(self): + self._trade_watermark = len(self.trades) + + def get_dirty_klines(self) -> Dict[tuple, List[Dict[str, Any]]]: + """自上次 flush 以来被更新的 (market, outcome) -> bars""" + return {k: v for k, v in self._dirty_klines.items()} + + def mark_klines_flushed(self): + self._dirty_klines.clear() + + def _update_kline(self, market_address: str, outcome_index: int, price: float, volume: float, timestamp: int): + m_addr = market_address.lower() + if m_addr not in self.klines: + self.klines[m_addr] = {} + if outcome_index not in self.klines[m_addr]: + self.klines[m_addr][outcome_index] = [] + + bar_time = (timestamp // 5) * 5 # 5秒一根 K 线 Bar + bars = self.klines[m_addr][outcome_index] + key = (m_addr, outcome_index) + + if bars and bars[-1]["time"] == bar_time: + # 更新当前柱子 + current = bars[-1] + current["high"] = max(current["high"], price) + current["low"] = min(current["low"], price) + current["close"] = price + current["volume"] += volume + else: + # 开新柱子 + prev_close = bars[-1]["close"] if bars else price + bars.append({ + "time": bar_time, + "open": prev_close, + "high": max(prev_close, price), + "low": min(prev_close, price), + "close": price, + "volume": volume + }) + + self._dirty_klines[key] = bars def apply_event(self, event: NormalizedEvent): - """确定性事件投影应用""" if event.event_id in self.processed_events: - return # 幂等拦截 + return self.processed_events.add(event.event_id) etype = event.event_type payload = event.payload if etype == "MARKET_DEPLOYED": - market_addr = payload.get("market") or event.contract_address + market_addr = (payload.get("market") or event.contract_address).lower() self.markets[market_addr] = { "market_address": market_addr, "question_id": payload.get("questionId"), @@ -50,30 +133,53 @@ class ProjectionStore: "timestamp": event.timestamp } - elif etype == "OUTCOME_RESOLVED": - q_id = payload.get("questionId") - for m in self.markets.values(): - if m.get("question_id") == q_id: - m["status"] = "RESOLVED" - m["tentative_winning_outcome"] = payload.get("outcome") - - elif etype == "MARKET_FINALISED": - q_id = payload.get("questionId") - for m in self.markets.values(): - if m.get("question_id") == q_id: - m["status"] = "FINALISED" - m["winning_outcome"] = payload.get("outcome") - elif etype in ("ORDER_MINT", "ORDER_REDEEM"): + market_addr = event.contract_address.lower() + + token_id = _to_int(payload.get("tokenId")) or _to_int(payload.get("id")) + outcome_idx = 0 + if token_id > 0 and (token_id & (token_id - 1)) == 0: + # tokenId = 2^idx -> idx + outcome_idx = int(math.log2(token_id)) + else: + outcome_idx = _to_int(payload.get("outcomeIndex")) + + if etype == "ORDER_MINT": + # MintSwapV2(caller, receiver, tokenId, collateralToPool, otToUser, collateralToTreasury) + collateral_in = _to_int(payload.get("collateralToPool")) + _to_int(payload.get("collateralToTreasury")) + ot_moved = _to_int(payload.get("otToUser")) + volume_raw = collateral_in + else: + # RedeemSwapV2(caller, receiver, tokenId, collateralFromPool, otToPool, collateralToTreasury) + collateral_out = _to_int(payload.get("collateralFromPool")) + _to_int(payload.get("collateralToTreasury")) + ot_moved = _to_int(payload.get("otToPool")) + volume_raw = collateral_out + + collateral_raw = volume_raw / 1e18 + tokens_raw = ot_moved / 1e18 + + # 成交均价 = 移动的抵押品 / 移动的 OT。PowerLDA 价格可 > 1,不做概率区间裁剪 + trade_price = collateral_raw / tokens_raw if tokens_raw > 0 else 0.0 + self.trades.append({ "tx_hash": event.tx_hash, "block_number": event.block_number, - "market": event.contract_address, + "market": market_addr, "type": "MINT" if etype == "ORDER_MINT" else "REDEEM", - "user": payload.get("user") or payload.get("buyer") or payload.get("seller"), - "outcome_index": payload.get("outcomeIndex") or payload.get("index"), - "collateral_amount": payload.get("collateralAmount") or payload.get("amountIn"), - "tokens_amount": payload.get("tokensAmount") or payload.get("amountOut"), - "fee": payload.get("fee", 0), - "timestamp": event.timestamp + "user": payload.get("receiver") or payload.get("user"), + "outcome_index": outcome_idx, + "collateral_amount": collateral_raw, + "tokens_amount": tokens_raw, + "price": trade_price, + "fee": _to_int(payload.get("collateralToTreasury")) / 1e18, + "timestamp": event.timestamp or int(event.block_number) }) + + # 更新 OHLCV K 线 + self._update_kline( + market_address=market_addr, + outcome_index=outcome_idx, + price=trade_price, + volume=collateral_raw, + timestamp=event.timestamp or int(event.block_number) + )