#!/usr/bin/env python3
"""Uptober eligibility engine: who BOUGHT $CARNAGE / $PEPENOM on or after the start
time, and kept holding it.

Standard library only. Read-only JSON-RPC against public Solana RPCs. It never
signs, sends or asks for anything; it only reads public chain data.

Subcommands
  check     <wallet> [<wallet> ...]   per-wallet report (JSON with --json)
  snapshot  --wallets FILE --out DIR  run check for a list, write JSON/CSV/eligible.txt + SHA256SUMS
  draw      --wallets FILE (--seed S | --drand-round N)   order wallets by sha256(seed + wallet)
  drand-round --at ISO-TIME           drand quicknet round number for a UTC time
  candidates --token SYM --out FILE   every wallet that bought SYM through its pool since start (no sign-up needed)
  hash      FILE [FILE ...]           sha256 of files (e.g. exclusions.json before launch)

Rule mapping (see research/engine.md):
  bought        = token increase in a tx the wallet SIGNED, that came out of a program-owned
                  (off-curve) account such as a pool vault, while the wallet paid SOL/wSOL/USDC/USDT
  transfers_in  = token increase that came from another normal wallet (on-curve owner)
  net_bought    = bought - (sold + transfers_out + burned + other outflows) since start
  hold          = aggregate balance never below the line from the moment buys reach the line
  --hold-end T  = one run checks both: buys reach the line by `end` (only buys at or before `end`
                  count), and the balance never drops below the line from then until T (config
                  "hold_until"; pass --hold-end config). Same verdict as the old two-run procedure
                  (binding run to `end` + hold run to T with qualified_at <= end).
"""
import argparse
import csv
import datetime as dt
import hashlib
import http.client
import io
import json
import os
import random
import sys
import threading
import time
import urllib.error
import urllib.request
from collections import defaultdict
from concurrent.futures import ThreadPoolExecutor

VERSION = "1.1.0"
HERE = os.path.dirname(os.path.abspath(__file__))
UA = "uptober-eligibility/1.0 (+read-only)"

TOKEN_PROGRAM = "TokenkegQfeZyiNwAJbNbGKPFXCWuBvf9Ss623VQ5DA"
TOKEN_2022_PROGRAM = "TokenzQdBNbLqP5VEhdkAS6EPFLC1PHnBqCXEpPxuEb"
ATA_PROGRAM = "ATokenGPvbdGVxr1b2hvZbsiqW5xWH25efTNsLJA8knL"
WSOL = "So11111111111111111111111111111111111111112"
USDC = "EPjFWdd5AufqSSqeM2qN1xzybapC8G4wEGGkZwyTDt1v"
USDT = "Es9vMFrzaCERmJfrF4H2FYD4KConky8sNbu6RQWgP3dQ"
STABLES = {USDC: "USDC", USDT: "USDT"}
FUNDING_MINTS = {USDC, USDT, WSOL}

DRAND_QUICKNET = "52db9ba70e0cc0f6eaf7803dd07447a1f5477735fd3f661792ba94600c84e971"
DRAND_GENESIS = 1692803367
DRAND_PERIOD = 3
DRAND_URLS = ["https://api.drand.sh", "https://drand.cloudflare.com"]

# A buy must cost at least this much (after fees and token-account rent are removed).
MIN_PAID_LAMPORTS = 100_000          # 0.0001 SOL
MIN_PAID_STABLE_RAW = 10_000         # 0.01 USDC/USDT (6 decimals)
ROUND_TRIP_SECONDS = 300             # sell within 5 min of a buy -> review flag
BOT_TRADES_PER_DAY = 50              # same bar as the Sep 12 bot rule

# ----------------------------------------------------------------------------- base58 / ed25519 / PDA

B58 = "123456789ABCDEFGHJKLMNPQRSTUVWXYZabcdefghijkmnopqrstuvwxyz"
B58_IDX = {c: i for i, c in enumerate(B58)}


def b58decode(s):
    n = 0
    for c in s:
        n = n * 58 + B58_IDX[c]
    body = n.to_bytes((n.bit_length() + 7) // 8, "big") if n else b""
    return b"\x00" * (len(s) - len(s.lstrip("1"))) + body


def b58encode(b):
    n = int.from_bytes(b, "big")
    out = []
    while n:
        n, r = divmod(n, 58)
        out.append(B58[r])
    return "1" * (len(b) - len(b.lstrip(b"\x00"))) + "".join(reversed(out))


def is_address(s):
    try:
        return isinstance(s, str) and 32 <= len(s) <= 44 and len(b58decode(s)) == 32
    except (KeyError, ValueError):
        return False


_P = 2 ** 255 - 19
_D = (-121665 * pow(121666, _P - 2, _P)) % _P
_on_curve_cache = {}


def is_on_curve(addr):
    """True if the 32-byte key decompresses to an ed25519 point (a normal wallet key).
    Program-derived addresses (pool vault owners, escrows) are off-curve. Mirrors
    curve25519-dalek CompressedEdwardsY::decompress().is_some()."""
    if addr in _on_curve_cache:
        return _on_curve_cache[addr]
    b = b58decode(addr) if isinstance(addr, str) else addr
    y = (int.from_bytes(b, "little") & ((1 << 255) - 1)) % _P
    u = (y * y - 1) % _P
    v = (_D * y * y + 1) % _P
    x2 = u * pow(v, _P - 2, _P) % _P
    ok = x2 == 0 or pow(x2, (_P - 1) // 2, _P) == 1
    if isinstance(addr, str):
        _on_curve_cache[addr] = ok
    return ok


def find_program_address(seeds, program_id):
    pid = b58decode(program_id)
    for bump in range(255, -1, -1):
        h = hashlib.sha256(b"".join(seeds) + bytes([bump]) + pid + b"ProgramDerivedAddress").digest()
        if not is_on_curve(h):
            return b58encode(h), bump
    raise ValueError("no PDA")


def derive_ata(wallet, mint, token_program):
    return find_program_address([b58decode(wallet), b58decode(token_program), b58decode(mint)], ATA_PROGRAM)[0]


# ----------------------------------------------------------------------------- time helpers

def parse_time(s):
    if s is None or s == "" or s == "now":
        return int(time.time())
    if isinstance(s, (int, float)):
        return int(s)
    s = s.strip()
    if s.isdigit():
        return int(s)
    if s.endswith("Z"):
        s = s[:-1] + "+00:00"
    t = dt.datetime.fromisoformat(s)
    if t.tzinfo is None:
        t = t.replace(tzinfo=dt.timezone.utc)
    return int(t.timestamp())


def iso(ts):
    if ts is None:
        return None
    return dt.datetime.fromtimestamp(ts, dt.timezone.utc).strftime("%Y-%m-%dT%H:%M:%SZ")


# ----------------------------------------------------------------------------- RPC client

class RpcError(Exception):
    pass


class Rpc:
    """Rotating JSON-RPC client with per-endpoint pacing and backoff.

    endpoints: list of {"url": ..., "history": "full"|"recent", "min_interval": seconds}.
    getSignaturesForAddress only goes to "full" endpoints: a "recent" node silently returns a
    short list for older history, which would look like the start of the wallet's life.
    """

    RATE_CODES = {-32005, -32429, 429}
    MISSING_CODES = {-32007, -32009, -32004}   # slot skipped / not in long-term storage / block unavailable

    def __init__(self, endpoints, timeout=60, max_retries=10, verbose=False):
        self.endpoints = [dict(e) for e in endpoints]
        for e in self.endpoints:
            e.setdefault("history", "full")
            e.setdefault("min_interval", 0.12)
        self.timeout = timeout
        self.max_retries = max_retries
        self.verbose = verbose
        self.lock = threading.Lock()
        self.next_ok = {e["url"]: 0.0 for e in self.endpoints}
        self.rr = 0
        self.calls = 0
        self.stats = {e["url"]: {"calls": 0, "rate_limited": 0, "errors": 0} for e in self.endpoints}

    def _pick(self, need_full, exclude=None):
        with self.lock:
            cands = [e for e in self.endpoints if (not need_full or e["history"] == "full")]
            if exclude:
                rest = [e for e in cands if e["url"] not in exclude]
                cands = rest or cands
            if not cands:
                raise RpcError("no endpoint with full history configured")
            self.rr = (self.rr + 1) % len(cands)
            order = cands[self.rr:] + cands[:self.rr]
            now = time.monotonic()
            e = min(order, key=lambda x: self.next_ok[x["url"]])
            wait = max(0.0, self.next_ok[e["url"]] - now)
            self.next_ok[e["url"]] = max(now, self.next_ok[e["url"]]) + e["min_interval"]
            self.calls += 1
            self.stats[e["url"]]["calls"] += 1
        if wait:
            time.sleep(wait)
        return e["url"]

    def _penalize(self, url, secs, kind="errors"):
        with self.lock:
            self.next_ok[url] = max(self.next_ok[url], time.monotonic() + secs)
            self.stats[url][kind] += 1

    def call(self, method, params, need_full=False, retry_null=False):
        last = None
        tried_null = set()
        for attempt in range(self.max_retries):
            url = self._pick(need_full, exclude=tried_null)
            body = json.dumps({"jsonrpc": "2.0", "id": 1, "method": method, "params": params}).encode()
            req = urllib.request.Request(url, data=body, headers={"Content-Type": "application/json", "User-Agent": UA})
            backoff = min(20.0, 0.6 * (2 ** attempt)) + random.random() * 0.5
            try:
                with urllib.request.urlopen(req, timeout=self.timeout) as r:
                    data = json.loads(r.read())
            except urllib.error.HTTPError as e:
                ra = e.headers.get("Retry-After") if e.headers else None
                wait = float(ra) if ra and ra.replace(".", "", 1).isdigit() else backoff
                self._penalize(url, wait, "rate_limited" if e.code == 429 else "errors")
                last = "HTTP %s from %s" % (e.code, url)
                if self.verbose:
                    print("  rpc %s %s -> %s, wait %.1fs" % (method, url, last, wait), file=sys.stderr)
                continue
            except (urllib.error.URLError, TimeoutError, ConnectionError, http.client.HTTPException,
                    json.JSONDecodeError, OSError) as e:
                self._penalize(url, backoff)
                last = "%s from %s" % (type(e).__name__, url)
                continue
            if "error" in data:
                err = data["error"] or {}
                code, msg = err.get("code"), str(err.get("message", ""))
                if code in self.RATE_CODES or "rate" in msg.lower() or "too many" in msg.lower():
                    self._penalize(url, backoff, "rate_limited")
                    last = "rate limited (%s) at %s" % (code, url)
                    continue
                if code in self.MISSING_CODES:
                    return None
                if code in (-32603, -32000, -32001, -32002, -32003, -32014, -32016) or code is None:
                    self._penalize(url, backoff)
                    last = "%s %s at %s" % (code, msg[:120], url)
                    continue
                raise RpcError("%s: %s %s" % (method, code, msg[:300]))
            res = data.get("result")
            if res is None and retry_null and len(tried_null) < len(self.endpoints) - 1:
                tried_null.add(url)
                last = "null result at %s" % url
                continue
            return res
        raise RpcError("%s failed after %d tries: %s" % (method, self.max_retries, last))


# ----------------------------------------------------------------------------- chain reads

class Chain:
    def __init__(self, rpc, cache_dir=None, workers=4, max_sig_pages=50):
        self.rpc = rpc
        self.cache_dir = cache_dir
        self.workers = workers
        self.max_sig_pages = max_sig_pages
        self._mem = {}
        self._mint_info = {}
        recent = [e.get("recent_hours", 12) for e in rpc.endpoints if e["history"] != "full"]
        self.recent_window_s = 3600 * min(recent) if recent else 10 ** 12
        if cache_dir:
            os.makedirs(os.path.join(cache_dir, "tx"), exist_ok=True)

    # -- mint
    def mint_info(self, mint):
        if mint not in self._mint_info:
            res = self.rpc.call("getAccountInfo", [mint, {"encoding": "jsonParsed", "commitment": "finalized"}])
            if not res or not res.get("value"):
                raise RpcError("mint %s not found" % mint)
            v = res["value"]
            self._mint_info[mint] = {"program": v["owner"], "decimals": v["data"]["parsed"]["info"]["decimals"]}
        return self._mint_info[mint]

    # -- token accounts currently owned
    def token_accounts(self, wallet, mint):
        res = self.rpc.call("getTokenAccountsByOwner",
                            [wallet, {"mint": mint}, {"encoding": "jsonParsed", "commitment": "finalized"}])
        out = {}
        for it in (res or {}).get("value", []):
            info = it["account"]["data"]["parsed"]["info"]
            out[it["pubkey"]] = int(info["tokenAmount"]["amount"])
        return out

    # -- signatures newest -> oldest, stop below start_ts
    def signatures_since(self, address, start_ts, max_pages=None):
        out, before = [], None
        for _ in range(max_pages or self.max_sig_pages):
            p = {"limit": 1000, "commitment": "finalized"}
            if before:
                p["before"] = before
            res = self.rpc.call("getSignaturesForAddress", [address, p], need_full=True) or []
            for s in res:
                bt = s.get("blockTime")
                if bt is not None and bt < start_ts:
                    return out, True
                out.append(s)
            if len(res) < 1000:
                return out, True
            before = res[-1]["signature"]
        return out, False

    def history_summary(self, address, max_pages):
        """Walk the full signature history (capped). Returns count, oldest/newest and the oldest page."""
        count, before, newest, last_page, complete = 0, None, None, [], False
        for _ in range(max_pages):
            p = {"limit": 1000, "commitment": "finalized"}
            if before:
                p["before"] = before
            res = self.rpc.call("getSignaturesForAddress", [address, p], need_full=True) or []
            if not res:
                complete = True
                break
            if newest is None:
                newest = res[0].get("blockTime")
            count += len(res)
            last_page = res
            if len(res) < 1000:
                complete = True
                break
            before = res[-1]["signature"]
        oldest = last_page[-1] if last_page else None
        return {"count": count, "complete": complete, "newest_time": newest,
                "oldest_time": oldest.get("blockTime") if oldest else None,
                "oldest_page": last_page}

    # -- transactions (immutable once finalized -> disk cache, slimmed)
    def _cache_path(self, sig):
        return os.path.join(self.cache_dir, "tx", sig[:2], sig + ".json") if self.cache_dir else None

    def get_tx(self, sig, block_time=None):
        if isinstance(sig, dict):
            sig, block_time = sig["signature"], sig.get("blockTime")
        if sig in self._mem:
            return self._mem[sig]
        path = self._cache_path(sig)
        if path and os.path.exists(path):
            try:
                with open(path) as f:
                    tx = json.load(f)
                self._mem[sig] = tx
                return tx
            except (OSError, ValueError):
                pass
        # "recent" nodes return null for older txs: send those straight to full-history nodes
        old = block_time is not None and block_time < time.time() - self.recent_window_s
        raw = self.rpc.call("getTransaction", [sig, {"encoding": "jsonParsed", "maxSupportedTransactionVersion": 1,
                                                     "commitment": "finalized"}], need_full=old, retry_null=True)
        if raw is None:
            return None
        tx = slim_tx(raw, sig)
        self._mem[sig] = tx
        if path:
            os.makedirs(os.path.dirname(path), exist_ok=True)
            tmp = path + ".%d.tmp" % threading.get_ident()
            with open(tmp, "w") as f:
                json.dump(tx, f, separators=(",", ":"))
            os.replace(tmp, path)
        return tx

    def get_txs(self, sigs):
        """sigs: signature strings or getSignaturesForAddress entries (dicts carry blockTime for routing)."""
        uniq = {}
        for s in sigs:
            k = s["signature"] if isinstance(s, dict) else s
            uniq.setdefault(k, s)
        if not uniq:
            return {}
        with ThreadPoolExecutor(max_workers=self.workers) as ex:
            res = list(ex.map(self.get_tx, uniq.values()))
        return dict(zip(uniq.keys(), res))


def slim_tx(raw, sig):
    meta = raw.get("meta") or {}
    msg = raw["transaction"]["message"]
    keys, signers = [], []
    for k in msg["accountKeys"]:
        if isinstance(k, dict):
            keys.append(k["pubkey"])
            if k.get("signer"):
                signers.append(k["pubkey"])
        else:
            keys.append(k)
    if signers == [] and "header" in msg:   # "json" encoding fallback
        signers = keys[:msg["header"]["numRequiredSignatures"]]
    programs = set()
    for ix in msg.get("instructions", []):
        if ix.get("programId"):
            programs.add(ix["programId"])
    for inner in meta.get("innerInstructions") or []:
        for ix in inner.get("instructions", []):
            if ix.get("programId"):
                programs.add(ix["programId"])

    def tb(lst):
        return [{"i": b["accountIndex"], "mint": b["mint"], "owner": b.get("owner"),
                 "amount": b["uiTokenAmount"]["amount"]} for b in (lst or [])]

    return {"sig": sig, "slot": raw.get("slot"), "blockTime": raw.get("blockTime"), "version": raw.get("version"),
            "err": meta.get("err"), "fee": meta.get("fee", 0), "keys": keys, "signers": signers,
            "pre": meta.get("preBalances", []), "post": meta.get("postBalances", []),
            "preTok": tb(meta.get("preTokenBalances")), "postTok": tb(meta.get("postTokenBalances")),
            "programs": sorted(programs)}


# ----------------------------------------------------------------------------- per-transaction classification

def analyze_tx(tx, wallet, mint, flag_programs=(), count_token_swaps=False):
    """Classify one transaction from the point of view of `wallet` and `mint`.

    All amounts are raw integer token units of `mint`. Returns a dict with the wallet's
    delta split into categories that always add up to the delta.
    """
    keys = tx["keys"]
    pre = {b["i"]: b for b in tx["preTok"]}
    post = {b["i"]: b for b in tx["postTok"]}
    w_delta_by_mint = defaultdict(int)
    w_accounts = {}          # W-owned token accounts of `mint` -> post amount
    w_accounts_pre = {}      # -> pre amount
    others = defaultdict(int)   # other owners' delta of `mint`
    total = 0
    new_rent = closed_refund = 0
    for i in sorted(set(pre) | set(post)):
        a, b = pre.get(i), post.get(i)
        ref = b or a
        m = ref["mint"]
        owner = (b or {}).get("owner") or (a or {}).get("owner")
        pa = int(a["amount"]) if a else 0
        pb = int(b["amount"]) if b else 0
        d = pb - pa
        if owner == wallet:
            w_delta_by_mint[m] += d
            if a is None and b is not None and i < len(tx["post"]):
                lam = tx["post"][i] - tx["pre"][i] - (pb if m == WSOL else 0)
                new_rent += max(lam, 0)
            elif b is None and a is not None and i < len(tx["post"]):
                lam = tx["pre"][i] - tx["post"][i] - (pa if m == WSOL else 0)
                closed_refund += max(lam, 0)
        if m == mint:
            total += d
            if owner == wallet:
                w_accounts[keys[i]] = pb
                w_accounts_pre[keys[i]] = pa
            elif d:
                others[owner] += d

    delta = w_delta_by_mint.get(mint, 0)
    signed = wallet in tx["signers"]
    lam_delta = 0
    if wallet in keys:
        wi = keys.index(wallet)
        lam_delta = tx["post"][wi] - tx["pre"][wi]
    fee = tx["fee"] if keys and keys[0] == wallet else 0
    # SOL the wallet paid out, net of network fee and token-account rent, wSOL counted as SOL.
    # Rent for the wallet's new token accounts is only removed up to what the wallet itself paid
    # (gasless/relayed buys have the relayer pay the rent).
    raw = -lam_delta - fee - w_delta_by_mint.get(WSOL, 0)
    sol_spent = raw - min(new_rent, max(raw, 0)) + closed_refund
    stable_spent = {STABLES[m]: -w_delta_by_mint.get(m, 0) for m in STABLES if w_delta_by_mint.get(m, 0)}
    paid = sol_spent >= MIN_PAID_LAMPORTS or any(v >= MIN_PAID_STABLE_RAW for v in stable_spent.values())
    received = sol_spent <= -MIN_PAID_LAMPORTS or any(v <= -MIN_PAID_STABLE_RAW for v in stable_spent.values())
    other_spent = {m: -d for m, d in w_delta_by_mint.items() if d < 0 and m not in (mint, WSOL) and m not in STABLES}
    other_recv = {m: d for m, d in w_delta_by_mint.items() if d > 0 and m not in (mint, WSOL) and m not in STABLES}

    from_pda = sum(-d for o, d in others.items() if d < 0 and not is_on_curve(o))
    from_wallets = {o: -d for o, d in others.items() if d < 0 and is_on_curve(o)}
    to_pda = sum(d for o, d in others.items() if d > 0 and not is_on_curve(o))
    to_wallets = {o: d for o, d in others.items() if d > 0 and is_on_curve(o)}
    minted = max(0, total)
    burned_supply = max(0, -total)

    r = {"sig": tx["sig"], "time": tx["blockTime"], "slot": tx["slot"], "delta": delta, "signed": signed,
         "sol_spent_lamports": sol_spent, "stable_spent_raw": stable_spent,
         "bought": 0, "swap_in_other": 0, "otc_in": 0, "transfers_in": 0, "program_in": 0,
         "sold": 0, "swap_out_other": 0, "transfers_out": 0, "burned": 0, "program_out": 0,
         "in_from": {}, "out_to": {}, "paid_with": None,
         "flag_programs": sorted(set(tx.get("programs", [])) & set(flag_programs)),
         "w_accounts": w_accounts, "w_accounts_pre": w_accounts_pre}

    if delta > 0:
        rest = delta
        if signed and paid:
            b = min(rest, from_pda + minted)
            r["bought"] = b
            rest -= b
            if b:
                r["paid_with"] = "SOL" if sol_spent >= MIN_PAID_LAMPORTS else "/".join(sorted(stable_spent))
            o = min(rest, sum(from_wallets.values()))
            r["otc_in"] = o            # paid a wallet directly in the same tx (RFQ/OTC); not a pool buy
            rest -= o
        elif signed and other_spent and from_pda:
            s = min(rest, from_pda)
            if count_token_swaps:
                r["bought"] = s
                r["paid_with"] = "token:" + ",".join(sorted(other_spent))
            else:
                r["swap_in_other"] = s
            rest -= s
        if rest:
            t = min(rest, sum(from_wallets.values()))
            r["transfers_in"] = t
            rest -= t
            if t:
                r["in_from"] = from_wallets
        r["program_in"] = rest     # unsigned fills, claims, unlocks, airdrops from programs
    elif delta < 0:
        rest = -delta
        if received and to_pda:
            s = min(rest, to_pda)
            r["sold"] = s
            rest -= s
        elif other_recv and to_pda:
            s = min(rest, to_pda)
            r["swap_out_other"] = s
            rest -= s
        t = min(rest, sum(to_wallets.values()))
        r["transfers_out"] = t
        rest -= t
        if t:
            r["out_to"] = to_wallets
        bn = min(rest, burned_supply)
        r["burned"] = bn
        rest -= bn
        r["program_out"] = rest    # deposits, locks, staking, anything else
    kinds = [k for k in ("bought", "swap_in_other", "otc_in", "transfers_in", "program_in", "sold",
                         "swap_out_other", "transfers_out", "burned", "program_out") if r[k]]
    r["kind"] = "+".join(kinds) if kinds else ("none" if delta == 0 else "?")
    return r


# ----------------------------------------------------------------------------- per-wallet, per-mint

def token_report(chain, wallet, tok, start_ts, end_ts, flag_programs, count_token_swaps=False,
                 keep_events=True, hold_end_ts=None):
    """Replay the wallet's balance of one token. Buys count only in [start, end]. With hold_end_ts
    the replay continues to hold_end_ts: the hold (min balance since qualified_at) runs through it."""
    mint, threshold_ui = tok["mint"], tok["threshold"]
    replay_end = end_ts if hold_end_ts is None else max(end_ts, hold_end_ts)
    info = chain.mint_info(mint)
    dec = info["decimals"]
    unit = 10 ** dec
    threshold = int(round(threshold_ui * unit))

    current = chain.token_accounts(wallet, mint)
    accounts = set(current) | {derive_ata(wallet, mint, info["program"])}
    sig_lists, truncated = {}, False
    todo = list(accounts)
    txs = {}
    while todo:
        acct = todo.pop()
        sigs, complete = chain.signatures_since(acct, start_ts)
        truncated |= not complete
        sig_lists[acct] = [s for s in sigs if s.get("err") is None]
        fetched = chain.get_txs([s for s in sig_lists[acct] if s["signature"] not in txs])
        txs.update(fetched)
        # discovery: any other token account of this mint owned by the wallet seen in these txs
        for tx in fetched.values():
            if not tx:
                continue
            for b in tx["preTok"] + tx["postTok"]:
                if b["mint"] == mint and b.get("owner") == wallet:
                    a = tx["keys"][b["i"]]
                    if a not in accounts:
                        accounts.add(a)
                        todo.append(a)
    missing = [s for s, t in txs.items() if t is None]
    ordered = sorted((t for t in txs.values() if t and t.get("err") is None),
                     key=lambda t: (t["slot"] or 0, t["blockTime"] or 0))

    # starting balance of every account = pre-balance in its first tx since start, else current
    bal = {}
    for a in accounts:
        first = next((t for t in ordered if a in t["keys"]), None)
        if first is None:
            bal[a] = current.get(a, 0)
        else:
            idx = first["keys"].index(a)
            pre = next((b for b in first["preTok"] if b["i"] == idx and b["mint"] == mint), None)
            bal[a] = int(pre["amount"]) if pre else 0
    start_balance = sum(bal.values())

    agg = defaultdict(int)
    events = []
    cum_bought = 0
    first_buy = qualified_at = None
    min_since_first_buy = min_since_qualified = None
    min_since_qualified_at = None
    in_from, out_to = defaultdict(int), defaultdict(int)
    flagged_programs = set()
    buy_times, sell_times = [], []
    closed = False                 # True once the replay is past `end` (hold period only)
    end_balance = None
    agg_end = None
    for t in ordered:
        if t["blockTime"] is not None and t["blockTime"] > replay_end:
            break
        if not closed and t["blockTime"] is not None and t["blockTime"] > end_ts:
            closed = True
            end_balance = sum(bal.values())
            agg_end = dict(agg)
        r = analyze_tx(t, wallet, mint, flag_programs, count_token_swaps)
        for a, amt in r["w_accounts"].items():
            bal[a] = amt
        balance = sum(bal.values())
        for k in ("bought", "swap_in_other", "otc_in", "transfers_in", "program_in", "sold", "swap_out_other",
                  "transfers_out", "burned", "program_out"):
            agg[k] += r[k]
        for o, v in r["in_from"].items():
            in_from[o] += v
        for o, v in r["out_to"].items():
            out_to[o] += v
        if r["bought"]:
            flagged_programs.update(r["flag_programs"])
            buy_times.append(r["time"])
            if not closed:         # buys after the close never count towards the line
                cum_bought += r["bought"]
                if first_buy is None:
                    first_buy = r["time"]
                if qualified_at is None and cum_bought >= threshold:
                    qualified_at = r["time"]
        if r["sold"]:
            sell_times.append(r["time"])
        if first_buy is not None:
            min_since_first_buy = balance if min_since_first_buy is None else min(min_since_first_buy, balance)
        if qualified_at is not None:
            if min_since_qualified is None or balance < min_since_qualified:
                min_since_qualified = balance
                min_since_qualified_at = r["time"]
        if keep_events and r["delta"]:
            ev = {"time": iso(r["time"]), "sig": r["sig"], "kind": r["kind"],
                  "delta": r["delta"] / unit, "balance_after": balance / unit,
                  "paid_with": r["paid_with"],
                  "sol_spent": round(r["sol_spent_lamports"] / 1e9, 6) if r["bought"] else None}
            if closed:
                ev["after_close"] = True
            events.append(ev)
    last_balance = sum(bal.values())
    if not closed:
        end_balance = last_balance
        agg_end = dict(agg)
    agg_all, agg = agg, defaultdict(int, agg_end)
    current_balance = sum(current.values())

    def outflow(a):
        return a["sold"] + a["swap_out_other"] + a["transfers_out"] + a["burned"] + a["program_out"]

    net_bought = agg["bought"] - outflow(agg)

    round_trips = sum(1 for s in sell_times if any(0 <= s - b <= ROUND_TRIP_SECONDS for b in buy_times))
    span_days = max(1.0, ((min(replay_end, int(time.time())) - start_ts) / 86400.0))
    trades_per_day = (len(buy_times) + len(sell_times)) / span_days

    out = {
        "symbol": tok["symbol"], "mint": mint, "decimals": dec, "token_program": info["program"],
        "threshold": threshold_ui,
        "token_accounts": sorted(accounts),
        "qualifying_bought": agg["bought"] / unit,
        "buys": len(buy_times),
        "first_buy": iso(first_buy),
        "qualified_at": iso(qualified_at),
        "net_bought": net_bought / unit,
        "transfers_in": agg["transfers_in"] / unit,
        "transfers_in_from": {o: v / unit for o, v in in_from.items()},
        "transfers_out": agg["transfers_out"] / unit,
        "transfers_out_to": {o: v / unit for o, v in out_to.items()},
        "sold": agg["sold"] / unit,
        "sells": len(sell_times),
        "swap_in_from_other_token": agg["swap_in_other"] / unit,
        "swap_out_to_other_token": agg["swap_out_other"] / unit,
        "otc_in": agg["otc_in"] / unit,
        "program_in": agg["program_in"] / unit,
        "program_out": agg["program_out"] / unit,
        "burned": agg["burned"] / unit,
        "balance_at_start": start_balance / unit,
        "balance_at_end": end_balance / unit,
        "current_balance": current_balance / unit,
        "min_balance_since_first_buy": None if min_since_first_buy is None else min_since_first_buy / unit,
        "min_balance_since_qualified": None if min_since_qualified is None else min_since_qualified / unit,
        "min_balance_since_qualified_at": iso(min_since_qualified_at),
        "txs_scanned": len(ordered),
        "history_truncated": truncated,
        "missing_txs": missing,
        "round_trips": round_trips,
        "trades_per_day": round(trades_per_day, 2),
        "flagged_programs": sorted(flagged_programs),
    }
    passes = {
        "bought_enough": agg["bought"] >= threshold,
        "net_bought_enough": net_bought >= threshold,
        "held_continuously": qualified_at is not None and min_since_qualified is not None
                              and min_since_qualified >= threshold,
        "holding_at_end": end_balance >= threshold,
    }
    if hold_end_ts is not None:
        # the hold run of the old two-run procedure, folded in: every buy and outflow up to hold_end
        net_hold = agg_all["bought"] - outflow(agg_all)
        out["hold_end"] = iso(hold_end_ts)
        out["hold_complete"] = hold_end_ts <= int(time.time())
        out["balance_at_hold_end"] = last_balance / unit
        out["net_bought_through_hold_end"] = net_hold / unit
        out["after_close"] = {k: (agg_all[k] - agg[k]) / unit for k in sorted(agg_all) if agg_all[k] != agg[k]}
        passes["holding_at_hold_end"] = last_balance >= threshold
        passes["net_bought_enough_through_hold_end"] = net_hold >= threshold
    out["passes"] = passes
    out["dropped_below_line"] = qualified_at is not None and not passes["held_continuously"]
    if replay_end >= int(time.time()) - 60 and last_balance != current_balance and not truncated:
        out["balance_check_mismatch"] = True
    if keep_events:
        out["events"] = events
    return out


# ----------------------------------------------------------------------------- wallet age, funder, exclusions

def funding_of(chain, address, max_pages, look=12):
    """First funder of `address`, from its oldest successful txs:
    - SOL: the earliest tx where its lamports go 0 -> >0; funder = the account that lost the most lamports;
    - else a token (gasless/USDC-only wallets): the earliest tx it did not sign in which one of its token
      accounts went up; funder = the owner of the account that lost that token (or the fee payer)."""
    h = chain.history_summary(address, max_pages)
    res = {"address": address, "tx_count": h["count"], "history_complete": h["complete"],
           "first_activity": iso(h["oldest_time"]) if h["complete"] else None,
           "oldest_seen": iso(h["oldest_time"]), "funder": None, "funding_tx": None, "funding_asset": None,
           "funding_amount": None}
    if not h["complete"]:
        return res
    oldest_first = [s for s in reversed(h["oldest_page"]) if s.get("err") is None][:look]
    txs = chain.get_txs(oldest_first)
    token_hit = None
    for s in oldest_first:
        t = txs.get(s["signature"])
        if not t or address not in t["keys"]:
            continue
        i = t["keys"].index(address)
        if t["pre"][i] == 0 and t["post"][i] > 0:
            deltas = [(t["post"][j] - t["pre"][j], k) for j, k in enumerate(t["keys"]) if k != address]
            d, k = min(deltas) if deltas else (0, None)
            res.update(funder=k if d < 0 else None, funding_tx=t["sig"], funding_asset="SOL",
                       funding_amount=t["post"][i] / 1e9)
            return res
        if token_hit is None and address not in t["signers"]:
            pre = {b["i"]: int(b["amount"]) for b in t["preTok"]}
            for b in t["postTok"]:
                # only stablecoins/wSOL count as funding; random meme tokens are usually spam "dust"
                if b["mint"] not in FUNDING_MINTS:
                    continue
                if b.get("owner") == address and int(b["amount"]) > pre.get(b["i"], 0):
                    src = None
                    for a in t["preTok"]:
                        post_amt = next((int(x["amount"]) for x in t["postTok"] if x["i"] == a["i"]), 0)
                        if a["mint"] == b["mint"] and a.get("owner") not in (None, address) and post_amt < int(a["amount"]):
                            src = a["owner"]
                            break
                    if src is None or not is_on_curve(src):
                        src = t["keys"][0]
                    token_hit = {"funder": src, "funding_tx": t["sig"], "funding_asset": b["mint"],
                                 "funding_amount": (int(b["amount"]) - pre.get(b["i"], 0))}
                    break
    if token_hit:
        res.update(token_hit)
    return res


def load_exclusions(path):
    if not path or not os.path.exists(path):
        return {}, [], None
    with open(path, "rb") as f:
        raw = f.read()
    data = json.loads(raw)
    entries = data.get("entries", data.get("addresses", [])) if isinstance(data, dict) else data
    excl, programs = {}, []
    for e in entries:
        if isinstance(e, str):
            excl[e] = {"label": "listed", "category": "listed"}
        else:
            if e.get("category") == "program":
                programs.append(e["address"])
            else:
                excl[e["address"]] = {"label": e.get("label", ""), "category": e.get("category", "listed")}
    return excl, programs, hashlib.sha256(raw).hexdigest()


def check_wallet(chain, wallet, cfg, excl, flag_programs, opts):
    t0 = time.time()
    calls0 = chain.rpc.calls
    start_ts, end_ts = opts["start_ts"], opts["end_ts"]
    hold_end_ts = opts.get("hold_end_ts")
    out = {"wallet": wallet, "start": iso(start_ts), "end": iso(end_ts), "checked_at": iso(int(time.time()))}
    if hold_end_ts is not None:
        out["hold_end"] = iso(hold_end_ts)
        out["hold_complete"] = hold_end_ts <= int(time.time())
    if not is_address(wallet):
        out["error"] = "not a valid Solana address"
        return out
    reasons, flags = [], []
    if wallet in excl:
        e = excl[wallet]
        reasons.append("listed: %s (%s)" % (e["category"], e["label"]))
    if not is_on_curve(wallet):
        reasons.append("program-owned (off-curve) address, not a person's wallet")

    if reasons and opts.get("skip_excluded"):
        out.update(excluded=True, exclusion_reasons=reasons, eligible=False)
        out["runtime_s"] = round(time.time() - t0, 2)
        out["rpc_calls"] = chain.rpc.calls - calls0
        return out

    # wallet age + funder chain
    age = funding_of(chain, wallet, opts["age_pages"])
    out["wallet_tx_count"] = age["tx_count"] if age["history_complete"] else ">=%d" % age["tx_count"]
    out["first_activity"] = age["first_activity"]
    out["oldest_seen"] = age["oldest_seen"]
    out["first_funder"] = age["funder"]
    out["funding_tx"] = age["funding_tx"]
    out["funding_asset"] = age["funding_asset"]
    chain_list = []
    nxt = age["funder"]
    for hop in range(opts["hops"]):
        if not nxt:
            break
        node = {"hop": hop + 1, "address": nxt, "listed": excl.get(nxt)}
        chain_list.append(node)
        if nxt in excl:
            break
        if hop + 1 >= opts["hops"]:
            # last hop: no funder lookup, but still mark exchanges/services so funding groups skip them
            if not chain.history_summary(nxt, opts["hub_pages"])["complete"]:
                node["hub"] = True
            break
        f = funding_of(chain, nxt, opts["hub_pages"])
        if not f["history_complete"]:
            node["hub"] = True    # >= hub_pages*1000 txs: exchange, router or service; stop here
            break
        nxt = f["funder"]
    out["funder_chain"] = chain_list
    for node in chain_list:
        if node["listed"]:
            reasons.append("funded (hop %d) by listed %s %s (%s)" % (node["hop"], node["listed"]["category"],
                                                                      node["address"], node["listed"]["label"]))
            break

    if age["history_complete"] and age["first_activity"] and parse_time(age["first_activity"]) >= opts["fresh_after"]:
        flags.append("fresh_wallet: first activity %s" % age["first_activity"])
    if age["history_complete"] and not age["funder"]:
        flags.append("funder_unknown: no SOL/USDC funding found in the wallet's first txs")
    if not age["history_complete"]:
        flags.append("high_activity: more than %d lifetime txs (age/funder not traced)" % age["tx_count"])
    elif age["oldest_seen"]:
        life_days = max(1.0, (time.time() - parse_time(age["oldest_seen"])) / 86400.0)
        if age["tx_count"] / life_days > BOT_TRADES_PER_DAY:
            flags.append("high_activity: %.0f txs/day over wallet life" % (age["tx_count"] / life_days))

    tokens = {}
    for tok in cfg["tokens"]:
        rep = token_report(chain, wallet, tok, start_ts, end_ts, flag_programs,
                           opts.get("count_token_swaps", False), keep_events=opts.get("events", True),
                           hold_end_ts=hold_end_ts)
        tokens[tok["symbol"]] = rep
        for o in rep["transfers_in_from"]:
            if o in excl:
                reasons.append("received %s by transfer from listed %s %s" % (tok["symbol"], excl[o]["category"], o))
        for o in rep["transfers_out_to"]:
            if o in excl:
                flags.append("sent %s to listed %s %s" % (tok["symbol"], excl[o]["category"], o))
        if rep["flagged_programs"]:
            reasons.append("%s bought through flagged program %s" % (tok["symbol"], ",".join(rep["flagged_programs"])))
        if rep["round_trips"]:
            flags.append("%s: %d sell(s) within %ds of a buy" % (tok["symbol"], rep["round_trips"], ROUND_TRIP_SECONDS))
        if rep["trades_per_day"] > BOT_TRADES_PER_DAY:
            flags.append("%s: %.0f trades/day since start" % (tok["symbol"], rep["trades_per_day"]))
        if rep["program_in"] or rep["otc_in"] or rep["swap_in_from_other_token"]:
            flags.append("%s: inflow not counted as a buy (program %.0f, otc %.0f, token-swap %.0f) - review"
                         % (tok["symbol"], rep["program_in"], rep["otc_in"], rep["swap_in_from_other_token"]))
        if rep["history_truncated"] or rep["missing_txs"]:
            flags.append("%s: incomplete data (truncated=%s, missing=%d) - rerun" %
                         (tok["symbol"], rep["history_truncated"], len(rep["missing_txs"])))
        if rep.get("balance_check_mismatch"):
            flags.append("%s: replayed balance != current balance - review" % tok["symbol"])
        if rep["balance_at_end"] == rep["threshold"]:
            flags.append("%s: holds exactly the line" % tok["symbol"])

    if hold_end_ts is not None and hold_end_ts > int(time.time()):
        flags.append("hold_incomplete: hold period ends %s; checked through %s - rerun after it ends"
                     % (iso(hold_end_ts), iso(int(time.time()))))
    if opts.get("strict_bots"):
        for f in flags:
            if f.startswith("high_activity") or "trades/day since start" in f:
                reasons.append("bot rule (> %d txs/day, Sep 12 precedent): %s" % (BOT_TRADES_PER_DAY, f))
    out["tokens"] = tokens
    out["excluded"] = bool(reasons)
    out["exclusion_reasons"] = reasons
    out["review_flags"] = flags
    need = ["bought_enough", "held_continuously", "holding_at_end"] + (["net_bought_enough"] if opts.get("net_rule", True) else [])
    if hold_end_ts is not None:
        need += ["holding_at_hold_end"] + (["net_bought_enough_through_hold_end"] if opts.get("net_rule", True) else [])
    out["qualifies_on_chain"] = all(all(t["passes"][k] for k in need) for t in tokens.values())
    out["eligible"] = out["qualifies_on_chain"] and not out["excluded"] and not any(
        t["history_truncated"] or t["missing_txs"] for t in tokens.values())
    out["runtime_s"] = round(time.time() - t0, 2)
    out["rpc_calls"] = chain.rpc.calls - calls0
    return out


# ----------------------------------------------------------------------------- candidates from a pool

def pool_buyers(chain, pool, mint, start_ts, end_ts, max_pages=200):
    """Every wallet that bought `mint` through `pool` in [start, end], read from the pool's own history.
    Lets the team build the entrant list without ever asking anyone for an address."""
    sigs, complete = chain.signatures_since(pool, start_ts, max_pages)
    sigs = [s for s in sigs if s.get("err") is None and (s.get("blockTime") or 0) <= end_ts]
    txs = chain.get_txs(sigs)
    buyers = {}
    for t in txs.values():
        if not t:
            continue
        for w in t["signers"]:
            r = analyze_tx(t, w, mint)
            if r["bought"]:
                b = buyers.setdefault(w, {"bought_raw": 0, "buys": 0, "first_buy": r["time"]})
                b["bought_raw"] += r["bought"]
                b["buys"] += 1
                b["first_buy"] = min(b["first_buy"], r["time"])
    return buyers, complete, len(sigs), sum(1 for t in txs.values() if t is None)


# ----------------------------------------------------------------------------- draw (Sep 12 precedent)

def draw_order(seed, wallets):
    """Order by sha256(seed + wallet) as lowercase hex, ascending (Sep 12 2026 PEPENOM rule)."""
    return sorted(((hashlib.sha256((seed + w).encode()).hexdigest(), w) for w in wallets))


def list_sha256(wallets):
    """sha256 of the sorted list joined by '\\n', no trailing newline (Sep 12 eligibleSha256 format)."""
    return hashlib.sha256("\n".join(sorted(wallets)).encode()).hexdigest()


def drand_round_at(ts):
    return (ts - DRAND_GENESIS) // DRAND_PERIOD + 1


def drand_time_of(rnd):
    return DRAND_GENESIS + (rnd - 1) * DRAND_PERIOD


def drand_fetch(rnd):
    if drand_time_of(rnd) > time.time():
        raise RpcError("drand round %d is not out yet: it is emitted at %s" % (rnd, iso(drand_time_of(rnd))))
    last = None
    for base in DRAND_URLS:
        url = "%s/%s/public/%d" % (base, DRAND_QUICKNET, rnd)
        try:
            req = urllib.request.Request(url, headers={"User-Agent": UA})
            with urllib.request.urlopen(req, timeout=20) as r:
                d = json.loads(r.read())
            if int(d["round"]) != rnd:
                raise ValueError("round mismatch")
            # randomness = sha256(signature) for drand; check it so a proxy can't swap it
            if hashlib.sha256(bytes.fromhex(d["signature"])).hexdigest() != d["randomness"]:
                raise ValueError("randomness != sha256(signature)")
            d["source"] = url
            return d
        except Exception as e:  # noqa: BLE001 - try next mirror
            last = e
    raise RpcError("drand round %d unavailable: %s" % (rnd, last))


# ----------------------------------------------------------------------------- config / CLI

def load_config(path):
    with open(path) as f:
        cfg = json.load(f)
    env = os.environ.get("UPTOBER_RPC")
    if env:   # private full-history RPC first, e.g. a Helius/Triton URL
        cfg["rpcs"] = [{"url": u, "history": "full", "min_interval": 0.05} for u in env.split(",")] + cfg["rpcs"]
    return cfg


def build_chain(cfg, args):
    rpc = Rpc(cfg["rpcs"], verbose=getattr(args, "verbose", False))
    cache = None if getattr(args, "no_cache", False) else (args.cache or os.path.join(HERE, ".cache"))
    return Chain(rpc, cache_dir=cache, workers=args.workers)


def make_opts(cfg, args):
    if getattr(args, "final", False):     # the Nov 12 verification in one run
        args.end = args.end or "config"
        args.hold_end = args.hold_end or "config"
        args.strict_bots = True
    start_ts = parse_time(args.start or cfg["start"])
    end_raw = args.end if args.end else ("now" if args.cmd == "check" else cfg.get("end", "now"))
    if end_raw == "config":
        end_raw = cfg["end"]
    end_ts = parse_time(end_raw)
    hold_raw = getattr(args, "hold_end", None)
    if hold_raw in (None, "", "none"):
        hold_end_ts = None
    else:
        if hold_raw == "config":
            if not cfg.get("hold_until"):
                raise SystemExit("--hold-end config: config has no hold_until")
            hold_raw = cfg["hold_until"]
        hold_end_ts = parse_time(hold_raw)
        if hold_end_ts < end_ts:
            raise SystemExit("--hold-end (%s) is before --end (%s)" % (iso(hold_end_ts), iso(end_ts)))
    return {"start_ts": start_ts, "end_ts": end_ts, "hold_end_ts": hold_end_ts,
            "hops": args.hops, "age_pages": args.age_pages,
            "hub_pages": args.hub_pages, "fresh_after": parse_time(args.fresh_after or cfg.get("fresh_after") or cfg["start"]),
            "skip_excluded": args.skip_excluded, "count_token_swaps": args.count_token_swaps,
            "net_rule": not args.no_net_rule, "events": not getattr(args, "no_events", False),
            "strict_bots": args.strict_bots}


def apply_thresholds(cfg, args):
    for spec in args.threshold or []:
        sym, _, val = spec.partition("=")
        for t in cfg["tokens"]:
            if t["symbol"].upper() == sym.upper().lstrip("$"):
                t["threshold"] = float(val.replace(",", "").replace("_", ""))


def file_sha256(path):
    h = hashlib.sha256()
    with open(path, "rb") as f:
        for chunk in iter(lambda: f.read(1 << 16), b""):
            h.update(chunk)
    return h.hexdigest()


def summarize(r):
    lines = ["%s  eligible=%s  excluded=%s  (%.1fs, %s rpc calls)" % (
        r["wallet"], r.get("eligible"), r.get("excluded"), r.get("runtime_s", 0), r.get("rpc_calls"))]
    if r.get("error"):
        return r["wallet"] + "  ERROR " + r["error"]
    lines.append("  age: first activity %s, lifetime txs %s, first funder %s" % (
        r.get("first_activity") or ("before " + str(r.get("oldest_seen"))), r.get("wallet_tx_count"), r.get("first_funder")))
    for sym, t in (r.get("tokens") or {}).items():
        lines.append("  %-8s bought %s (%d buys, qualified %s)  net %s  in %s  out %s  sold %s" % (
            sym, fmt(t["qualifying_bought"]), t["buys"], t["qualified_at"], fmt(t["net_bought"]),
            fmt(t["transfers_in"]), fmt(t["transfers_out"]), fmt(t["sold"])))
        extra = ["%s %s" % (k, fmt(t[k])) for k in ("burned", "program_in", "program_out", "otc_in",
                                                    "swap_in_from_other_token", "swap_out_to_other_token") if t.get(k)]
        if extra:
            lines.append("           other: " + ", ".join(extra))
        lines.append("           balance start %s end %s now %s  min since qualified %s  passes %s" % (
            fmt(t["balance_at_start"]), fmt(t["balance_at_end"]), fmt(t["current_balance"]),
            fmt(t["min_balance_since_qualified"]), ",".join(k for k, v in t["passes"].items() if v) or "none"))
    for x in r.get("exclusion_reasons", []):
        lines.append("  EXCLUDED: " + x)
    for x in r.get("review_flags", []):
        lines.append("  review: " + x)
    return "\n".join(lines)


def fmt(x):
    if x is None:
        return "-"
    return "{:,.0f}".format(x) if abs(x) >= 100 else "{:,.4f}".format(x)


CSV_COLS = ["wallet", "eligible", "qualifies_on_chain", "excluded", "exclusion_reasons", "review_flags",
            "first_activity", "wallet_tx_count", "first_funder"]
TOKEN_COLS = ["qualifying_bought", "buys", "qualified_at", "net_bought", "sold", "transfers_in", "transfers_out",
              "program_in", "balance_at_end", "min_balance_since_qualified"]


def funding_groups(results, wallets=None):
    """One person = one funding group. Two wallets are in the same group when they share a non-exchange
    (non-hub) funder within the traced hops (default 2), or when one of them funded the other.
    results: check_wallet outputs. wallets: which wallets to group (default: the eligible ones).
    Returns every group (singletons included) as sorted lists, ordered by their first wallet."""
    if wallets is None:
        wallets = [r["wallet"] for r in results if r.get("eligible")]
    members = set(wallets)
    parent = {}

    def find(x):
        parent.setdefault(x, x)
        while parent[x] != x:
            parent[x] = parent[parent[x]]
            x = parent[x]
        return x

    def union(a, b):
        ra, rb = find(a), find(b)
        if ra != rb:
            parent[max(ra, rb)] = min(ra, rb)

    for r in results:
        w = r.get("wallet")
        if w not in members:
            continue
        find(w)
        for node in r.get("funder_chain") or []:
            if node.get("hub") or not node.get("address"):
                break
            union(w, node["address"])
    groups = defaultdict(list)
    for w in members:
        groups[find(w)].append(w)
    return sorted((sorted(g) for g in groups.values()), key=lambda g: g[0])


def write_snapshot(results, meta, outdir):
    os.makedirs(outdir, exist_ok=True)
    eligible = sorted(r["wallet"] for r in results if r.get("eligible"))
    # shared non-hub funders among wallets that qualify on chain -> one entry per funding source (review)
    by_funder = defaultdict(list)
    for r in results:
        if r.get("qualifies_on_chain") and r.get("first_funder"):
            hub = any(n.get("hub") for n in r.get("funder_chain", [])[:1])
            if not hub:
                by_funder[r["first_funder"]].append(r["wallet"])
    clusters = {f: sorted(ws) for f, ws in by_funder.items() if len(ws) > 1}
    for r in results:
        for f, ws in clusters.items():
            if r["wallet"] in ws:
                r.setdefault("review_flags", []).append("shared_funder %s with %d other entrant(s)" % (f, len(ws) - 1))
    groups = funding_groups(results, eligible)
    doc = dict(meta)
    doc.update({"eligible_count": len(eligible), "eligible_sha256": list_sha256(eligible),
                "qualified_count": len(groups),
                "eligible": eligible, "shared_funder_clusters": clusters,
                "funding_groups": [g for g in groups if len(g) > 1], "results": results})
    paths = {}
    paths["snapshot.json"] = os.path.join(outdir, "snapshot.json")
    with open(paths["snapshot.json"], "w") as f:
        json.dump(doc, f, indent=1, sort_keys=False)
    paths["snapshot.csv"] = os.path.join(outdir, "snapshot.csv")
    syms = [t["symbol"] for t in meta["config"]["tokens"]]
    with open(paths["snapshot.csv"], "w", newline="") as f:
        w = csv.writer(f)
        w.writerow(CSV_COLS + ["%s_%s" % (s, c) for s in syms for c in TOKEN_COLS])
        for r in results:
            row = [r.get("wallet"), r.get("eligible"), r.get("qualifies_on_chain"), r.get("excluded"),
                   " | ".join(r.get("exclusion_reasons", [])), " | ".join(r.get("review_flags", [])),
                   r.get("first_activity"), r.get("wallet_tx_count"), r.get("first_funder")]
            for s in syms:
                t = (r.get("tokens") or {}).get(s, {})
                row += [t.get(c) for c in TOKEN_COLS]
            w.writerow(row)
    paths["eligible.txt"] = os.path.join(outdir, "eligible.txt")
    with open(paths["eligible.txt"], "w") as f:
        f.write("\n".join(eligible))
    sums = os.path.join(outdir, "SHA256SUMS")
    with open(sums, "w") as f:
        for name, p in paths.items():
            f.write("%s  %s\n" % (file_sha256(p), name))
    return doc, sums


def read_wallet_file(path):
    out = []
    with open(path) as f:
        for line in f:
            s = line.split("#", 1)[0].strip().strip(",")
            if s:
                out.append(s)
    return list(dict.fromkeys(out))


def build_parser():
    ap = argparse.ArgumentParser(description=__doc__, formatter_class=argparse.RawDescriptionHelpFormatter)
    ap.add_argument("--config", default=os.path.join(HERE, "config.json"))
    ap.add_argument("--exclusions", default=None, help="exclusions JSON (default: from config)")
    sub = ap.add_subparsers(dest="cmd", required=True)

    def common(p):
        p.add_argument("--start", help="ISO time; default from config (2026-10-07T00:00:00Z)")
        p.add_argument("--end", help="ISO time, 'now' or 'config' (check: now; snapshot: config end)")
        p.add_argument("--hold-end", help="ISO time or 'config' (= config hold_until): buys still count only up to "
                                          "--end, but the hold must last until this time (one-run final check)")
        p.add_argument("--final", action="store_true",
                       help="shorthand for --end config --hold-end config --strict-bots")
        p.add_argument("--threshold", action="append", help="SYMBOL=amount, e.g. CARNAGE=666666")
        p.add_argument("--hops", type=int, default=2, help="funder hops to trace (default 2)")
        p.add_argument("--age-pages", type=int, default=20, help="max 1000-sig pages for wallet age (default 20)")
        p.add_argument("--hub-pages", type=int, default=3, help="funder with more pages than this = hub, stop")
        p.add_argument("--fresh-after", help="first activity at/after this -> fresh_wallet review flag")
        p.add_argument("--skip-excluded", action="store_true", help="do not scan wallets already on the list")
        p.add_argument("--count-token-swaps", action="store_true", help="count token-for-token swaps as buys")
        p.add_argument("--no-net-rule", action="store_true", help="do not require net_bought >= line")
        p.add_argument("--strict-bots", action="store_true",
                       help="exclude (not just flag) wallets over %d txs/day (Sep 12 bot rule)" % BOT_TRADES_PER_DAY)
        p.add_argument("--workers", type=int, default=4)
        p.add_argument("--cache", help="tx cache dir (default engine/.cache)")
        p.add_argument("--no-cache", action="store_true")
        p.add_argument("--verbose", action="store_true")

    p = sub.add_parser("check")
    p.add_argument("wallets", nargs="+")
    p.add_argument("--json", action="store_true")
    p.add_argument("--no-events", action="store_true")
    common(p)
    p = sub.add_parser("snapshot")
    p.add_argument("--wallets", required=True, help="file, one wallet per line (# comments ok)")
    p.add_argument("--out", required=True)
    p.add_argument("--no-events", action="store_true")
    common(p)
    p = sub.add_parser("candidates", help="list wallets that bought through a token's pool since start")
    p.add_argument("--token", required=True, help="symbol from config, e.g. PEPENOM")
    p.add_argument("--pool", help="pool address (default: from config)")
    p.add_argument("--min-bought", type=float, default=0, help="only wallets that bought at least this many tokens")
    p.add_argument("--out", required=True, help="writes OUT (one wallet per line) and OUT.json (amounts)")
    p.add_argument("--max-pages", type=int, default=200)
    common(p)
    p = sub.add_parser("draw")
    p.add_argument("--wallets", required=True, help="eligible.txt or any one-per-line file")
    p.add_argument("--seed")
    p.add_argument("--drand-round", type=int)
    p.add_argument("--winners", type=int, default=1)
    p.add_argument("--alternates", type=int, default=3)
    p.add_argument("--expect-sha256", help="refuse to draw if the list hash differs (pre-committed hash)")
    p = sub.add_parser("drand-round")
    p.add_argument("--at", required=True)
    p = sub.add_parser("hash")
    p.add_argument("files", nargs="+")
    return ap


def main(argv=None):
    args = build_parser().parse_args(argv)

    if args.cmd == "hash":
        for f in args.files:
            print("%s  %s" % (file_sha256(f), f))
        return 0
    if args.cmd == "drand-round":
        ts = parse_time(args.at)
        rnd = drand_round_at(ts)
        print(json.dumps({"at": iso(ts), "round": rnd, "round_time": iso(drand_time_of(rnd)),
                          "chain": DRAND_QUICKNET, "url": "%s/%s/public/%d" % (DRAND_URLS[0], DRAND_QUICKNET, rnd)}))
        return 0
    if args.cmd == "draw":
        wallets = read_wallet_file(args.wallets)
        bad = [w for w in wallets if not is_address(w)]
        if bad:
            print("invalid addresses: %s" % bad, file=sys.stderr)
            return 2
        h = list_sha256(wallets)
        if args.expect_sha256 and args.expect_sha256 != h:
            print("list sha256 %s != expected %s; refusing" % (h, args.expect_sha256), file=sys.stderr)
            return 3
        seed_info = {}
        if args.drand_round:
            try:
                d = drand_fetch(args.drand_round)
            except RpcError as e:
                print(str(e), file=sys.stderr)
                return 4
            seed = d["randomness"]
            seed_info = {"drand_round": d["round"], "drand_round_time": iso(drand_time_of(d["round"])),
                         "randomness": d["randomness"], "signature": d["signature"], "source": d["source"]}
        elif args.seed:
            seed = args.seed
        else:
            print("need --seed or --drand-round", file=sys.stderr)
            return 2
        order = draw_order(seed, wallets)
        n, a = args.winners, args.alternates
        print(json.dumps({"rule": "order = sha256(seed + wallet) hex ascending; first N win; next are alternates",
                          "seed": seed, **seed_info, "list_count": len(wallets), "list_sha256": h,
                          "winners": [w for _, w in order[:n]], "alternates": [w for _, w in order[n:n + a]],
                          "order": [{"rank": i + 1, "wallet": w, "hash": x} for i, (x, w) in enumerate(order)]},
                         indent=1))
        return 0

    cfg = load_config(args.config)
    apply_thresholds(cfg, args)
    excl_path = args.exclusions or os.path.join(HERE, cfg.get("exclusions", "exclusions.json"))
    excl, flag_programs, excl_sha = load_exclusions(excl_path)
    chain = build_chain(cfg, args)
    opts = make_opts(cfg, args)

    if args.cmd == "candidates":
        tok = next(t for t in cfg["tokens"] if t["symbol"].upper() == args.token.upper().lstrip("$"))
        pool = args.pool or tok["pool"]
        t0 = time.time()
        dec = chain.mint_info(tok["mint"])["decimals"]
        buyers, complete, nsig, missing = pool_buyers(chain, pool, tok["mint"], opts["start_ts"], opts["end_ts"],
                                                      args.max_pages)
        rows = sorted(({"wallet": w, "bought": b["bought_raw"] / 10 ** dec, "buys": b["buys"],
                        "first_buy": iso(b["first_buy"]), "listed": w in excl} for w, b in buyers.items()),
                      key=lambda r: -r["bought"])
        rows = [r for r in rows if r["bought"] >= args.min_bought]
        with open(args.out, "w") as f:
            f.write("\n".join(r["wallet"] for r in rows if not r["listed"]))
        with open(args.out + ".json", "w") as f:
            json.dump({"token": tok["symbol"], "pool": pool, "start": iso(opts["start_ts"]), "end": iso(opts["end_ts"]),
                       "pool_txs": nsig, "history_complete": complete, "missing_txs": missing,
                       "runtime_s": round(time.time() - t0, 1), "buyers": rows}, f, indent=1)
        print(json.dumps({"token": tok["symbol"], "pool_txs": nsig, "buyers": len(rows),
                          "unlisted_written": sum(1 for r in rows if not r["listed"]), "history_complete": complete,
                          "missing_txs": missing, "runtime_s": round(time.time() - t0, 1), "out": args.out}, indent=1))
        return 0

    if args.cmd == "check":
        results = []
        for w in args.wallets:
            try:
                r = check_wallet(chain, w.strip(), cfg, excl, flag_programs, opts)
            except RpcError as e:
                r = {"wallet": w, "error": str(e)}
            results.append(r)
            if not args.json:
                print(summarize(r), flush=True)
        if args.json:
            print(json.dumps({"engine_version": VERSION, "exclusions_sha256": excl_sha,
                              "config": {k: cfg[k] for k in ("tokens", "start")}, "end": iso(opts["end_ts"]),
                              "hold_end": iso(opts["hold_end_ts"]), "strict_bots": opts["strict_bots"],
                              "results": results}, indent=1))
        return 0

    if args.cmd == "snapshot":
        wallets = read_wallet_file(args.wallets)
        results = []
        t0 = time.time()
        for i, w in enumerate(wallets):
            try:
                r = check_wallet(chain, w, cfg, excl, flag_programs, opts)
            except RpcError as e:
                r = {"wallet": w, "error": str(e), "eligible": False}
            results.append(r)
            print("[%d/%d] %s eligible=%s excluded=%s %.1fs" % (i + 1, len(wallets), w, r.get("eligible"),
                                                                r.get("excluded"), r.get("runtime_s", 0)),
                  file=sys.stderr, flush=True)
        meta = {"engine_version": VERSION, "engine_sha256": file_sha256(os.path.abspath(__file__)),
                "generated_at": iso(int(time.time())), "start": iso(opts["start_ts"]), "end": iso(opts["end_ts"]),
                "hold_end": iso(opts["hold_end_ts"]),
                "config": {"tokens": cfg["tokens"], "net_rule": opts["net_rule"], "strict_bots": opts["strict_bots"],
                           "count_token_swaps": opts["count_token_swaps"], "hops": opts["hops"]},
                "exclusions_file": os.path.basename(excl_path), "exclusions_sha256": excl_sha,
                "wallets_in": len(wallets), "runtime_s": round(time.time() - t0, 1),
                "errors": [r["wallet"] for r in results if r.get("error")]}
        doc, sums = write_snapshot(results, meta, args.out)
        print(json.dumps({"eligible_count": doc["eligible_count"], "qualified_count": doc["qualified_count"],
                          "eligible_sha256": doc["eligible_sha256"],
                          "sha256sums": sums, "errors": meta["errors"]}, indent=1))
        return 0
    return 1


if __name__ == "__main__":
    sys.exit(main())
