#!/usr/bin/env python3 # Author: Simon-Pierre Boucher # Mail: contact@spboucher.ai """Run the indexer: one worker thread per chain + price/supply workers, all under a watchdog that restarts any thread that dies. python3 run_indexer.py # all configured chains python3 run_indexer.py --chains ethereum base # a subset """ import argparse import logging import threading from indexer import config, db from indexer.adapters.bitcoin import BitcoinIndexer from indexer.adapters.solana import SolanaIndexer from indexer.adapters.tron import TronIndexer from indexer.enrich import PriceWorker, StatsWorker, SupplyWorker from indexer.ingest import ChainIndexer ADAPTERS = { "evm": ChainIndexer, "tron": TronIndexer, "solana": SolanaIndexer, "bitcoin": BitcoinIndexer, } def main(): ap = argparse.ArgumentParser() ap.add_argument("--chains", nargs="*", help="subset of chains (default: all)") ap.add_argument("--db", default=config.db_path()) args = ap.parse_args() logging.basicConfig( level=logging.INFO, format="%(asctime)s %(levelname)-7s %(message)s", datefmt="%H:%M:%S", ) chains = config.load_chains() tokens = config.load_tokens() selected = config.selected_chains(chains, tokens, args.chains, set(ADAPTERS)) conn = db.connect(args.db) # main-thread conn: schema + token metadata only for chain in selected: db.upsert_chain(conn, chain, chains[chain].get("family", "evm"), chains[chain].get("chain_id")) db.upsert_tokens(conn, chain, tokens[chain]) conn.close() stop = threading.Event() # every worker is a (name, factory) pair so the watchdog can rebuild it — # workers keep their own retry loops; this is the belt AND the suspenders def chain_factory(c): return lambda: ADAPTERS[chains[c].get("family", "evm")]( c, chains[c], tokens[c], args.db).run(stop) factories = {c: chain_factory(c) for c in selected} factories["prices"] = lambda: PriceWorker( {c: tokens[c] for c in selected}, args.db).run(stop) factories["supply"] = lambda: SupplyWorker( chains, {c: tokens[c] for c in selected}, args.db).run(stop) factories["stats"] = lambda: StatsWorker(args.db).run(stop) def spawn(name): t = threading.Thread(target=guarded(name), name=name, daemon=True) t.start() return t def guarded(name): def inner(): try: factories[name]() except Exception as e: # last-resort: workers shouldn't get here logging.error("worker %s died: %s", name, e) return inner threads = {name: spawn(name) for name in factories} logging.info("started %d chain workers + prices + supply (watchdog on)", len(selected)) try: while not stop.is_set(): stop.wait(30) for name, t in list(threads.items()): if not t.is_alive() and not stop.is_set(): logging.warning("watchdog: restarting dead worker %s", name) threads[name] = spawn(name) except KeyboardInterrupt: logging.info("stopping...") stop.set() for t in threads.values(): t.join(timeout=10) if __name__ == "__main__": main()