titikshaha's picture
new files
8f82127
Raw
History Blame Contribute Delete
8.69 kB
#old version, do not refer
# indexer.py
import time
import json
import argparse
from decimal import Decimal, getcontext
from web3 import Web3
import sqlite3
from typing import List
# --- Configuration ---
RPC_URL = "https://polygon-rpc.com/"
POL_TOKEN = Web3.to_checksum_address("0x455e53CBB86018Ac2B8092FdCd39d8444aFFC3F6")
# Binance addresses provided (checksum)
BINANCE_ADDRESSES = {
Web3.to_checksum_address("0xF977814e90dA44bFA03b6295A0616a897441aceC"),
Web3.to_checksum_address("0xe7804c37c13166fF0b37F5aE0BB07A3aEbb6e245"),
Web3.to_checksum_address("0x505e71695E9bc45943c58adEC1650577BcA68fD9"),
Web3.to_checksum_address("0x290275e3db66394C52272398959845170E4DCb88"),
Web3.to_checksum_address("0xD5C08681719445A5Fdce2Bda98b341A49050d821"),
Web3.to_checksum_address("0x082489A616aB4D46d1947eE3F912e080815b08DA"),
}
# Decimal precision (enough for 18 decimals)
getcontext().prec = 40
# ERC20 Transfer event signature topic
TRANSFER_TOPIC = Web3.keccak(text="Transfer(address,address,uint256)").hex()
# --- DB helpers ---
DB_PATH = "src/pol_indexer.sqlite"
SCHEMA_SQL = open("src/schema.sql", "r").read()
def get_db_conn():
conn = sqlite3.connect(DB_PATH, timeout=30)
conn.row_factory = sqlite3.Row
return conn
def init_db():
conn = get_db_conn()
cur = conn.cursor()
cur.executescript(SCHEMA_SQL)
# initialize meta last_processed_block to current block if missing
cur.execute("SELECT value FROM meta WHERE key = 'last_processed_block'")
row = cur.fetchone()
if row is None:
# set last_processed_block to current chain head to avoid backfill
w3 = Web3(Web3.HTTPProvider(RPC_URL))
head = w3.eth.block_number
cur.execute("INSERT OR REPLACE INTO meta (key, value) VALUES (?, ?)", ("last_processed_block", str(head)))
print(f"[DB] Initialized last_processed_block = {head} (real-time only, no backfill)")
conn.commit()
conn.close()
def get_last_processed_block() -> int:
conn = get_db_conn()
cur = conn.cursor()
cur.execute("SELECT value FROM meta WHERE key = 'last_processed_block'")
v = int(cur.fetchone()["value"])
conn.close()
return v
def set_last_processed_block(n: int):
conn = get_db_conn()
cur = conn.cursor()
cur.execute("INSERT OR REPLACE INTO meta (key, value) VALUES (?, ?)", ("last_processed_block", str(n)))
conn.commit()
conn.close()
def insert_raw_log(tx_hash, log_index, block_number, address, topics, data):
conn = get_db_conn()
cur = conn.cursor()
try:
cur.execute("""
INSERT OR IGNORE INTO raw_logs (tx_hash, log_index, block_number, address, topics, data, received_at)
VALUES (?, ?, ?, ?, ?, ?, ?)
""", (tx_hash, log_index, block_number, address, json.dumps(topics), data, int(time.time())))
conn.commit()
finally:
conn.close()
def insert_transfer(tx_hash, block_number, log_index, from_addr, to_addr, amount_raw, amount_dec, token_address):
conn = get_db_conn()
cur = conn.cursor()
try:
cur.execute("""
INSERT OR IGNORE INTO transfers
(tx_hash, block_number, log_index, from_address, to_address, amount_raw, amount_dec, token_address, received_at)
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?)
""", (tx_hash, block_number, log_index, from_addr, to_addr, str(amount_raw), str(amount_dec), token_address, int(time.time())))
conn.commit()
finally:
conn.close()
def get_latest_cumulative() -> Decimal:
conn = get_db_conn()
cur = conn.cursor()
cur.execute("SELECT cumulative_value FROM netflow ORDER BY created_at DESC LIMIT 1")
row = cur.fetchone()
conn.close()
if row:
return Decimal(row["cumulative_value"])
return Decimal("0")
def append_cumulative(new_cum: Decimal):
conn = get_db_conn()
cur = conn.cursor()
cur.execute("INSERT INTO netflow (created_at, cumulative_value) VALUES (?, ?)", (int(time.time()), str(new_cum)))
conn.commit()
conn.close()
# --- Web3 setup ---
w3 = Web3(Web3.HTTPProvider(RPC_URL))
# Minimal ABI for Transfer event decoding (so we can use contract.events)
ERC20_TRANSFER_ABI = [{
"anonymous": False,
"inputs": [
{"indexed": True, "name": "from", "type": "address"},
{"indexed": True, "name": "to", "type": "address"},
{"indexed": False, "name": "value", "type": "uint256"}
],
"name": "Transfer",
"type": "event"
}]
token_contract = w3.eth.contract(address=POL_TOKEN, abi=ERC20_TRANSFER_ABI)
# --- Processing logic ---
def decode_transfer_log(log) -> dict:
"""
Accepts a log (dict) and decodes it to a Transfer event using web3.
Returns dict with keys: from, to, value (int)
"""
# use process_log which needs the log in web3 format
ev = token_contract.events.Transfer().process_log(log)
return {
"from": Web3.to_checksum_address(ev["args"]["from"]),
"to": Web3.to_checksum_address(ev["args"]["to"]),
"value": int(ev["args"]["value"])
}
def process_logs_for_block(block_number: int) -> Decimal:
"""
Fetch logs for POL token for a single block, store raw + normalized transfers,
and return the net delta (Decimal) to add to cumulative netflow for Binance:
+value if to_address in BINANCE_ADDRESSES
-value if from_address in BINANCE_ADDRESSES
"""
delta = Decimal("0")
try:
logs = w3.eth.get_logs({
"fromBlock": block_number,
"toBlock": block_number,
"address": POL_TOKEN
})
except Exception as e:
print(f"[RPC] get_logs failed for block {block_number}: {e}")
return delta
for log in logs:
# insert raw
tx_hash = log["transactionHash"].hex()
log_index = log["logIndex"]
insert_raw_log(tx_hash, log_index, block_number, log["address"], [t.hex() for t in log["topics"]], log["data"])
# decode transfer
try:
decoded = decode_transfer_log(log)
except Exception as e:
print(f"[Decode] Failed to decode log {tx_hash} idx {log_index}: {e}")
continue
amt_raw = decoded["value"] # integer
# convert to decimal string (POL has 18 decimals). Use Decimal for precision.
amt_dec = (Decimal(amt_raw) / Decimal(10**18)).normalize()
insert_transfer(tx_hash, block_number, log_index, decoded["from"], decoded["to"], amt_raw, str(amt_dec), POL_TOKEN)
# compute delta for Binance
if decoded["to"] in BINANCE_ADDRESSES:
delta += (Decimal(amt_raw) / Decimal(10**18))
if decoded["from"] in BINANCE_ADDRESSES:
delta -= (Decimal(amt_raw) / Decimal(10**18))
print(f"[POL Transfer] blk {block_number} | {decoded['from']} -> {decoded['to']} | {amt_dec} POL")
return delta
# Main real-time loop
def start_realtime_loop(poll_interval: float = 1.5):
print("[Indexer] Starting real-time loop")
init_db()
last = get_last_processed_block()
while True:
try:
head = w3.eth.block_number
except Exception as e:
print("[RPC] cannot get block_number:", e)
time.sleep(5)
continue
if head > last:
# process blocks one by one (small windows)
for blk in range(last + 1, head + 1):
print(f"[Indexer] Processing block {blk} / {head}")
delta = process_logs_for_block(blk) # Decimal
if delta != Decimal("0"):
# update cumulative netflow incrementally
current = get_latest_cumulative()
new_cum = (current + delta).normalize()
append_cumulative(new_cum)
print(f"[Netflow] Delta {delta} POL applied. New cumulative: {new_cum} POL")
# update last processed block even if no logs, so we don't reprocess
set_last_processed_block(blk)
last = blk
else:
# nothing new
time.sleep(poll_interval)
# CLI to get current netflow
def get_current_netflow() -> Decimal:
init_db()
return get_latest_cumulative()
if __name__ == "__main__":
parser = argparse.ArgumentParser(description="POL -> Binance real-time indexer")
parser.add_argument("cmd", choices=["start", "netflow"], help="start = run indexer, netflow = print current cumulative")
args = parser.parse_args()
if args.cmd == "start":
start_realtime_loop()
elif args.cmd == "netflow":
v = get_current_netflow()
print(f"Cumulative POL netflow to Binance: {v} POL")