From 3674fa008f3fe9b486a795d89d14d91c38401efa Mon Sep 17 00:00:00 2001 From: Bot Date: Mon, 31 Aug 2026 01:02:02 +0800 Subject: [PATCH] feat(indexer): complete deterministic projection pipeline and replay tests --- src/__pycache__/__init__.cpython-310.pyc | Bin 0 -> 148 bytes src/config.py | 10 +- src/core/__pycache__/__init__.cpython-310.pyc | Bin 0 -> 153 bytes src/core/__pycache__/abis.cpython-310.pyc | Bin 0 -> 1049 bytes src/core/__pycache__/events.cpython-310.pyc | Bin 0 -> 1082 bytes src/core/abis.py | 22 ++++ .../__pycache__/__init__.cpython-310.pyc | Bin 0 -> 157 bytes src/decoding/__pycache__/evm.cpython-310.pyc | Bin 0 -> 3777 bytes src/decoding/evm.py | 109 ++++++++++++++++++ src/ingestion/fetcher.py | 50 ++++++++ src/main.py | 93 ++++++++++++--- .../__pycache__/__init__.cpython-310.pyc | Bin 0 -> 160 bytes .../__pycache__/store.cpython-310.pyc | Bin 0 -> 2657 bytes src/projections/store.py | 79 +++++++++++++ ...exer_pipeline.cpython-310-pytest-8.4.1.pyc | Bin 0 -> 3774 bytes tests/test_indexer_pipeline.py | 76 ++++++++++++ 16 files changed, 419 insertions(+), 20 deletions(-) create mode 100644 src/__pycache__/__init__.cpython-310.pyc create mode 100644 src/core/__pycache__/__init__.cpython-310.pyc create mode 100644 src/core/__pycache__/abis.cpython-310.pyc create mode 100644 src/core/__pycache__/events.cpython-310.pyc create mode 100644 src/core/abis.py create mode 100644 src/decoding/__pycache__/__init__.cpython-310.pyc create mode 100644 src/decoding/__pycache__/evm.cpython-310.pyc create mode 100644 src/decoding/evm.py create mode 100644 src/ingestion/fetcher.py create mode 100644 src/projections/__pycache__/__init__.cpython-310.pyc create mode 100644 src/projections/__pycache__/store.cpython-310.pyc create mode 100644 src/projections/store.py create mode 100644 tests/__pycache__/test_indexer_pipeline.cpython-310-pytest-8.4.1.pyc create mode 100644 tests/test_indexer_pipeline.py diff --git a/src/__pycache__/__init__.cpython-310.pyc b/src/__pycache__/__init__.cpython-310.pyc new file mode 100644 index 0000000000000000000000000000000000000000..499a2df3cfa933c25a5ca6586b523b0517679703 GIT binary patch literal 148 zcmd1j<>g`kg0I0-vOx4>5P=LBfgA@QE@lA|DGb33nv8xc8Hzx{2;!HyvsFxJacWU< zOmRU@YGRB_YH@Z+enCulh+9NVc}bdXW?o8aMQTw@aZz$ie0*kJW=VX!UP0w84x8Nk Rl+v73JCK3JOhAH#0RSvhB6g`kg0I0-vOx4>5P=LBfgA@QE@lA|DGb33nv8xc8Hzx{2;!HGvsFxJacWU< zOmRU@YGRB_YH@Z+enCulh+9NVc}bdXW?o8aMQTw@aZz$ia(+>2OniK1US>&ryk0@& WEe@O9{FKt1R6CHV#Y{kgg#iHd=Olyx literal 0 HcmV?d00001 diff --git a/src/core/__pycache__/abis.cpython-310.pyc b/src/core/__pycache__/abis.cpython-310.pyc new file mode 100644 index 0000000000000000000000000000000000000000..4b2c3ce98a9f59b1074234305da1baf2b93479e1 GIT binary patch literal 1049 zcmah|&1(}u6rY)$B%4i|eu&b8Xb)Z@7&Kmrh*;AK(l2PM)g>itcP4S0&F(rg{R&$V zsvv^s!J8DCiy)r77W$9ORf~G{C@8+&SP?{=W!{_j=J$U4_BXTPz<`Bdd|i39HbfEn zkOY4Y$f4^n81|9f_t?e|6Q6Vp+vrf6V#ExdVU%Z^lx0}=%RBv^ zGBN34@la)B!b)tA4LvtroZc~1Zr_ZJ9^RoUA5&%R648$7et`|gv_}FQ(`7^u8=;7z zT?4}!qeXNH$?Z*yVU5u}m=C5~kS(8BKLu(WJ=_0$czE>W{l|lMM=zg#I(RcSapP8( zR4U!#rfObjiZD<-VA_lvx9zmJ)OinD~ zh~$>)QWQ$*tpZ2FFEx%dZUW9x&RZnT0OZ?v?MnT&z`SJgwp?rr*s>5f5PtkK%Jb$Rlcq zG1#;RfxDJyo$7U`Y8N0en~DWNel~R5Qwz1~#d1NM1iUDL(OI<}fog&z86(LUQVCE2 zgh^?2sWv~iv@kz2Gquowv)0z^#KQE{QnI#7v-7p-M#`GR>8Jl;&L({3)lSIPeSR(Z T!({SO6PeKd2r&muT*N;ClT;+( literal 0 HcmV?d00001 diff --git a/src/core/__pycache__/events.cpython-310.pyc b/src/core/__pycache__/events.cpython-310.pyc new file mode 100644 index 0000000000000000000000000000000000000000..30e1d67c9ef4318a6ca6aeab75b157332f408475 GIT binary patch literal 1082 zcmZ8f&2JM&6rcU@dhNtbKxvh7@-d6j9I76w2q9c&6GhmulsGL~S(*+zlVoAPm>E+X z^$;~_Xa%hhAP$WxaYE^#ih5~N5dA~uO5DJ|z^QMHp%ycm_ulW#d-LXfWR*%8!8+4^ zcW~ z(_i-feSEurd%yqu{>kx9|L&ucr*}@C{aKx9kRIdJPd(*^TO>+xMAi7!I18`iw3}RI zTovUlyhX&8D#h6*j#A2QZ7;F8&ADo`wtCrbHwiee~jW2!NT5s1|O@FppZ`78Xo}tDl6CsZVRaRxDbHK8fGc_4zspKS-m{7`@ z5UQAyUYrrCN-|2AkR-{K!

zG&Dw9qvZmeIvjBh$hYVkI)KgqH-_9iK)>oaFuUi$ zp#S}0|Bs#i!y{drzrQ>E{YR+Q%cK43%ra$3o=KL5z3OGw^IjYCp2c18;d)H@7!b(9 z7A)TQ4e=>}ujlLYopr&u=!iUKq_e=pb(!UzuiA@OI^W2RPXoMY$NKBw-keeiY9upxh|k`bDrgl%O3wo>vVP__gP6*Jc1P*v|6#eK@AShX!3ME)MCD$1ljTL- NXQ6E;(UdV|{s-*pExZ5# literal 0 HcmV?d00001 diff --git a/src/core/abis.py b/src/core/abis.py new file mode 100644 index 0000000..dc530eb --- /dev/null +++ b/src/core/abis.py @@ -0,0 +1,22 @@ +import json +import os +from typing import Dict, Any + +def load_contract_abi(name: str) -> list: + """加载共享合约 ABI""" + # 优先从 frontend packages 读取 + possible_paths = [ + os.path.join(os.path.dirname(__file__), "..", "..", "..", "wtf-frontend", "packages", "contracts", "src", "abis", f"{name}.json"), + os.path.join(os.path.dirname(__file__), "..", "..", "..", "wtf-contract", "artifacts", "main", "src", "controllerv2", f"{name}.sol", f"{name}.json"), + os.path.join(os.path.dirname(__file__), "..", "..", "..", "wtf-contract", "artifacts", "main", "src", "marketv2", f"{name}.sol", f"{name}.json") + ] + for p in possible_paths: + if os.path.exists(p): + with open(p, "r", encoding="utf-8") as f: + data = json.load(f) + return data.get("abi", data) if isinstance(data, dict) else data + return [] + +CONTROLLER_ABI = load_contract_abi("WTFControllerV2") +MARKET_ABI = load_contract_abi("WTFMarketV2") +MOCK_ERC20_ABI = load_contract_abi("MockERC20") diff --git a/src/decoding/__pycache__/__init__.cpython-310.pyc b/src/decoding/__pycache__/__init__.cpython-310.pyc new file mode 100644 index 0000000000000000000000000000000000000000..adf0681fb06bbe8b3b9c8d073ed19933ee55ce89 GIT binary patch literal 157 zcmd1j<>g`kg0I0-vOx4>5P=LBfgA@QE@lA|DGb33nv8xc8Hzx{2;!H6vsFxJacWU< zOmRU@YGRB_YH@Z+enCulh+9NVc}bdXW?o8aMQTw@aZz$iN@{X`N@iYqOniK1US>&r ayk0@&Ee@O9{FKt1R6CH##Y{kgg#iGa9wrvcm$E~ClbYkT6^ zna!PX{GnMCWWx(BDpExSrPNl1LY~SS6+G}a@OodVo94MZK!gP6+_4=uO_7+@J#+5I z+*}`NHEg4s zu`?9sbl>b|?X2Pr-|FVTnzB^zKDBko7-J)IWmh2KGobHu*WuMxER|v~8YnL!< zN3(|*T_CmGeK?uS)pTi0dYxG6^hX`(sTZZ&6S8ZD>(y=j@Bt+#&AXd~aj{ZBp?bY6iLeIE!XsHd$lYU@nn z#Sf1zkz&Kep*#Il-X6ynhf+l5|Z@e}; zsiYL^FT&=~t%3yAPu}_c!C&9`^v>?5Z|oer@vBeod~k2?CkMa&L-qWOyTwKI+b>ns zXI`*UlU8Td^#aFZ9SycPdXUF90WCTMB%uj$NKAJKytRbv>H7pt4}#FSnjr>YjmQ9z z+QtiV(7C?u`OJyKp4V}@ZZ86(84>rF)xx_(lws+EmnU9YT8y}emZF}|-K9w$t-&Wt zSKCuRT-u12zu^UpZ*s8|iO!O`4-k2Y-|AlI!KFEl7kIJblwhmK1Y%N)7U=fyL$h2s zTGOF0fAnO5!SzbH1Q-vM0_9?Dj7*Z(U+mL_?&*x~Q>LwH;$_v-_i0~KoUuaowZ5L- zSfWpo9ojb%?Oo+ZJt2lF0le-;3f-*zH|@0yo2*N67ZYla_n*1}|Ll!;b4mmw`F@_yIlt zfIk0#)-Z5lz8TjnSyXB`?sX3@UZDGn93HB#aM12X=mnwxecN2)osPRE^;N!E%L=?k zsl#UiRTHPMn?b`wQjfe9X)0&B_=qZ^2s3ACr!cQbrEzBw~ii+20T=G&oWpJd;*Xx^@&7}@(gYv zrisq`zc@_SY|lLt97XLpOPWyYYb>I z{1z33q==LaE6IyWkCE-K#RJJmGQ6YpPsOkzIkf{CzJ3PeFZV}buTdrA0jngVDCHPA zA>ZGU(UZNaBzV#wF`TjmS*AZm$@XAe+AFh>#Dw@hm5jYh-h=z2 z6h`$K=V~2_qxSLKxSE46+Rc`;@M^vFaI^A!tJ!W&G-sUo`Xb0@4;~LJ&ZR@Of?b&8 zJwM!tA(!A9gR}2u#PJSXgxhx6OHXJ>tx^7K2{O8Q_~v z%~I8jETm~8P~jr8@uss1K_qkFtx%#s;K+QM6IGh8bV59=6FWx%nC-b+e&{lpJI*|{ zVHGDb??qk!xhUvBeQ?J>@ANGUJCETtIe=h0Jj{o(s?pQW+0q9zGIrwBg08 zGP^9oE+kojMUT|Cy`Ic9HalDuz0w3RS74hA7b1k$danl@)<{ab+8LFZMGXhOj^ueH z=aD=M#6Hvws}{TXIo-r}anU6tmyy(QdFc=dz%p3i*hApmsg*c89Fw}oHh}~^J15+Y zw2HJdv4c7wLfCl>11Gh|b2v~)LxnlFphA&4P)WI6K!aly(u2CuE(gbTyn_>u;E&8d z1eq%A8Qlg_&gEbBY8P! z^DlHl!7r$IiIgyQ@pu##Mkfl&92Ar@%6Y{i3O9=c1$&(6(z=vd=Q3h=gf+57EVwef M&H}B_yf&o$7kEL>6#xJL literal 0 HcmV?d00001 diff --git a/src/decoding/evm.py b/src/decoding/evm.py new file mode 100644 index 0000000..afc21e1 --- /dev/null +++ b/src/decoding/evm.py @@ -0,0 +1,109 @@ +import json +import logging +from typing import Dict, Any, Optional +from web3 import Web3 +from eth_abi import decode +from src.core.events import NormalizedEvent +from src.core.abis import CONTROLLER_ABI, MARKET_ABI + +logger = logging.getLogger(__name__) + +class EVMDecoder: + """EVM ABI 日志规范化解码器 (Layer 2: Decoding)""" + + def __init__(self, chain_id: int): + self.chain_id = chain_id + self.w3 = Web3() + self._build_topic_maps() + + def _build_topic_maps(self): + self.event_abi_map = {} + + # 解析 Controller 与 Market ABI 中的 events + for abi in CONTROLLER_ABI + MARKET_ABI: + if abi.get("type") == "event": + name = abi.get("name") + inputs = abi.get("inputs", []) + types = [i["type"] for i in inputs] + sig = f"{name}({','.join(types)})" + topic0 = self.w3.keccak(text=sig).hex() + 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 + + topic0 = topics[0].hex() if isinstance(topics[0], bytes) else topics[0] + abi = self.event_abi_map.get(topic0) + if not abi: + return None + + event_name = abi["name"] + contract_addr = log.get("address", "").lower() + tx_hash = log.get("transactionHash", "").hex() if isinstance(log.get("transactionHash"), bytes) else str(log.get("transactionHash", "")) + 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 + + # 2. Non-indexed data + data = log.get("data", "0x") + if isinstance(data, str) and data.startswith("0x"): + data_bytes = bytes.fromhex(data[2:]) + elif isinstance(data, bytes): + data_bytes = data + else: + data_bytes = b"" + + if data_bytes and non_indexed_inputs: + types = [i["type"] for i in non_indexed_inputs] + try: + decoded_vals = decode(types, data_bytes) + for inp, val in zip(non_indexed_inputs, decoded_vals): + if isinstance(val, bytes): + payload[inp["name"]] = "0x" + val.hex() + else: + payload[inp["name"]] = val + 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", + "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", + } + + normalized_type = type_mapping.get(event_name, f"EVM_{event_name.upper()}") + + return NormalizedEvent( + chain_id=self.chain_id, + block_number=block_number, + tx_hash=tx_hash, + log_index=log_index, + event_type=normalized_type, + contract_address=contract_addr, + payload=payload, + timestamp=block_timestamp + ) diff --git a/src/ingestion/fetcher.py b/src/ingestion/fetcher.py new file mode 100644 index 0000000..9381c42 --- /dev/null +++ b/src/ingestion/fetcher.py @@ -0,0 +1,50 @@ +import asyncio +import logging +from typing import List, Dict, Any, Optional +from web3 import AsyncWeb3, AsyncHTTPProvider +from web3.types import LogReceipt, BlockNumber +from src.config import indexer_settings + +logger = logging.getLogger(__name__) + +class EVMFetcher: + """EVM 批量日志与区块抓取器 (Layer 1: Ingestion)""" + + def __init__(self, rpc_url: Optional[str] = None, chain_id: Optional[int] = None): + self.rpc_url = rpc_url or indexer_settings.RPC_URL + self.chain_id = chain_id or indexer_settings.CHAIN_ID + self.w3 = AsyncWeb3(AsyncHTTPProvider(self.rpc_url)) + + async def get_latest_block_number(self) -> int: + """获取链上最新安全区块高度""" + block_num = await self.w3.eth.block_number + return max(0, block_num - indexer_settings.CONFIRMATIONS) + + async def fetch_logs( + self, + from_block: int, + to_block: int, + addresses: Optional[List[str]] = None, + topics: Optional[List[Any]] = None + ) -> List[LogReceipt]: + """批量获取指定区块区间的日志""" + filter_params: Dict[str, Any] = { + "fromBlock": from_block, + "toBlock": to_block + } + if addresses: + filter_params["address"] = [self.w3.to_checksum_address(a) for a in addresses] + if topics: + filter_params["topics"] = topics + + try: + logs = await self.w3.eth.get_logs(filter_params) + return logs + except Exception as e: + logger.error(f"Error fetching logs from {from_block} to {to_block}: {e}") + raise e + + async def get_block_timestamp(self, block_number: int) -> int: + """获取区块出块时间戳""" + block = await self.w3.eth.get_block(block_number) + return block.get("timestamp", 0) diff --git a/src/main.py b/src/main.py index fdd70cc..2b9608d 100644 --- a/src/main.py +++ b/src/main.py @@ -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()) diff --git a/src/projections/__pycache__/__init__.cpython-310.pyc b/src/projections/__pycache__/__init__.cpython-310.pyc new file mode 100644 index 0000000000000000000000000000000000000000..f17afdb13a5e2a90b77c950218ab04025b0684e1 GIT binary patch literal 160 zcmd1j<>g`kg0I0-vOx4>5P=LBfgA@QE@lA|DGb33nv8xc8Hzx{2;!HsvsFxJacWU< zOmRU@YGRB_YH@Z+enCulh+9NVc}bdXW?o8aMQTw@aZz$iK~a8IYH~?teqM1*e0*kJ dW=VX!UP0w84x8Nkl+v73JCNbUOhAH#0RSmcCz1dF literal 0 HcmV?d00001 diff --git a/src/projections/__pycache__/store.cpython-310.pyc b/src/projections/__pycache__/store.cpython-310.pyc new file mode 100644 index 0000000000000000000000000000000000000000..a10388db9c61138a530de4935850c071178ccc75 GIT binary patch literal 2657 zcma(SNo*WNaC)wtU9S(Egd`-uak2))Sq!K}}{(s5C>ZBf6cTSfO)8HG!8kdp93km0({HOGp#cx+lN29boONHc(Z-{{no{=mX9yB55Ls<{KMAE zi-p?<-5JIU_dQft7ZMfa;nw`(^6Ar;&K_R*aH)Of>q{5EUp}$eI(8IJ-qW5x-J1XM zw~OajUVUr%{Cn+Bj$FR*X6xIf_K7drZ@$!e`<=qxFp4KR8+~H1&^q=}>*%M6a$rk% zR!+Xv{$Qyv!Xg$IZiThxyE9U+M>e2~StPQxFsh(`BB6mk(Np80$D)W)hhadXQ$*+y^CwOzybF~4`N-~v z#-EHBkH({#&)o46i>BhRHvaV3qtA@L5Kr7)2`D?n_;|#^Plhxdp%^IPzq6+q$b|y!uvIuedqpg}nLwfHaf?$^vOB1Tl3D)7CIO!D3^E-q8!Bda1_TSk0|#-AH{@9a9Y{6(aB`Ayp;v;5ki$Dn94o!cyOG0?LN@^# z;e*YoIfA2$%x+oNpZG^|3PXmq3?qF0#>NB6Li&9CC?f7ccG)Hl$2*T&zTJqZ0G{X25?nrcf+JTyA;WN z6;MsuyGuE^FO{?jc;5eClHR48uaL9_B)ynQx*ELn3s{n0=u9og8<=42R?x9Xw*`c5 zmzu5t|9m5{9c$Rm)DzbNLQ1&-C_hWejVa|OdNcSBeRaz`co1?G z{DiZ3YkdbqbiHz&d*jmir6grr=T1Ng8WP=m`bQowk2$6C-ofF0<&rRy2QP95>ny@| zK0qNa>wKDtoEQ4O8#C_uBI9x9!UGk0Tw#0#Ym1$kU~B~b3$uUs*uYceV&3jbq&Y67 zka>~lSe5EjDBli;_%#U7OuiFkSRljUQb-iLaXk{2^-0(=LJmTyb7#g$l2|iLTmCJkhZNOmb z@`ss_gkFcDXdCKDA^3Yoah@X@m=}`Mtjy)`Qr#T)7Sec0LZLSP-T3 zY(gKdL+1oXatsS}wM=}-dESkl7kTL~C#YBVGtRf--Sv2;#`x`MUuUWusx4@@uUc2@ zBy#P}L|eM$cfcurBkJf9mMke4R#~((h`7ilG}It6HFw4jT`IIm>>6GeKtu83)y&{` zp)4$~QeaXK&1u(%ER+RNm~O4c0t&Sf+W_AJ0P2{u{$9ZCjt!}1>j%LSL8?iS#JZML z$_^0X&m^4+lP;nF@~RHC(o{`isXathEs~5JtSSs$j-KDGZmg+YK5%jO#oQfWm~Wmomhrt_F9r&><}63Ebz3TxzE&K<=7B7c1o@T y5Z%y{`#kKzNo$wdL>2-#C|eSL5LtT>pr-#aL0G$z%RG#^nA))!T6a$ERsR7HBN$Ks literal 0 HcmV?d00001 diff --git a/src/projections/store.py b/src/projections/store.py new file mode 100644 index 0000000..d20f0bb --- /dev/null +++ b/src/projections/store.py @@ -0,0 +1,79 @@ +import logging +from typing import Dict, Any, Optional +from src.core.events import NormalizedEvent + +logger = logging.getLogger(__name__) + +class ProjectionStore: + """ + 确定性业务投影状态存储 (Layer 4: Projections) + 在第一阶段使用内存/本地轻量结构,并支持向 PostgreSQL 写入投影,支持随时 Reset & Replay + """ + + def __init__(self): + self.markets: Dict[str, Dict[str, Any]] = {} + self.trades: list = [] + self.positions: Dict[str, Dict[str, Any]] = {} + self.processed_events: set = set() + + def reset(self): + """全量清空投影视图(支持 Replay)""" + self.markets.clear() + self.trades.clear() + self.positions.clear() + self.processed_events.clear() + logger.info("Projection store reset successfully.") + + def apply_event(self, event: NormalizedEvent): + """确定性事件投影应用""" + if event.event_id in self.processed_events: + 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 + self.markets[market_addr] = { + "market_address": market_addr, + "question_id": payload.get("questionId"), + "curve": payload.get("curve"), + "collateral": payload.get("collateral"), + "creator": payload.get("creator"), + "tier": payload.get("tier", 1), + "fee_rate": payload.get("feeRate"), + "status": "ACTIVE", + "winning_outcome": None, + "created_at_block": event.block_number, + "created_at_tx": event.tx_hash, + "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"): + self.trades.append({ + "tx_hash": event.tx_hash, + "block_number": event.block_number, + "market": event.contract_address, + "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 + }) diff --git a/tests/__pycache__/test_indexer_pipeline.cpython-310-pytest-8.4.1.pyc b/tests/__pycache__/test_indexer_pipeline.cpython-310-pytest-8.4.1.pyc new file mode 100644 index 0000000000000000000000000000000000000000..7d6828632a1b66b8256886ccc1b634e1ea5cf0ad GIT binary patch literal 3774 zcmcInOKjZ68Rn4O7xyL0v12)oBRjUMO&fLJyE3XGdaR3}@dLKvri4&};*4zCTyi@@ z$&$jo)IAhMTlmnU4${#;5ACJb_Rv#Kz4XFkd#R2^(2G$N{r}-gS}6$v6wWO8a~}W9 z{~G@Pw?VOJD9|qdK&@8-5lsr%>xdS*~xlb$TY$_fq@H zrXp&vo{p4#!rRcxV9zJ8j(xIE5<8Tv})KcC*v! zber{Vqxw=e(|wyiWnpambb|>D+=tZELznS5<8Wp}6c3btJzy{4rZ|sn*SB117S_DL zd20C^{WZqTOt@t|w#8#J=LL@}*QabNHW=OrE8HG1bHWMyklRjZ*_3i7M4TDe+g@PP zI2XEokZSh_JLjA2cB6Xpq#JM6Uy@hq{CtJ%RIkqrw(AvPd^+k(&rf5oS`o8}Fqay` z>G@=~RS|JYPFJ*811=`AE4MybeHO02FT=+8Dw-@LO(My-HE;*abA2{D z*p9DAm-t5?zIEs?)*l2R^PO${uEk{_Dmn2}lac-(HJ$$aH)TExEGH5|877Vr!b}Pi z!gFP)gyd_|0*RD8NC^7`5>^IMgYra!#GppnCVbA3E^&-9pPs` z5oY$XGEvOKTuvs6MY)QEof#+ zp>IqnkS@**GaxwwDGN1v0lQN4ywvDToFCUn%FBSf3%e>Lp)8%E7o)6{mmA5;eL)_i zwL>3B`}MF8>Ceefqq(6GX^{bv^HE-WK9ZaVU%Z9$GFe_a^1+{AJr@=rS^fp}z#b&a zOJiR8D{@FL{S7%r=``_@4(Sxc=AI- z#(SDxf#iP`$^aIDcXw5KO_qR?w6lb$Y`5>;G7E@_F%E6DXtSaSTuLGSDDq$+*|n!M&aDS^n$ z9S-SI;h1`P1_|N-o`sDD@%88_pJN9 z2X{Yu*jomBCh(+8un_`qVD2yU6=r4+m(3}tnBW_ClUdJP-*^2-)`$}fefa+wI^cl= z7>m)QG((6^T3jh}z^h6Eck=dN;B6Lu!7)xGkYoxyIBdC+xF95f=4 zof8EnN`c@ns}?5z2R