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):
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"
+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
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}")
+65 -21
View File
@@ -8,8 +8,42 @@ from src.core.abis import CONTROLLER_ABI, MARKET_ABI
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:
"""EVM ABI 日志规范化解码器 (Layer 2: Decoding)"""
"""EVM ABI 日志规范化解码器 (包含 Controller, Market 与 User Vault)"""
def __init__(self, chain_id: int):
self.chain_id = chain_id
@@ -19,8 +53,7 @@ class EVMDecoder:
def _build_topic_maps(self):
self.event_abi_map = {}
# 解析 Controller 与 Market ABI 中的 events
for abi in CONTROLLER_ABI + MARKET_ABI:
for abi in CONTROLLER_ABI + MARKET_ABI + VAULT_EVENTS_ABI:
if abi.get("type") == "event":
name = abi.get("name")
inputs = abi.get("inputs", [])
@@ -30,7 +63,6 @@ class EVMDecoder:
self.event_abi_map[topic0] = abi
def decode_log(self, log: Dict[str, Any], block_timestamp: int = 0) -> Optional[NormalizedEvent]:
"""将原始 EVM 日志解析为统一的 NormalizedEvent"""
topics = log.get("topics", [])
if not topics:
return None
@@ -46,19 +78,24 @@ class EVMDecoder:
block_number = log.get("blockNumber", 0)
log_index = log.get("logIndex", 0)
# 解码 Indexed 参数与 Non-Indexed Data
payload: Dict[str, Any] = {}
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")]
# 1. Indexed topics
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()
# 2. Non-indexed data
data = log.get("data", "0x")
if isinstance(data, str) and data.startswith("0x"):
data_bytes = bytes.fromhex(data[2:])
@@ -79,20 +116,27 @@ class EVMDecoder:
except Exception as e:
logger.error(f"Failed to decode non-indexed data for {event_name}: {e}")
# 标准化 Event Type
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",
"SetProtocolFeeRate": "GOV_FEE_RATE_UPDATED",
"SetTreasury": "GOV_TREASURY_UPDATED",
"SetCentralWallet": "GOV_CENTRAL_WALLET_UPDATED",
"SetCreatorShare": "GOV_CREATOR_SHARE_UPDATED",
"Paused": "PROTOCOL_PAUSED",
"Unpaused": "PROTOCOL_UNPAUSED",
"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"
}
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 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)
+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
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)
)