SPB Git

spb/coinexplorer Public MIT

Self-hosted, zero-API-key explorer for stablecoins and major crypto.

Python 60.3% HTML 23.6% JavaScript 8.1% CSS 6.8% SQL 1%
29.9 KB · 833 lines python
Raw Blame History
1# Author: Simon-Pierre Boucher2# Mail: contact@spboucher.ai3"""Stablecoin explorer API — FastAPI over the indexer's SQLite DB,4plus live pass-through reads (blocks, txs) straight from free RPCs.56    uvicorn api.main:app --port 80807"""89import asyncio10import functools11import os12import time1314import pathlib1516from fastapi import FastAPI, HTTPException, Query, WebSocket, WebSocketDisconnect17from fastapi.staticfiles import StaticFiles1819from indexer import config, db20from indexer.decode import TRANSFER_TOPIC, decode_transfer, format_amount21from indexer.enrich import ZERO_ADDRESSES22from indexer.rpc import RpcPool2324DB_PATH = config.db_path()25CHAINS = config.load_chains()26TOKENS = config.load_tokens()2728app = FastAPI(29    title="coinexplorer API",30    description="Self-hosted, free-RPC-only explorer for stablecoins and major crypto.",31    version="0.3.0",32)333435@app.exception_handler(Exception)36async def unhandled_exception(request, exc):37    # never leak a stack trace; always answer JSON38    from fastapi.responses import JSONResponse39    return JSONResponse(status_code=500, content={"error": str(exc)[:200]})4041_pools = {}424344def pool(chain):45    if chain not in CHAINS:46        raise HTTPException(404, f"unknown chain '{chain}'")47    if CHAINS[chain].get("family", "evm") != "evm":48        raise HTTPException(501, f"live reads for '{chain}' land with its adapter (Deliverable 4)")49    if chain not in _pools:50        _pools[chain] = RpcPool(CHAINS[chain]["rpcs"])51    return _pools[chain]525354def dbc():55    return db.connect(DB_PATH, readonly=True)565758# stablecoins trade ~$1, so face value ≈ USD; CASE avoids relying on POW()59# (not compiled into every SQLite build)60USD_EXPR = (61    "CAST(amount AS DOUBLE PRECISION) / (CASE decimals "62    + " ".join(f"WHEN {d} THEN 1e{d}" for d in range(19))63    + " ELSE 1e6 END)"64)6566# price-aware USD: stablecoins are seeded at 1.0 by the PriceWorker; crypto67# assets get their CoinGecko price; unknown symbols fall back to face value68PRICE_JOIN = "LEFT JOIN prices ON prices.symbol = transfers.symbol"69USD_PRICED = f"(({USD_EXPR}) * COALESCE(prices.usd, 1.0))"70SELECT_T = f"SELECT transfers.*, COALESCE(prices.usd, 1.0) AS usd_rate FROM transfers {PRICE_JOIN}"7172DEFAULT_WHALE_USD = float(os.environ.get("WHALE_MIN_USD", 1_000_000))737475def price_of(conn, symbol):76    r = conn.execute("SELECT usd FROM prices WHERE symbol = ?", (symbol,)).fetchone()77    return r["usd"] if r and r["usd"] is not None else 1.07879_cache = {}808182def cached(ttl):83    """Tiny in-process TTL cache for aggregate endpoints."""84    def deco(fn):85        @functools.wraps(fn)86        def wrap(*a, **k):87            key = (fn.__name__, a, tuple(sorted(k.items())))88            hit = _cache.get(key)89            now = time.time()90            if hit and now - hit[0] < ttl:91                return hit[1]92            val = fn(*a, **k)93            _cache[key] = (now, val)94            return val95        return wrap96    return deco979899def parse_window(window):100    units = {"m": 60, "h": 3600, "d": 86400}101    try:102        return int(window[:-1]) * units[window[-1]]103    except (KeyError, ValueError, IndexError):104        raise HTTPException(400, "window must look like 15m, 24h or 7d")105106107def with_value(row):108    d = dict(row)109    rate = d.pop("usd_rate", None)110    if d.get("amount") is not None and d.get("decimals") is not None:111        d["value"] = format_amount(d["amount"], d["decimals"])112        if rate is not None:113            d["usd"] = round(float(d["value"]) * rate, 2)114    return d115116117def token_meta(chain, address):118    for t in TOKENS.get(chain, []):119        if t.get("address", "").lower() == address.lower():120            return t121    return None122123124def block_summary(blk):125    return {126        "number": int(blk["number"], 16),127        "hash": blk["hash"],128        "parent_hash": blk["parentHash"],129        "timestamp": int(blk["timestamp"], 16),130        "tx_count": len(blk.get("transactions", [])),131        "gas_used": int(blk.get("gasUsed", "0x0"), 16),132    }133134135# -- meta ---------------------------------------------------------------136137138@app.get("/v1/chains")139def chains():140    return {141        chain: {142            "family": cfg.get("family", "evm"),143            "chain_id": cfg.get("chain_id"),144            "block_time_s": cfg.get("block_time"),145            "confirmations": cfg.get("confirmations"),146            "tokens": TOKENS.get(chain, []),147        }148        for chain, cfg in CHAINS.items()149    }150151152@app.get("/v1/status")153def status():154    conn = dbc()155    try:156        cursors = {157            r["chain"]: {158                "last_block": r["last_block"],159                "last_hash": r["last_hash"],160                "backfill_block": r["backfill_block"],161                "head_block": r["head_block"],162                "lag": (r["head_block"] - r["last_block"])163                       if r["head_block"] is not None else None,164            }165            for r in conn.execute("SELECT * FROM cursors")166            if not r["chain"].startswith("_")  # internal scan cursors167        }168        counts = {169            r["chain"]: r["n"]170            for r in conn.execute(171                "SELECT chain, COUNT(*) AS n FROM transfers GROUP BY chain"172            )173        }174        return {"cursors": cursors, "indexed_transfers": counts}175    finally:176        conn.close()177178179@app.get("/v1/rpc/health")180def rpc_health():181    """Per-endpoint health as last reported by the indexer workers:182    score = success-rate EMA, cooldown_s > 0 = currently benched."""183    conn = dbc()184    try:185        rows = conn.execute(186            "SELECT * FROM rpc_health ORDER BY chain, score DESC"187        ).fetchall()188    finally:189        conn.close()190    out = {}191    for r in rows:192        out.setdefault(r["chain"], []).append(193            {k: r[k] for k in ("url", "score", "ok", "fail", "latency_ms", "cooldown_s", "updated")}194        )195    return out196197198@app.get("/metrics")199def metrics():200    """Prometheus text exposition: per-chain progress/lag + RPC endpoint health.201    Free-RPC degradation shows up here as rising lag and falling rpc scores."""202    from fastapi.responses import PlainTextResponse203    conn = dbc()204    try:205        lines = [206            "# HELP explorer_transfers_total Indexed transfers per chain",207            "# TYPE explorer_transfers_total counter",208        ]209        for r in conn.execute("SELECT chain, COUNT(*) n FROM transfers GROUP BY chain"):210            lines.append(f'explorer_transfers_total{{chain="{r["chain"]}"}} {r["n"]}')211        lines += [212            "# HELP explorer_cursor_block Last indexed block per chain",213            "# TYPE explorer_cursor_block gauge",214            "# HELP explorer_lag_blocks Head minus cursor per chain",215            "# TYPE explorer_lag_blocks gauge",216        ]217        for r in conn.execute("SELECT * FROM cursors"):218            lines.append(f'explorer_cursor_block{{chain="{r["chain"]}"}} {r["last_block"]}')219            if r["head_block"] is not None:220                lines.append(f'explorer_lag_blocks{{chain="{r["chain"]}"}} '221                             f'{max(0, r["head_block"] - r["last_block"])}')222        lines += [223            "# HELP explorer_rpc_score Endpoint health score (success-rate EMA, 0-1)",224            "# TYPE explorer_rpc_score gauge",225            "# HELP explorer_rpc_calls_total Successful/failed calls per endpoint",226            "# TYPE explorer_rpc_calls_total counter",227        ]228        for r in conn.execute("SELECT * FROM rpc_health"):229            lbl = f'chain="{r["chain"]}",url="{r["url"]}"'230            lines.append(f"explorer_rpc_score{{{lbl}}} {r['score']}")231            lines.append(f'explorer_rpc_calls_total{{{lbl},result="ok"}} {r["ok"]}')232            lines.append(f'explorer_rpc_calls_total{{{lbl},result="fail"}} {r["fail"]}')233        return PlainTextResponse("\n".join(lines) + "\n")234    finally:235        conn.close()236237238# -- live chain reads (straight off the free RPCs) -----------------------239240241@app.get("/v1/{chain}/block/latest")242def latest_block(chain: str):243    blk = pool(chain).call("eth_getBlockByNumber", ["latest", False])244    return block_summary(blk)245246247@app.get("/v1/{chain}/block/{number}")248def block_by_number(chain: str, number: int):249    p = pool(chain)250    blk = p.call("eth_getBlockByNumber", [hex(number), False])251    if blk is None:252        raise HTTPException(404, "block not found")253    out = block_summary(blk)254    # stablecoin transfers inside this block, decoded live255    addrs = [t["address"] for t in TOKENS.get(chain, [])]256    if addrs:257        logs = p.call(258            "eth_getLogs",259            [{"fromBlock": hex(number), "toBlock": hex(number),260              "address": addrs, "topics": [TRANSFER_TOPIC]}],261        )262        transfers = []263        for l in logs or []:264            meta = token_meta(chain, l["address"])265            if meta and len(l.get("topics", [])) == 3:266                transfers.append(267                    with_value(decode_transfer(chain, l, meta, out["timestamp"]))268                )269        out["stablecoin_transfers"] = transfers270    return out271272273@app.get("/v1/{chain}/tx/{tx_hash}")274def tx(chain: str, tx_hash: str):275    if chain in CHAINS and CHAINS[chain].get("family", "evm") != "evm":276        # non-EVM: serve what the index knows about this tx277        conn = dbc()278        try:279            rows = conn.execute(280                "SELECT * FROM transfers WHERE chain = ? AND tx_hash IN (?, ?) "281                "ORDER BY log_index", (chain, tx_hash, tx_hash.lower()),282            ).fetchall()283        finally:284            conn.close()285        if not rows:286            raise HTTPException(404, "transaction not in index")287        first = rows[0]288        return {"chain": chain, "hash": first["tx_hash"], "block": first["block"],289                "timestamp": first["timestamp"], "source": "index",290                "stablecoin_transfers": [with_value(r) for r in rows]}291    p = pool(chain)292    t = p.call("eth_getTransactionByHash", [tx_hash])293    if t is None:294        raise HTTPException(404, "transaction not found")295    receipt = p.call("eth_getTransactionReceipt", [tx_hash])296    transfers = []297    for l in (receipt or {}).get("logs", []):298        meta = token_meta(chain, l["address"])299        if (300            meta301            and l.get("topics")302            and l["topics"][0].lower() == TRANSFER_TOPIC303            and len(l["topics"]) == 3304        ):305            transfers.append(with_value(decode_transfer(chain, l, meta, None)))306    return {307        "chain": chain,308        "hash": t["hash"],309        "block": int(t["blockNumber"], 16) if t.get("blockNumber") else None,310        "from": t["from"],311        "to": t.get("to"),312        "native_value_wei": str(int(t.get("value", "0x0"), 16)),313        "status": int(receipt["status"], 16) if receipt and receipt.get("status") else None,314        "gas_used": int(receipt["gasUsed"], 16) if receipt else None,315        "stablecoin_transfers": transfers,316    }317318319# -- indexed queries (from SQLite) ---------------------------------------320321322@app.get("/v1/{chain}/token/{address}/transfers")323def token_transfers(324    chain: str,325    address: str,326    from_addr: str | None = Query(None, alias="from"),327    to_addr: str | None = Query(None, alias="to"),328    min_amount: float | None = None,329    since: int | None = Query(None, description="unix timestamp lower bound"),330    before_block: int | None = Query(None, description="pagination cursor: only blocks below this"),331    limit: int = Query(50, le=500),332):333    if chain not in CHAINS:334        raise HTTPException(404, f"unknown chain '{chain}'")335    # EVM identifiers are case-insensitive hex (stored lowercase); Base58 /336    # coin-type identifiers on other chains are case-sensitive — keep verbatim337    evm = CHAINS[chain].get("family", "evm") == "evm"338    norm = (lambda s: s.lower()) if evm else (lambda s: s)339    q = SELECT_T + " WHERE chain = ? AND token = ?"340    args = [chain, norm(address)]341    if from_addr:342        q += ' AND "from" = ?'343        args.append(norm(from_addr))344    if to_addr:345        q += ' AND "to" = ?'346        args.append(norm(to_addr))347    if min_amount is not None:348        meta = token_meta(chain, address)349        decimals = meta["decimals"] if meta else 6350        # CAST for comparison: amounts are TEXT (uint256 overflows int64);351        # DOUBLE PRECISION works on PG and maps to REAL affinity on SQLite352        q += " AND CAST(amount AS DOUBLE PRECISION) >= ?"353        args.append(min_amount * 10**decimals)354    if since is not None:355        q += " AND timestamp >= ?"356        args.append(since)357    if before_block is not None:358        q += " AND block < ?"359        args.append(before_block)360    q += " ORDER BY block DESC, log_index DESC LIMIT ?"361    args.append(limit)362    conn = dbc()363    try:364        return [with_value(r) for r in conn.execute(q, args)]365    finally:366        conn.close()367368369@app.get("/v1/address/{addr}/transfers")370def address_transfers(addr: str, limit: int = Query(50, le=500)):371    """Cross-chain: every indexed stablecoin transfer touching this address."""372    conn = dbc()373    try:374        # match both verbatim (Base58/native formats) and lowercased (EVM hex)375        rows = conn.execute(376            SELECT_T + ' WHERE "from" IN (?, ?) OR "to" IN (?, ?) '377            "ORDER BY timestamp DESC LIMIT ?",378            (addr, addr.lower(), addr, addr.lower(), limit),379        )380        return [with_value(r) for r in rows]381    finally:382        conn.close()383384385@app.get("/v1/stablecoins/volume")386@cached(ttl=30)387def volume(token: str = "USDT", window: str = "24h"):388    seconds = parse_window(window)389    since = int(time.time()) - seconds390    conn = dbc()391    try:392        # serve from the StatsWorker's precomputed table when fresh — the393        # COUNT(DISTINCT) work does not scale on the request path394        agg = conn.execute(395            "SELECT * FROM agg_volume WHERE window = ? AND symbol = ?",396            (window, token.upper()),397        ).fetchall()398        if agg and int(time.time()) - agg[0]["updated"] < 900:399            per_chain = {400                r["chain"]: {"transfers": r["transfers"], "volume": r["volume"],401                             "senders": r["senders"], "receivers": r["receivers"]}402                for r in agg403            }404            return {405                "token": token.upper(), "window": window, "since_unix": since,406                "as_of": agg[0]["updated"],407                "total_volume": round(sum(r["volume"] for r in agg), 2),408                "chains": per_chain,409                "note": "volume only covers blocks this instance has indexed",410            }411        # fallback (fresh deploy / unusual window): live scan WITHOUT the412        # distinct-address counts, which are what makes this query heavy413        rows = conn.execute(414            "SELECT chain, decimals, COUNT(*) AS transfers, "415            "SUM(CAST(amount AS DOUBLE PRECISION)) AS raw_sum "416            "FROM transfers WHERE symbol = ? AND timestamp >= ? "417            "GROUP BY chain, decimals",418            (token.upper(), since),419        ).fetchall()420        price = price_of(conn, token.upper())421    finally:422        conn.close()423    per_chain = {}424    total = 0.0425    for r in rows:426        vol = (r["raw_sum"] or 0) / 10 ** r["decimals"] * price427        entry = per_chain.setdefault(428            r["chain"], {"transfers": 0, "volume": 0.0, "senders": None, "receivers": None}429        )430        entry["transfers"] += r["transfers"]431        entry["volume"] = round(entry["volume"] + vol, 2)432        total += vol433    return {434        "token": token.upper(),435        "window": window,436        "since_unix": since,437        "total_volume": round(total, 2),438        "chains": per_chain,439        "note": "volume only covers blocks this instance has indexed",440    }441442443@app.get("/v1/tokens")444@cached(ttl=300)445def tokens():446    """Every configured stablecoin, grouped by canonical symbol."""447    out = {}448    for chain, toks in TOKENS.items():449        for t in toks:450            entry = out.setdefault(t["symbol"], {451                "category": t.get("category", "stablecoin"), "chains": {}})452            entry["chains"][chain] = {453                "id": t.get("address") or t.get("id"),454                "decimals": t.get("decimals"),455                "native": bool(t.get("native")),456                "discontinued": bool(t.get("discontinued", False)),457                "onchain_symbol": t.get("onchain_symbol", t["symbol"]),458            }459    return out460461462@app.get("/v1/prices")463@cached(ttl=30)464def prices():465    """Latest USD prices used for valuation (stablecoins seeded at 1.0)."""466    conn = dbc()467    try:468        return {469            r["symbol"]: {"usd": r["usd"], "updated": r["updated"]}470            for r in conn.execute("SELECT * FROM prices ORDER BY symbol")471        }472    finally:473        conn.close()474475476@app.get("/v1/stablecoins/supply")477@cached(ttl=60)478def supply(token: str | None = None):479    """Latest on-chain supply snapshot per (chain, token), + totals."""480    conn = dbc()481    try:482        q = (483            "SELECT s.chain, s.token, s.symbol, s.supply, s.decimals, s.timestamp "484            "FROM supply_snapshots s JOIN ("485            "  SELECT chain, token, MAX(timestamp) AS mt FROM supply_snapshots "486            "  GROUP BY chain, token) m "487            "ON s.chain = m.chain AND s.token = m.token AND s.timestamp = m.mt"488        )489        args = ()490        if token:491            q += " WHERE s.symbol = ?"492            args = (token.upper(),)493        rows = conn.execute(q, args).fetchall()494        prices = {r["symbol"]: r["usd"] for r in conn.execute("SELECT symbol, usd FROM prices")}495    finally:496        conn.close()497    per_symbol = {}498    for r in rows:499        supply_h = float(format_amount(r["supply"], r["decimals"] or 6))500        price = prices.get(r["symbol"]) or 1.0501        sym = per_symbol.setdefault(r["symbol"], {"total": 0.0, "total_usd": 0.0, "chains": {}})502        sym["chains"][r["chain"]] = {503            "supply": supply_h, "usd": round(supply_h * price, 2), "raw": r["supply"],504            "decimals": r["decimals"], "as_of": r["timestamp"],505        }506        sym["total"] = round(sym["total"] + supply_h, 2)507        sym["total_usd"] = round(sym["total_usd"] + supply_h * price, 2)508    return per_symbol509510511@app.get("/v1/stablecoins/whales")512def whales(token: str | None = None, window: str = "24h",513           min_usd: float | None = None, limit: int = Query(50, le=500)):514    """Largest transfers in the window (face value ≈ USD for stablecoins)."""515    threshold = min_usd if min_usd is not None else DEFAULT_WHALE_USD516    since = int(time.time()) - parse_window(window)517    # whale_events is maintained incrementally by StatsWorker (floor $100K) —518    # reading it is O(window) instead of a full transfers scan519    q = "SELECT * FROM whale_events WHERE timestamp >= ? AND usd >= ?"520    args = [since, threshold]521    if token:522        q += " AND symbol = ?"523        args.append(token.upper())524    q += " ORDER BY usd DESC LIMIT ?"525    args.append(limit)526    conn = dbc()527    try:528        return [with_value(r) for r in conn.execute(q, args)]529    finally:530        conn.close()531532533@app.get("/v1/stablecoins/mints-burns")534def mints_burns(token: str | None = None, window: str = "24h",535                limit: int = Query(100, le=500)):536    """Issuance events: transfers from the zero address are mints, to it are537    burns (EVM + Tron; Solana mints don't transit an address — see docs)."""538    zeros = tuple(ZERO_ADDRESSES.values())539    since = int(time.time()) - parse_window(window)540    zp = ", ".join("?" for _ in zeros)541    q = (SELECT_T + ' WHERE timestamp >= ? '542         f'AND ("from" IN ({zp}) OR "to" IN ({zp}))')543    args = [since, *zeros, *zeros]544    if token:545        q += " AND transfers.symbol = ?"546        args.append(token.upper())547    q += " ORDER BY timestamp DESC LIMIT ?"548    args.append(limit)549    conn = dbc()550    try:551        out = []552        for r in conn.execute(q, args):553            d = with_value(r)554            d["direction"] = "mint" if r["from"] in zeros else "burn"555            out.append(d)556        return out557    finally:558        conn.close()559560561@app.get("/v1/stablecoins/volume/series")562@cached(ttl=60)563def volume_series(token: str = "USDT", window: str = "24h",564                  interval: str = "1h", chain: str | None = None):565    """Bucketed transfer volume over time — per chain, for line charts.566    Served from StatsWorker's precomputed buckets for the standard windows."""567    since = int(time.time()) - parse_window(window)568    step = parse_window(interval)569    conn0 = dbc()570    try:571        agg = conn0.execute(572            "SELECT * FROM agg_series WHERE window = ? AND symbol = ? "573            + ("AND chain = ? " if chain else "") + "ORDER BY t",574            (window, token.upper(), *([chain] if chain else [])),575        ).fetchall()576    finally:577        conn0.close()578    if agg and int(time.time()) - agg[0]["updated"] < 900:579        out = {}580        for r in agg:581            out.setdefault(r["chain"], []).append(582                {"t": r["t"], "volume": r["volume"], "transfers": r["transfers"]})583        worker_step = {"1h": 300, "24h": 3600, "7d": 21600}.get(window, step)584        return {"token": token.upper(), "since": since, "interval_s": worker_step,585                "as_of": agg[0]["updated"], "chains": out}586    q = (f"SELECT (timestamp / {step}) * {step} AS t, chain, decimals, "587         "COUNT(*) AS transfers, SUM(CAST(amount AS DOUBLE PRECISION)) AS raw_sum "588         "FROM transfers WHERE symbol = ? AND timestamp >= ?")589    args = [token.upper(), since]590    if chain:591        q += " AND chain = ?"592        args.append(chain)593    q += " GROUP BY t, chain, decimals ORDER BY t"594    conn = dbc()595    try:596        rows = conn.execute(q, args).fetchall()597        price = price_of(conn, token.upper())598    finally:599        conn.close()600    out = {}601    for r in rows:602        c = out.setdefault(r["chain"], {})603        b = c.setdefault(r["t"], {"volume": 0.0, "transfers": 0})604        b["volume"] = round(b["volume"] + (r["raw_sum"] or 0) / 10 ** r["decimals"] * price, 2)605        b["transfers"] += r["transfers"]606    return {607        "token": token.upper(), "since": since, "interval_s": step,608        "chains": {609            c: [{"t": t, **v} for t, v in sorted(buckets.items())]610            for c, buckets in out.items()611        },612    }613614615@app.get("/v1/stablecoins/supply/series")616@cached(ttl=120)617def supply_series(token: str = "USDT", window: str = "7d"):618    """Supply snapshots over time per chain (hourly cadence)."""619    since = int(time.time()) - parse_window(window)620    conn = dbc()621    try:622        rows = conn.execute(623            "SELECT chain, supply, decimals, timestamp FROM supply_snapshots "624            "WHERE symbol = ? AND timestamp >= ? ORDER BY timestamp",625            (token.upper(), since),626        ).fetchall()627    finally:628        conn.close()629    out = {}630    for r in rows:631        out.setdefault(r["chain"], []).append(632            {"t": r["timestamp"],633             "supply": float(format_amount(r["supply"], r["decimals"] or 6))}634        )635    return {"token": token.upper(), "since": since, "chains": out}636637638@app.get("/v1/stablecoins/flows")639@cached(ttl=60)640def flows(token: str | None = None, window: str = "7d", interval: str = "1d"):641    """Net issuance over time: mints (from zero addr) minus burns (to zero)."""642    zeros = tuple(ZERO_ADDRESSES.values())643    since = int(time.time()) - parse_window(window)644    step = parse_window(interval)645    zp = ", ".join("?" for _ in zeros)646    q = (f"SELECT (timestamp / {step}) * {step} AS t, transfers.symbol AS symbol, "647         f'SUM(CASE WHEN "from" IN ({zp}) THEN {USD_PRICED} ELSE 0 END) AS minted, '648         f'SUM(CASE WHEN "to" IN ({zp}) THEN {USD_PRICED} ELSE 0 END) AS burned '649         f'FROM transfers {PRICE_JOIN} '650         f'WHERE timestamp >= ? AND ("from" IN ({zp}) OR "to" IN ({zp}))')651    args = [*zeros, *zeros, since, *zeros, *zeros]652    if token:653        q += " AND transfers.symbol = ?"654        args.append(token.upper())655    q += " GROUP BY t, transfers.symbol ORDER BY t"656    conn = dbc()657    try:658        rows = conn.execute(q, args).fetchall()659    finally:660        conn.close()661    out = {}662    for r in rows:663        out.setdefault(r["symbol"], []).append({664            "t": r["t"], "minted": round(r["minted"] or 0, 2),665            "burned": round(r["burned"] or 0, 2),666            "net": round((r["minted"] or 0) - (r["burned"] or 0), 2),667        })668    return {"since": since, "interval_s": step, "tokens": out}669670671@app.get("/v1/{chain}/transfers")672def chain_transfers(chain: str, symbol: str | None = None,673                    before_block: int | None = None,674                    min_amount: float | None = None,675                    limit: int = Query(50, le=500)):676    """Recent indexed transfers on one chain, all tokens (paginated)."""677    if chain not in CHAINS:678        raise HTTPException(404, f"unknown chain '{chain}'")679    q = SELECT_T + " WHERE chain = ?"680    args = [chain]681    if symbol:682        q += " AND transfers.symbol = ?"683        args.append(symbol.upper())684    if min_amount is not None:685        q += f" AND {USD_PRICED} >= ?"686        args.append(min_amount)687    if before_block is not None:688        q += " AND block < ?"689        args.append(before_block)690    q += " ORDER BY block DESC, log_index DESC LIMIT ?"691    args.append(limit)692    conn = dbc()693    try:694        return [with_value(r) for r in conn.execute(q, args)]695    finally:696        conn.close()697698699@app.get("/v1/{chain}/summary")700@cached(ttl=30)701def chain_summary(chain: str):702    if chain not in CHAINS:703        raise HTTPException(404, f"unknown chain '{chain}'")704    since = int(time.time()) - 86400705    conn = dbc()706    try:707        vol = conn.execute(708            f"SELECT transfers.symbol AS symbol, COUNT(*) AS transfers, "709            f"SUM({USD_PRICED}) AS volume "710            f"FROM transfers {PRICE_JOIN} "711            "WHERE chain = ? AND timestamp >= ? GROUP BY transfers.symbol "712            "ORDER BY volume DESC",713            (chain, since),714        ).fetchall()715        cur = conn.execute("SELECT * FROM cursors WHERE chain = ?", (chain,)).fetchone()716        total = conn.execute(717            "SELECT COUNT(*) AS n FROM transfers WHERE chain = ?", (chain,)718        ).fetchone()719    finally:720        conn.close()721    cfg = CHAINS[chain]722    return {723        "chain": chain,724        "family": cfg.get("family", "evm"),725        "chain_id": cfg.get("chain_id"),726        "block_time_s": cfg.get("block_time"),727        "tokens": TOKENS.get(chain, []),728        "indexed_transfers": total["n"] if total else 0,729        "volume_24h": [730            {"symbol": r["symbol"], "transfers": r["transfers"],731             "volume": round(r["volume"] or 0, 2)} for r in vol732        ],733        "cursor": dict(cur) if cur else None,734    }735736737@app.get("/v1/search")738def search(q: str):739    """Classify a query: tx hash / address / token symbol / chain name.740    Checks the index first (so we can say WHICH chain a tx lives on)."""741    s = q.strip()742    if not s:743        raise HTTPException(400, "empty query")744    if s.upper() in {t["symbol"] for toks in TOKENS.values() for t in toks}:745        return {"type": "token", "symbol": s.upper()}746    if s.lower() in CHAINS:747        return {"type": "chain", "chain": s.lower()}748    conn = dbc()749    try:750        needles = (s, s.lower())751        r = conn.execute(752            "SELECT chain, tx_hash FROM transfers WHERE tx_hash IN (?, ?) LIMIT 1", needles753        ).fetchone()754        if r:755            return {"type": "tx", "chain": r["chain"], "hash": r["tx_hash"]}756        r = conn.execute(757            'SELECT COUNT(*) AS n FROM transfers WHERE "from" IN (?, ?) OR "to" IN (?, ?)',758            needles + needles,759        ).fetchone()760        if r and r["n"]:761            return {"type": "address", "address": s, "indexed_transfers": r["n"]}762    finally:763        conn.close()764    # shape-based fallback for things we haven't indexed (yet)765    import re766    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):767        return {"type": "tx", "chain": None, "hash": s}768    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) \769       or re.fullmatch(r"[1-9A-HJ-NP-Za-km-z]{32,44}", s):770        return {"type": "address", "address": s, "indexed_transfers": 0}771    return {"type": "unknown"}772773774# -- live stream ---------------------------------------------------------775776777def _poll_transfers(since_ts):778    conn = dbc()779    try:780        rows = conn.execute(781            SELECT_T + " WHERE timestamp >= ? ORDER BY timestamp ASC LIMIT 1000",782            (since_ts,)783        ).fetchall()784        return [with_value(r) for r in rows]785    finally:786        conn.close()787788789@app.websocket("/v1/stream/transfers")790async def stream_transfers(ws: WebSocket):791    """Pushes new transfers every ~2s. Optional query params:792    token=USDT  chain=ethereum  min_usd=1000"""793    await ws.accept()794    token = (ws.query_params.get("token") or "").upper() or None795    chain = ws.query_params.get("chain")796    min_usd = float(ws.query_params.get("min_usd") or 0)797    last_ts = int(time.time()) - 2798    seen = set()799    try:800        while True:801            rows = await asyncio.to_thread(_poll_transfers, last_ts)802            batch = []803            for r in rows:804                key = (r["chain"], r["tx_hash"], r["log_index"])805                if key in seen:806                    continue807                if token and r["symbol"] != token:808                    seen.add(key)809                    continue810                if chain and r["chain"] != chain:811                    seen.add(key)812                    continue813                if min_usd and r.get("usd", float(r["value"])) < min_usd:814                    seen.add(key)815                    continue816                seen.add(key)817                batch.append(r)818            if batch:819                await ws.send_json(batch)820            if rows:821                new_last = max(r["timestamp"] or last_ts for r in rows)822                if new_last > last_ts:823                    last_ts = new_last824                    # only remember keys that can still reappear in queries825                    seen = {k for k in seen} if len(seen) < 50_000 else set()826            await asyncio.sleep(2)827    except WebSocketDisconnect:828        pass829830831# -- dashboard (must mount last: "/" catches everything below the API routes) --832app.mount("/", StaticFiles(directory=str(pathlib.Path(__file__).parent.parent / "ui"), html=True), name="ui")833