Compare commits

..
3 Commits
Author SHA1 Message Date
Bot 0352802b77 chore: remove tracked __pycache__ artifacts, add gitignore 2026-08-31 02:44:10 +08:00
Bot e23225f112 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
2026-08-31 02:41:21 +08:00
Bot 8e5df08f3a feat(vault): add user vault events decoding in indexer 2026-08-31 01:42:38 +08:00
7 changed files with 423 additions and 67 deletions
+6
View File
@@ -0,0 +1,6 @@
__pycache__/
*.pyc
.venv/
.env
data/
checkpoint_*.json
+10 -5
View File
@@ -3,19 +3,24 @@ from pydantic_settings import BaseSettings
class IndexerSettings(BaseSettings): class IndexerSettings(BaseSettings):
CHAIN_ID: int = 46630 # Robinhood Testnet CHAIN_ID: int = 46630 # Robinhood Testnet
RPC_URL: str = "https://rpc.testnet.chain.robinhood.com" RPC_URL: str = "https://rpc.testnet.chain.robinhood.com"
START_BLOCK: int = 0 # 回放起点:默认为测试网控制器/market 部署后的区块,可通过环境变量覆盖
BATCH_SIZE: int = 100 START_BLOCK: int = 110000000
BATCH_SIZE: int = 500
POLL_INTERVAL_SECONDS: float = 2.0 POLL_INTERVAL_SECONDS: float = 2.0
CONFIRMATIONS: int = 1 # Testnet 确认数 CONFIRMATIONS: int = 1 # Testnet 确认数
# 数据库与 Redis # 数据库与 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" REDIS_URL: str = "redis://localhost:6379/0"
# checkpoint 持久化文件
CHECKPOINT_FILE: str = "./data/checkpoint_46630.json"
# 已部署合约地址 (Robinhood Testnet) # 已部署合约地址 (Robinhood Testnet)
CONTROLLER_ADDRESS: str = "0xc0E24E152771C588B21AEB654b30B1cBAf381c1a" CONTROLLER_ADDRESS: str = "0xc0E24E152771C588B21AEB654b30B1cBAf381c1a"
COLLATERAL_ADDRESS: str = "0xe776e957953EA69b7Eaa9d7d4098aBC076bDD5E7" COLLATERAL_ADDRESS: str = "0xe776e957953EA69b7Eaa9d7d4098aBC076bDD5E7"
CURVE_ADDRESS: str = "0xF0E189974c413506098AFB6Bc6Bf4C6715fBf7B1" CURVE_ADDRESS: str = "0xF0E189974c413506098AFB6Bc6Bf4C6715fBf7B1"
VAULT_FACTORY_ADDRESS: str = "0x89401e07296267c01017cA150D5Ec8883a78e0B2"
class Config: class Config:
env_file = ".env" env_file = ".env"
+31
View File
@@ -1,7 +1,38 @@
import json
import logging
import os
from pathlib import Path
from typing import Optional
from pydantic import BaseModel from pydantic import BaseModel
logger = logging.getLogger(__name__)
class CheckpointState(BaseModel): class CheckpointState(BaseModel):
chain_id: int chain_id: int
last_processed_block: int last_processed_block: int
last_processed_tx: str = "" last_processed_tx: str = ""
updated_at: int = 0 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}")
+65 -21
View File
@@ -8,8 +8,42 @@ from src.core.abis import CONTROLLER_ABI, MARKET_ABI
logger = logging.getLogger(__name__) logger = logging.getLogger(__name__)
# Vault Factory & User Vault Events ABI
VAULT_EVENTS_ABI = [
{
"anonymous": False,
"inputs": [
{"indexed": True, "name": "user", "type": "address"},
{"indexed": True, "name": "vault", "type": "address"},
{"indexed": False, "name": "vaultIndex", "type": "uint256"}
],
"name": "VaultCreated",
"type": "event"
},
{
"anonymous": False,
"inputs": [
{"indexed": True, "name": "token", "type": "address"},
{"indexed": True, "name": "from", "type": "address"},
{"indexed": False, "name": "amount", "type": "uint256"}
],
"name": "Deposited",
"type": "event"
},
{
"anonymous": False,
"inputs": [
{"indexed": True, "name": "token", "type": "address"},
{"indexed": True, "name": "to", "type": "address"},
{"indexed": False, "name": "amount", "type": "uint256"}
],
"name": "Withdrawn",
"type": "event"
}
]
class EVMDecoder: class EVMDecoder:
"""EVM ABI 日志规范化解码器 (Layer 2: Decoding)""" """EVM ABI 日志规范化解码器 (包含 Controller, Market 与 User Vault)"""
def __init__(self, chain_id: int): def __init__(self, chain_id: int):
self.chain_id = chain_id self.chain_id = chain_id
@@ -19,8 +53,7 @@ class EVMDecoder:
def _build_topic_maps(self): def _build_topic_maps(self):
self.event_abi_map = {} self.event_abi_map = {}
# 解析 Controller 与 Market ABI 中的 events for abi in CONTROLLER_ABI + MARKET_ABI + VAULT_EVENTS_ABI:
for abi in CONTROLLER_ABI + MARKET_ABI:
if abi.get("type") == "event": if abi.get("type") == "event":
name = abi.get("name") name = abi.get("name")
inputs = abi.get("inputs", []) inputs = abi.get("inputs", [])
@@ -30,7 +63,6 @@ class EVMDecoder:
self.event_abi_map[topic0] = abi self.event_abi_map[topic0] = abi
def decode_log(self, log: Dict[str, Any], block_timestamp: int = 0) -> Optional[NormalizedEvent]: def decode_log(self, log: Dict[str, Any], block_timestamp: int = 0) -> Optional[NormalizedEvent]:
"""将原始 EVM 日志解析为统一的 NormalizedEvent"""
topics = log.get("topics", []) topics = log.get("topics", [])
if not topics: if not topics:
return None return None
@@ -46,19 +78,24 @@ class EVMDecoder:
block_number = log.get("blockNumber", 0) block_number = log.get("blockNumber", 0)
log_index = log.get("logIndex", 0) log_index = log.get("logIndex", 0)
# 解码 Indexed 参数与 Non-Indexed Data
payload: Dict[str, Any] = {} payload: Dict[str, Any] = {}
indexed_inputs = [i for i in abi.get("inputs", []) if i.get("indexed")] indexed_inputs = [i for i in abi.get("inputs", []) if i.get("indexed")]
non_indexed_inputs = [i for i in abi.get("inputs", []) if not i.get("indexed")] non_indexed_inputs = [i for i in abi.get("inputs", []) if not i.get("indexed")]
# 1. Indexed topics
for idx, inp in enumerate(indexed_inputs): for idx, inp in enumerate(indexed_inputs):
if idx + 1 < len(topics): if idx + 1 < len(topics):
raw_topic = topics[idx + 1] raw_topic = topics[idx + 1]
t_hex = raw_topic.hex() if isinstance(raw_topic, bytes) else raw_topic if isinstance(raw_topic, bytes):
payload[inp["name"]] = t_hex 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()
# 2. Non-indexed data
data = log.get("data", "0x") data = log.get("data", "0x")
if isinstance(data, str) and data.startswith("0x"): if isinstance(data, str) and data.startswith("0x"):
data_bytes = bytes.fromhex(data[2:]) data_bytes = bytes.fromhex(data[2:])
@@ -79,20 +116,27 @@ class EVMDecoder:
except Exception as e: except Exception as e:
logger.error(f"Failed to decode non-indexed data for {event_name}: {e}") logger.error(f"Failed to decode non-indexed data for {event_name}: {e}")
# 标准化 Event Type
type_mapping = { type_mapping = {
"DeployMarket": "MARKET_DEPLOYED", # WTFControllerV2 / MarketFactory
"SetOutcome": "OUTCOME_RESOLVED", "CreateNewMarket": "MARKET_DEPLOYED",
"CreateNewQuestionV2": "QUESTION_CREATED",
"Resolve": "OUTCOME_RESOLVED",
"Unresolve": "OUTCOME_UNRESOLVED",
"Finalise": "MARKET_FINALISED", "Finalise": "MARKET_FINALISED",
"Mint": "ORDER_MINT", "OverrideFinalise": "OUTCOME_OVERRIDDEN",
"Redeem": "ORDER_REDEEM", "ManuallyFinalise": "MARKET_FINALISED",
"Claim": "POSITION_CLAIMED", # WTFMarketV2 (V2 真实成交事件)
"SetProtocolFeeRate": "GOV_FEE_RATE_UPDATED", "MintSwapV2": "ORDER_MINT",
"SetTreasury": "GOV_TREASURY_UPDATED", "RedeemSwapV2": "ORDER_REDEEM",
"SetCentralWallet": "GOV_CENTRAL_WALLET_UPDATED", "MintSwap": "ORDER_MINT",
"SetCreatorShare": "GOV_CREATOR_SHARE_UPDATED", "RedeemSwap": "ORDER_REDEEM",
"Paused": "PROTOCOL_PAUSED", "ClaimPayout": "POSITION_CLAIMED",
"Unpaused": "PROTOCOL_UNPAUSED", "GraduateMarket": "MARKET_GRADUATED",
"MarketRefunded": "MARKET_REFUNDED",
# Vault
"VaultCreated": "VAULT_CREATED",
"Deposited": "VAULT_DEPOSITED",
"Withdrawn": "VAULT_WITHDRAWN"
} }
normalized_type = type_mapping.get(event_name, f"EVM_{event_name.upper()}") normalized_type = type_mapping.get(event_name, f"EVM_{event_name.upper()}")
+49 -12
View File
@@ -2,7 +2,7 @@ import asyncio
import logging import logging
import signal import signal
import sys import sys
from typing import Optional from typing import Optional, Set
import sys import sys
import os import os
sys.path.insert(0, os.path.abspath(os.path.join(os.path.dirname(__file__), ".."))) 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.ingestion.fetcher import EVMFetcher
from src.decoding.evm import EVMDecoder from src.decoding.evm import EVMDecoder
from src.projections.store import ProjectionStore from src.projections.store import ProjectionStore
from src.projections.repository import KlineRepository
from src.core.checkpoint import CheckpointState from src.core.checkpoint import CheckpointState
logging.basicConfig(level=logging.INFO, format="%(asctime)s [%(levelname)s] %(name)s: %(message)s") 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.fetcher = EVMFetcher(indexer_settings.RPC_URL, indexer_settings.CHAIN_ID)
self.decoder = EVMDecoder(indexer_settings.CHAIN_ID) self.decoder = EVMDecoder(indexer_settings.CHAIN_ID)
self.projection_store = ProjectionStore() self.projection_store = ProjectionStore()
self.checkpoint = CheckpointState( self.checkpoint = CheckpointState.load(
chain_id=indexer_settings.CHAIN_ID, indexer_settings.CHECKPOINT_FILE,
last_processed_block=indexer_settings.START_BLOCK 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 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): async def run(self):
self.is_running = True 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"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: while self.is_running:
try: try:
@@ -43,23 +59,44 @@ class IndexerService:
to_block = min(current_block + indexer_settings.BATCH_SIZE, 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})...") logger.info(f"Scanning blocks {current_block + 1} -> {to_block} (Latest: {latest_block})...")
# 1. Fetch raw logs
logs = await self.fetcher.fetch_logs( logs = await self.fetcher.fetch_logs(
from_block=current_block + 1, 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: 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: if normalized:
logger.info(f"Decoded Event: {normalized.event_type} (Tx: {normalized.tx_hash[:10]}...)") decoded_count += 1
# 3. Apply to Projection 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) 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.last_processed_block = to_block
self.checkpoint.updated_at = int(asyncio.get_event_loop().time()) 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: else:
await asyncio.sleep(indexer_settings.POLL_INTERVAL_SECONDS) await asyncio.sleep(indexer_settings.POLL_INTERVAL_SECONDS)
+127
View File
@@ -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
]
+135 -29
View File
@@ -1,13 +1,34 @@
import logging import logging
from typing import Dict, Any, Optional import math
from typing import Dict, Any, List, Optional
from src.core.events import NormalizedEvent from src.core.events import NormalizedEvent
logger = logging.getLogger(__name__) 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: class ProjectionStore:
""" """
确定性业务投影状态存储 (Layer 4: Projections) 确定性业务投影与秒级 OHLCV K 线聚合引擎】
在第一阶段使用内存/本地轻量结构,并支持向 PostgreSQL 写入投影,支持随时 Reset & Replay - 记录逐笔成交 (Trades)
- 动态维护市场各 Token 瞬时价格与流动性
- 按照 5s 粒度聚合实时 OHLCV K 线(适配 TradingView 图表)
- 内存投影 + 脏数据队列,由 main 定期 flush 到 PostgreSQL
""" """
def __init__(self): def __init__(self):
@@ -16,25 +37,87 @@ class ProjectionStore:
self.positions: Dict[str, Dict[str, Any]] = {} self.positions: Dict[str, Dict[str, Any]] = {}
self.processed_events: set = set() 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): def reset(self):
"""全量清空投影视图(支持 Replay"""
self.markets.clear() self.markets.clear()
self.trades.clear() self.trades.clear()
self.positions.clear() self.positions.clear()
self.processed_events.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): def apply_event(self, event: NormalizedEvent):
"""确定性事件投影应用"""
if event.event_id in self.processed_events: if event.event_id in self.processed_events:
return # 幂等拦截 return
self.processed_events.add(event.event_id) self.processed_events.add(event.event_id)
etype = event.event_type etype = event.event_type
payload = event.payload payload = event.payload
if etype == "MARKET_DEPLOYED": 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] = { self.markets[market_addr] = {
"market_address": market_addr, "market_address": market_addr,
"question_id": payload.get("questionId"), "question_id": payload.get("questionId"),
@@ -50,30 +133,53 @@ class ProjectionStore:
"timestamp": event.timestamp "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"): 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({ self.trades.append({
"tx_hash": event.tx_hash, "tx_hash": event.tx_hash,
"block_number": event.block_number, "block_number": event.block_number,
"market": event.contract_address, "market": market_addr,
"type": "MINT" if etype == "ORDER_MINT" else "REDEEM", "type": "MINT" if etype == "ORDER_MINT" else "REDEEM",
"user": payload.get("user") or payload.get("buyer") or payload.get("seller"), "user": payload.get("receiver") or payload.get("user"),
"outcome_index": payload.get("outcomeIndex") or payload.get("index"), "outcome_index": outcome_idx,
"collateral_amount": payload.get("collateralAmount") or payload.get("amountIn"), "collateral_amount": collateral_raw,
"tokens_amount": payload.get("tokensAmount") or payload.get("amountOut"), "tokens_amount": tokens_raw,
"fee": payload.get("fee", 0), "price": trade_price,
"timestamp": event.timestamp "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)
)