# Author: Simon-Pierre Boucher # Mail: contact@spboucher.ai """Stablecoin explorer API — FastAPI over the indexer's SQLite DB, plus live pass-through reads (blocks, txs) straight from free RPCs. uvicorn api.main:app --port 8080 """ import asyncio import functools import os import time import pathlib from fastapi import FastAPI, HTTPException, Query, WebSocket, WebSocketDisconnect from fastapi.staticfiles import StaticFiles from indexer import config, db from indexer.decode import TRANSFER_TOPIC, decode_transfer, format_amount from indexer.enrich import ZERO_ADDRESSES from indexer.rpc import RpcPool DB_PATH = config.db_path() CHAINS = config.load_chains() TOKENS = config.load_tokens() app = FastAPI( title="coinexplorer API", description="Self-hosted, free-RPC-only explorer for stablecoins and major crypto.", version="0.3.0", ) @app.exception_handler(Exception) async def unhandled_exception(request, exc): # never leak a stack trace; always answer JSON from fastapi.responses import JSONResponse return JSONResponse(status_code=500, content={"error": str(exc)[:200]}) _pools = {} def pool(chain): if chain not in CHAINS: raise HTTPException(404, f"unknown chain '{chain}'") if CHAINS[chain].get("family", "evm") != "evm": raise HTTPException(501, f"live reads for '{chain}' land with its adapter (Deliverable 4)") if chain not in _pools: _pools[chain] = RpcPool(CHAINS[chain]["rpcs"]) return _pools[chain] def dbc(): return db.connect(DB_PATH, readonly=True) # stablecoins trade ~$1, so face value ≈ USD; CASE avoids relying on POW() # (not compiled into every SQLite build) USD_EXPR = ( "CAST(amount AS DOUBLE PRECISION) / (CASE decimals " + " ".join(f"WHEN {d} THEN 1e{d}" for d in range(19)) + " ELSE 1e6 END)" ) # price-aware USD: stablecoins are seeded at 1.0 by the PriceWorker; crypto # assets get their CoinGecko price; unknown symbols fall back to face value PRICE_JOIN = "LEFT JOIN prices ON prices.symbol = transfers.symbol" USD_PRICED = f"(({USD_EXPR}) * COALESCE(prices.usd, 1.0))" SELECT_T = f"SELECT transfers.*, COALESCE(prices.usd, 1.0) AS usd_rate FROM transfers {PRICE_JOIN}" DEFAULT_WHALE_USD = float(os.environ.get("WHALE_MIN_USD", 1_000_000)) def price_of(conn, symbol): r = conn.execute("SELECT usd FROM prices WHERE symbol = ?", (symbol,)).fetchone() return r["usd"] if r and r["usd"] is not None else 1.0 _cache = {} def cached(ttl): """Tiny in-process TTL cache for aggregate endpoints.""" def deco(fn): @functools.wraps(fn) def wrap(*a, **k): key = (fn.__name__, a, tuple(sorted(k.items()))) hit = _cache.get(key) now = time.time() if hit and now - hit[0] < ttl: return hit[1] val = fn(*a, **k) _cache[key] = (now, val) return val return wrap return deco def parse_window(window): units = {"m": 60, "h": 3600, "d": 86400} try: return int(window[:-1]) * units[window[-1]] except (KeyError, ValueError, IndexError): raise HTTPException(400, "window must look like 15m, 24h or 7d") def with_value(row): d = dict(row) rate = d.pop("usd_rate", None) if d.get("amount") is not None and d.get("decimals") is not None: d["value"] = format_amount(d["amount"], d["decimals"]) if rate is not None: d["usd"] = round(float(d["value"]) * rate, 2) return d def token_meta(chain, address): for t in TOKENS.get(chain, []): if t.get("address", "").lower() == address.lower(): return t return None def block_summary(blk): return { "number": int(blk["number"], 16), "hash": blk["hash"], "parent_hash": blk["parentHash"], "timestamp": int(blk["timestamp"], 16), "tx_count": len(blk.get("transactions", [])), "gas_used": int(blk.get("gasUsed", "0x0"), 16), } # -- meta --------------------------------------------------------------- @app.get("/v1/chains") def chains(): return { chain: { "family": cfg.get("family", "evm"), "chain_id": cfg.get("chain_id"), "block_time_s": cfg.get("block_time"), "confirmations": cfg.get("confirmations"), "tokens": TOKENS.get(chain, []), } for chain, cfg in CHAINS.items() } @app.get("/v1/status") def status(): conn = dbc() try: cursors = { r["chain"]: { "last_block": r["last_block"], "last_hash": r["last_hash"], "backfill_block": r["backfill_block"], "head_block": r["head_block"], "lag": (r["head_block"] - r["last_block"]) if r["head_block"] is not None else None, } for r in conn.execute("SELECT * FROM cursors") if not r["chain"].startswith("_") # internal scan cursors } counts = { r["chain"]: r["n"] for r in conn.execute( "SELECT chain, COUNT(*) AS n FROM transfers GROUP BY chain" ) } return {"cursors": cursors, "indexed_transfers": counts} finally: conn.close() @app.get("/v1/rpc/health") def rpc_health(): """Per-endpoint health as last reported by the indexer workers: score = success-rate EMA, cooldown_s > 0 = currently benched.""" conn = dbc() try: rows = conn.execute( "SELECT * FROM rpc_health ORDER BY chain, score DESC" ).fetchall() finally: conn.close() out = {} for r in rows: out.setdefault(r["chain"], []).append( {k: r[k] for k in ("url", "score", "ok", "fail", "latency_ms", "cooldown_s", "updated")} ) return out @app.get("/metrics") def metrics(): """Prometheus text exposition: per-chain progress/lag + RPC endpoint health. Free-RPC degradation shows up here as rising lag and falling rpc scores.""" from fastapi.responses import PlainTextResponse conn = dbc() try: lines = [ "# HELP explorer_transfers_total Indexed transfers per chain", "# TYPE explorer_transfers_total counter", ] for r in conn.execute("SELECT chain, COUNT(*) n FROM transfers GROUP BY chain"): lines.append(f'explorer_transfers_total{{chain="{r["chain"]}"}} {r["n"]}') lines += [ "# HELP explorer_cursor_block Last indexed block per chain", "# TYPE explorer_cursor_block gauge", "# HELP explorer_lag_blocks Head minus cursor per chain", "# TYPE explorer_lag_blocks gauge", ] for r in conn.execute("SELECT * FROM cursors"): lines.append(f'explorer_cursor_block{{chain="{r["chain"]}"}} {r["last_block"]}') if r["head_block"] is not None: lines.append(f'explorer_lag_blocks{{chain="{r["chain"]}"}} ' f'{max(0, r["head_block"] - r["last_block"])}') lines += [ "# HELP explorer_rpc_score Endpoint health score (success-rate EMA, 0-1)", "# TYPE explorer_rpc_score gauge", "# HELP explorer_rpc_calls_total Successful/failed calls per endpoint", "# TYPE explorer_rpc_calls_total counter", ] for r in conn.execute("SELECT * FROM rpc_health"): lbl = f'chain="{r["chain"]}",url="{r["url"]}"' lines.append(f"explorer_rpc_score{{{lbl}}} {r['score']}") lines.append(f'explorer_rpc_calls_total{{{lbl},result="ok"}} {r["ok"]}') lines.append(f'explorer_rpc_calls_total{{{lbl},result="fail"}} {r["fail"]}') return PlainTextResponse("\n".join(lines) + "\n") finally: conn.close() # -- live chain reads (straight off the free RPCs) ----------------------- @app.get("/v1/{chain}/block/latest") def latest_block(chain: str): blk = pool(chain).call("eth_getBlockByNumber", ["latest", False]) return block_summary(blk) @app.get("/v1/{chain}/block/{number}") def block_by_number(chain: str, number: int): p = pool(chain) blk = p.call("eth_getBlockByNumber", [hex(number), False]) if blk is None: raise HTTPException(404, "block not found") out = block_summary(blk) # stablecoin transfers inside this block, decoded live addrs = [t["address"] for t in TOKENS.get(chain, [])] if addrs: logs = p.call( "eth_getLogs", [{"fromBlock": hex(number), "toBlock": hex(number), "address": addrs, "topics": [TRANSFER_TOPIC]}], ) transfers = [] for l in logs or []: meta = token_meta(chain, l["address"]) if meta and len(l.get("topics", [])) == 3: transfers.append( with_value(decode_transfer(chain, l, meta, out["timestamp"])) ) out["stablecoin_transfers"] = transfers return out @app.get("/v1/{chain}/tx/{tx_hash}") def tx(chain: str, tx_hash: str): if chain in CHAINS and CHAINS[chain].get("family", "evm") != "evm": # non-EVM: serve what the index knows about this tx conn = dbc() try: rows = conn.execute( "SELECT * FROM transfers WHERE chain = ? AND tx_hash IN (?, ?) " "ORDER BY log_index", (chain, tx_hash, tx_hash.lower()), ).fetchall() finally: conn.close() if not rows: raise HTTPException(404, "transaction not in index") first = rows[0] return {"chain": chain, "hash": first["tx_hash"], "block": first["block"], "timestamp": first["timestamp"], "source": "index", "stablecoin_transfers": [with_value(r) for r in rows]} p = pool(chain) t = p.call("eth_getTransactionByHash", [tx_hash]) if t is None: raise HTTPException(404, "transaction not found") receipt = p.call("eth_getTransactionReceipt", [tx_hash]) transfers = [] for l in (receipt or {}).get("logs", []): meta = token_meta(chain, l["address"]) if ( meta and l.get("topics") and l["topics"][0].lower() == TRANSFER_TOPIC and len(l["topics"]) == 3 ): transfers.append(with_value(decode_transfer(chain, l, meta, None))) return { "chain": chain, "hash": t["hash"], "block": int(t["blockNumber"], 16) if t.get("blockNumber") else None, "from": t["from"], "to": t.get("to"), "native_value_wei": str(int(t.get("value", "0x0"), 16)), "status": int(receipt["status"], 16) if receipt and receipt.get("status") else None, "gas_used": int(receipt["gasUsed"], 16) if receipt else None, "stablecoin_transfers": transfers, } # -- indexed queries (from SQLite) --------------------------------------- @app.get("/v1/{chain}/token/{address}/transfers") def token_transfers( chain: str, address: str, from_addr: str | None = Query(None, alias="from"), to_addr: str | None = Query(None, alias="to"), min_amount: float | None = None, since: int | None = Query(None, description="unix timestamp lower bound"), before_block: int | None = Query(None, description="pagination cursor: only blocks below this"), limit: int = Query(50, le=500), ): if chain not in CHAINS: raise HTTPException(404, f"unknown chain '{chain}'") # EVM identifiers are case-insensitive hex (stored lowercase); Base58 / # coin-type identifiers on other chains are case-sensitive — keep verbatim evm = CHAINS[chain].get("family", "evm") == "evm" norm = (lambda s: s.lower()) if evm else (lambda s: s) q = SELECT_T + " WHERE chain = ? AND token = ?" args = [chain, norm(address)] if from_addr: q += ' AND "from" = ?' args.append(norm(from_addr)) if to_addr: q += ' AND "to" = ?' args.append(norm(to_addr)) if min_amount is not None: meta = token_meta(chain, address) decimals = meta["decimals"] if meta else 6 # CAST for comparison: amounts are TEXT (uint256 overflows int64); # DOUBLE PRECISION works on PG and maps to REAL affinity on SQLite q += " AND CAST(amount AS DOUBLE PRECISION) >= ?" args.append(min_amount * 10**decimals) if since is not None: q += " AND timestamp >= ?" args.append(since) if before_block is not None: q += " AND block < ?" args.append(before_block) q += " ORDER BY block DESC, log_index DESC LIMIT ?" args.append(limit) conn = dbc() try: return [with_value(r) for r in conn.execute(q, args)] finally: conn.close() @app.get("/v1/address/{addr}/transfers") def address_transfers(addr: str, limit: int = Query(50, le=500)): """Cross-chain: every indexed stablecoin transfer touching this address.""" conn = dbc() try: # match both verbatim (Base58/native formats) and lowercased (EVM hex) rows = conn.execute( SELECT_T + ' WHERE "from" IN (?, ?) OR "to" IN (?, ?) ' "ORDER BY timestamp DESC LIMIT ?", (addr, addr.lower(), addr, addr.lower(), limit), ) return [with_value(r) for r in rows] finally: conn.close() @app.get("/v1/stablecoins/volume") @cached(ttl=30) def volume(token: str = "USDT", window: str = "24h"): seconds = parse_window(window) since = int(time.time()) - seconds conn = dbc() try: # serve from the StatsWorker's precomputed table when fresh — the # COUNT(DISTINCT) work does not scale on the request path agg = conn.execute( "SELECT * FROM agg_volume WHERE window = ? AND symbol = ?", (window, token.upper()), ).fetchall() if agg and int(time.time()) - agg[0]["updated"] < 900: per_chain = { r["chain"]: {"transfers": r["transfers"], "volume": r["volume"], "senders": r["senders"], "receivers": r["receivers"]} for r in agg } return { "token": token.upper(), "window": window, "since_unix": since, "as_of": agg[0]["updated"], "total_volume": round(sum(r["volume"] for r in agg), 2), "chains": per_chain, "note": "volume only covers blocks this instance has indexed", } # fallback (fresh deploy / unusual window): live scan WITHOUT the # distinct-address counts, which are what makes this query heavy rows = conn.execute( "SELECT chain, decimals, COUNT(*) AS transfers, " "SUM(CAST(amount AS DOUBLE PRECISION)) AS raw_sum " "FROM transfers WHERE symbol = ? AND timestamp >= ? " "GROUP BY chain, decimals", (token.upper(), since), ).fetchall() price = price_of(conn, token.upper()) finally: conn.close() per_chain = {} total = 0.0 for r in rows: vol = (r["raw_sum"] or 0) / 10 ** r["decimals"] * price entry = per_chain.setdefault( r["chain"], {"transfers": 0, "volume": 0.0, "senders": None, "receivers": None} ) entry["transfers"] += r["transfers"] entry["volume"] = round(entry["volume"] + vol, 2) total += vol return { "token": token.upper(), "window": window, "since_unix": since, "total_volume": round(total, 2), "chains": per_chain, "note": "volume only covers blocks this instance has indexed", } @app.get("/v1/tokens") @cached(ttl=300) def tokens(): """Every configured stablecoin, grouped by canonical symbol.""" out = {} for chain, toks in TOKENS.items(): for t in toks: entry = out.setdefault(t["symbol"], { "category": t.get("category", "stablecoin"), "chains": {}}) entry["chains"][chain] = { "id": t.get("address") or t.get("id"), "decimals": t.get("decimals"), "native": bool(t.get("native")), "discontinued": bool(t.get("discontinued", False)), "onchain_symbol": t.get("onchain_symbol", t["symbol"]), } return out @app.get("/v1/prices") @cached(ttl=30) def prices(): """Latest USD prices used for valuation (stablecoins seeded at 1.0).""" conn = dbc() try: return { r["symbol"]: {"usd": r["usd"], "updated": r["updated"]} for r in conn.execute("SELECT * FROM prices ORDER BY symbol") } finally: conn.close() @app.get("/v1/stablecoins/supply") @cached(ttl=60) def supply(token: str | None = None): """Latest on-chain supply snapshot per (chain, token), + totals.""" conn = dbc() try: q = ( "SELECT s.chain, s.token, s.symbol, s.supply, s.decimals, s.timestamp " "FROM supply_snapshots s JOIN (" " SELECT chain, token, MAX(timestamp) AS mt FROM supply_snapshots " " GROUP BY chain, token) m " "ON s.chain = m.chain AND s.token = m.token AND s.timestamp = m.mt" ) args = () if token: q += " WHERE s.symbol = ?" args = (token.upper(),) rows = conn.execute(q, args).fetchall() prices = {r["symbol"]: r["usd"] for r in conn.execute("SELECT symbol, usd FROM prices")} finally: conn.close() per_symbol = {} for r in rows: supply_h = float(format_amount(r["supply"], r["decimals"] or 6)) price = prices.get(r["symbol"]) or 1.0 sym = per_symbol.setdefault(r["symbol"], {"total": 0.0, "total_usd": 0.0, "chains": {}}) sym["chains"][r["chain"]] = { "supply": supply_h, "usd": round(supply_h * price, 2), "raw": r["supply"], "decimals": r["decimals"], "as_of": r["timestamp"], } sym["total"] = round(sym["total"] + supply_h, 2) sym["total_usd"] = round(sym["total_usd"] + supply_h * price, 2) return per_symbol @app.get("/v1/stablecoins/whales") def whales(token: str | None = None, window: str = "24h", min_usd: float | None = None, limit: int = Query(50, le=500)): """Largest transfers in the window (face value ≈ USD for stablecoins).""" threshold = min_usd if min_usd is not None else DEFAULT_WHALE_USD since = int(time.time()) - parse_window(window) # whale_events is maintained incrementally by StatsWorker (floor $100K) — # reading it is O(window) instead of a full transfers scan q = "SELECT * FROM whale_events WHERE timestamp >= ? AND usd >= ?" args = [since, threshold] if token: q += " AND symbol = ?" args.append(token.upper()) q += " ORDER BY usd DESC LIMIT ?" args.append(limit) conn = dbc() try: return [with_value(r) for r in conn.execute(q, args)] finally: conn.close() @app.get("/v1/stablecoins/mints-burns") def mints_burns(token: str | None = None, window: str = "24h", limit: int = Query(100, le=500)): """Issuance events: transfers from the zero address are mints, to it are burns (EVM + Tron; Solana mints don't transit an address — see docs).""" zeros = tuple(ZERO_ADDRESSES.values()) since = int(time.time()) - parse_window(window) zp = ", ".join("?" for _ in zeros) q = (SELECT_T + ' WHERE timestamp >= ? ' f'AND ("from" IN ({zp}) OR "to" IN ({zp}))') args = [since, *zeros, *zeros] if token: q += " AND transfers.symbol = ?" args.append(token.upper()) q += " ORDER BY timestamp DESC LIMIT ?" args.append(limit) conn = dbc() try: out = [] for r in conn.execute(q, args): d = with_value(r) d["direction"] = "mint" if r["from"] in zeros else "burn" out.append(d) return out finally: conn.close() @app.get("/v1/stablecoins/volume/series") @cached(ttl=60) def volume_series(token: str = "USDT", window: str = "24h", interval: str = "1h", chain: str | None = None): """Bucketed transfer volume over time — per chain, for line charts. Served from StatsWorker's precomputed buckets for the standard windows.""" since = int(time.time()) - parse_window(window) step = parse_window(interval) conn0 = dbc() try: agg = conn0.execute( "SELECT * FROM agg_series WHERE window = ? AND symbol = ? " + ("AND chain = ? " if chain else "") + "ORDER BY t", (window, token.upper(), *([chain] if chain else [])), ).fetchall() finally: conn0.close() if agg and int(time.time()) - agg[0]["updated"] < 900: out = {} for r in agg: out.setdefault(r["chain"], []).append( {"t": r["t"], "volume": r["volume"], "transfers": r["transfers"]}) worker_step = {"1h": 300, "24h": 3600, "7d": 21600}.get(window, step) return {"token": token.upper(), "since": since, "interval_s": worker_step, "as_of": agg[0]["updated"], "chains": out} q = (f"SELECT (timestamp / {step}) * {step} AS t, chain, decimals, " "COUNT(*) AS transfers, SUM(CAST(amount AS DOUBLE PRECISION)) AS raw_sum " "FROM transfers WHERE symbol = ? AND timestamp >= ?") args = [token.upper(), since] if chain: q += " AND chain = ?" args.append(chain) q += " GROUP BY t, chain, decimals ORDER BY t" conn = dbc() try: rows = conn.execute(q, args).fetchall() price = price_of(conn, token.upper()) finally: conn.close() out = {} for r in rows: c = out.setdefault(r["chain"], {}) b = c.setdefault(r["t"], {"volume": 0.0, "transfers": 0}) b["volume"] = round(b["volume"] + (r["raw_sum"] or 0) / 10 ** r["decimals"] * price, 2) b["transfers"] += r["transfers"] return { "token": token.upper(), "since": since, "interval_s": step, "chains": { c: [{"t": t, **v} for t, v in sorted(buckets.items())] for c, buckets in out.items() }, } @app.get("/v1/stablecoins/supply/series") @cached(ttl=120) def supply_series(token: str = "USDT", window: str = "7d"): """Supply snapshots over time per chain (hourly cadence).""" since = int(time.time()) - parse_window(window) conn = dbc() try: rows = conn.execute( "SELECT chain, supply, decimals, timestamp FROM supply_snapshots " "WHERE symbol = ? AND timestamp >= ? ORDER BY timestamp", (token.upper(), since), ).fetchall() finally: conn.close() out = {} for r in rows: out.setdefault(r["chain"], []).append( {"t": r["timestamp"], "supply": float(format_amount(r["supply"], r["decimals"] or 6))} ) return {"token": token.upper(), "since": since, "chains": out} @app.get("/v1/stablecoins/flows") @cached(ttl=60) def flows(token: str | None = None, window: str = "7d", interval: str = "1d"): """Net issuance over time: mints (from zero addr) minus burns (to zero).""" zeros = tuple(ZERO_ADDRESSES.values()) since = int(time.time()) - parse_window(window) step = parse_window(interval) zp = ", ".join("?" for _ in zeros) q = (f"SELECT (timestamp / {step}) * {step} AS t, transfers.symbol AS symbol, " f'SUM(CASE WHEN "from" IN ({zp}) THEN {USD_PRICED} ELSE 0 END) AS minted, ' f'SUM(CASE WHEN "to" IN ({zp}) THEN {USD_PRICED} ELSE 0 END) AS burned ' f'FROM transfers {PRICE_JOIN} ' f'WHERE timestamp >= ? AND ("from" IN ({zp}) OR "to" IN ({zp}))') args = [*zeros, *zeros, since, *zeros, *zeros] if token: q += " AND transfers.symbol = ?" args.append(token.upper()) q += " GROUP BY t, transfers.symbol ORDER BY t" conn = dbc() try: rows = conn.execute(q, args).fetchall() finally: conn.close() out = {} for r in rows: out.setdefault(r["symbol"], []).append({ "t": r["t"], "minted": round(r["minted"] or 0, 2), "burned": round(r["burned"] or 0, 2), "net": round((r["minted"] or 0) - (r["burned"] or 0), 2), }) return {"since": since, "interval_s": step, "tokens": out} @app.get("/v1/{chain}/transfers") def chain_transfers(chain: str, symbol: str | None = None, before_block: int | None = None, min_amount: float | None = None, limit: int = Query(50, le=500)): """Recent indexed transfers on one chain, all tokens (paginated).""" if chain not in CHAINS: raise HTTPException(404, f"unknown chain '{chain}'") q = SELECT_T + " WHERE chain = ?" args = [chain] if symbol: q += " AND transfers.symbol = ?" args.append(symbol.upper()) if min_amount is not None: q += f" AND {USD_PRICED} >= ?" args.append(min_amount) if before_block is not None: q += " AND block < ?" args.append(before_block) q += " ORDER BY block DESC, log_index DESC LIMIT ?" args.append(limit) conn = dbc() try: return [with_value(r) for r in conn.execute(q, args)] finally: conn.close() @app.get("/v1/{chain}/summary") @cached(ttl=30) def chain_summary(chain: str): if chain not in CHAINS: raise HTTPException(404, f"unknown chain '{chain}'") since = int(time.time()) - 86400 conn = dbc() try: vol = conn.execute( f"SELECT transfers.symbol AS symbol, COUNT(*) AS transfers, " f"SUM({USD_PRICED}) AS volume " f"FROM transfers {PRICE_JOIN} " "WHERE chain = ? AND timestamp >= ? GROUP BY transfers.symbol " "ORDER BY volume DESC", (chain, since), ).fetchall() cur = conn.execute("SELECT * FROM cursors WHERE chain = ?", (chain,)).fetchone() total = conn.execute( "SELECT COUNT(*) AS n FROM transfers WHERE chain = ?", (chain,) ).fetchone() finally: conn.close() cfg = CHAINS[chain] return { "chain": chain, "family": cfg.get("family", "evm"), "chain_id": cfg.get("chain_id"), "block_time_s": cfg.get("block_time"), "tokens": TOKENS.get(chain, []), "indexed_transfers": total["n"] if total else 0, "volume_24h": [ {"symbol": r["symbol"], "transfers": r["transfers"], "volume": round(r["volume"] or 0, 2)} for r in vol ], "cursor": dict(cur) if cur else None, } @app.get("/v1/search") def search(q: str): """Classify a query: tx hash / address / token symbol / chain name. Checks the index first (so we can say WHICH chain a tx lives on).""" s = q.strip() if not s: raise HTTPException(400, "empty query") if s.upper() in {t["symbol"] for toks in TOKENS.values() for t in toks}: return {"type": "token", "symbol": s.upper()} if s.lower() in CHAINS: return {"type": "chain", "chain": s.lower()} conn = dbc() try: needles = (s, s.lower()) r = conn.execute( "SELECT chain, tx_hash FROM transfers WHERE tx_hash IN (?, ?) LIMIT 1", needles ).fetchone() if r: return {"type": "tx", "chain": r["chain"], "hash": r["tx_hash"]} r = conn.execute( 'SELECT COUNT(*) AS n FROM transfers WHERE "from" IN (?, ?) OR "to" IN (?, ?)', needles + needles, ).fetchone() if r and r["n"]: return {"type": "address", "address": s, "indexed_transfers": r["n"]} finally: conn.close() # shape-based fallback for things we haven't indexed (yet) import re if re.fullmatch(r"(0x)?[0-9a-fA-F]{64}", s) or re.fullmatch(r"[1-9A-HJ-NP-Za-km-z]{80,90}", s): return {"type": "tx", "chain": None, "hash": s} if re.fullmatch(r"0x[0-9a-fA-F]{40}", s) or re.fullmatch(r"T[1-9A-HJ-NP-Za-km-z]{33}", s) \ or re.fullmatch(r"[1-9A-HJ-NP-Za-km-z]{32,44}", s): return {"type": "address", "address": s, "indexed_transfers": 0} return {"type": "unknown"} # -- live stream --------------------------------------------------------- def _poll_transfers(since_ts): conn = dbc() try: rows = conn.execute( SELECT_T + " WHERE timestamp >= ? ORDER BY timestamp ASC LIMIT 1000", (since_ts,) ).fetchall() return [with_value(r) for r in rows] finally: conn.close() @app.websocket("/v1/stream/transfers") async def stream_transfers(ws: WebSocket): """Pushes new transfers every ~2s. Optional query params: token=USDT chain=ethereum min_usd=1000""" await ws.accept() token = (ws.query_params.get("token") or "").upper() or None chain = ws.query_params.get("chain") min_usd = float(ws.query_params.get("min_usd") or 0) last_ts = int(time.time()) - 2 seen = set() try: while True: rows = await asyncio.to_thread(_poll_transfers, last_ts) batch = [] for r in rows: key = (r["chain"], r["tx_hash"], r["log_index"]) if key in seen: continue if token and r["symbol"] != token: seen.add(key) continue if chain and r["chain"] != chain: seen.add(key) continue if min_usd and r.get("usd", float(r["value"])) < min_usd: seen.add(key) continue seen.add(key) batch.append(r) if batch: await ws.send_json(batch) if rows: new_last = max(r["timestamp"] or last_ts for r in rows) if new_last > last_ts: last_ts = new_last # only remember keys that can still reappear in queries seen = {k for k in seen} if len(seen) < 50_000 else set() await asyncio.sleep(2) except WebSocketDisconnect: pass # -- dashboard (must mount last: "/" catches everything below the API routes) -- app.mount("/", StaticFiles(directory=str(pathlib.Path(__file__).parent.parent / "ui"), html=True), name="ui")