#!/usr/bin/env python3
"""Group-join reader: extract full log history from stealer exfil groups/channels.

Mechanism (VERIFIED 2026-07-27, docs/log-source-discovery.md §8.2):
If a stealer bot is admin in its exfil group/channel, its own token can call
exportChatInviteLink → our MTProto user session joins via the invite link →
get_chat_history reads the FULL history including the bot's own posts
(per-victim logs back to group creation — no 24h getUpdates window).

Inputs (in priority order):
  1. findings/data/tg_bots_state.json → discovered_chats (from my_chat_member)
  2. findings/data/tg_bots_validated.tsv → valid rows with chat_id set

OPSEC: joining a group is visible to its admins. Every join requires --yes
(operator confirmation) and is logged to findings/data/tg_group_joins.log.

Usage:
  .venv/bin/python scripts/tg_group_read.py --list            # show candidates, no action
  .venv/bin/python scripts/tg_group_read.py --yes             # process all candidates
  .venv/bin/python scripts/tg_group_read.py --chat -1001234567890 --token <tok> --yes
"""
import argparse
import asyncio
import csv
import json
import os
import sys
import time
import urllib.parse
import urllib.request
from datetime import datetime, timezone
from pathlib import Path

sys.path.insert(0, str(Path(__file__).parent))
from tg_session_load import build_session_file, load_dotenv, ENV_FILE, SESSION_DIR

PROJECT_DIR = Path(__file__).resolve().parent.parent
DATA_DIR = PROJECT_DIR / "findings" / "data"
BREACH_DIR = PROJECT_DIR / "findings" / "breaches"
IN_TSV = DATA_DIR / "tg_bots_validated.tsv"
STATE_JSON = DATA_DIR / "tg_bots_state.json"
JOIN_LOG = DATA_DIR / "tg_group_joins.log"

TG_API = "https://api.telegram.org/bot{token}/{method}"


def bot_api(token: str, method: str, **params) -> dict:
    url = TG_API.format(token=token, method=method)
    if params:
        url += "?" + urllib.parse.urlencode(params)
    try:
        with urllib.request.urlopen(url, timeout=30) as r:
            return json.loads(r.read())
    except Exception as e:
        return {"ok": False, "description": str(e)}


def now_iso() -> str:
    return datetime.now(timezone.utc).strftime("%Y-%m-%dT%H:%M:%SZ")


def log_join(line: str):
    DATA_DIR.mkdir(parents=True, exist_ok=True)
    with JOIN_LOG.open("a") as f:
        f.write(f"{now_iso()} {line}\n")


def collect_candidates() -> list[dict]:
    """[{token, bot_username, chat_id, chat_title, source}]"""
    out = []
    if STATE_JSON.exists():
        state = json.loads(STATE_JSON.read_text())
        for token, st in state.items():
            for cid, info in (st.get("discovered_chats") or {}).items():
                out.append({
                    "token": token, "bot_username": st.get("bot_username", ""),
                    "chat_id": cid, "chat_title": info.get("title") or "",
                    "source": "my_chat_member",
                })
    if IN_TSV.exists():
        with IN_TSV.open() as f:
            for row in csv.DictReader(f, delimiter="\t"):
                if row["status"] == "valid" and row.get("chat_id"):
                    out.append({
                        "token": row["token"],
                        "bot_username": row["bot_username"],
                        "chat_id": row["chat_id"], "chat_title": "",
                        "source": "ioc_url",
                    })
    # dedup by (token, chat_id)
    seen, dedup = set(), []
    for c in out:
        k = (c["token"], c["chat_id"])
        if k not in seen:
            seen.add(k)
            dedup.append(c)
    return dedup


async def read_group(app, token: str, chat_id: str, limit: int) -> dict:
    """getChat → exportChatInviteLink → join → dump history. Returns stats."""
    info = bot_api(token, "getChat", chat_id=chat_id)
    if not info.get("ok"):
        return {"error": f"getChat: {info.get('description')}"}
    chat = info["result"]
    ctype = chat.get("type")
    title = chat.get("title") or chat.get("first_name")
    print(f"    chat {chat_id}: type={ctype} title={title!r}")

    link = bot_api(token, "exportChatInviteLink", chat_id=chat_id)
    if not link.get("ok"):
        return {"error": f"exportChatInviteLink: {link.get('description')}",
                "chat_type": ctype, "title": title}
    invite = link["result"]
    print(f"    invite link: {invite}")
    log_join(f"JOIN chat_id={chat_id} title={title!r} link={invite}")

    joined = await app.join_chat(invite)
    print(f"    joined as user, chat id={joined.id}")
    time.sleep(2)

    bot_dir = BREACH_DIR / f"TGBOT-{token.split(':')[0]}"
    hist_dir = bot_dir / f"history_{chat_id}"
    hist_dir.mkdir(parents=True, exist_ok=True)
    n_msgs, n_docs = 0, 0
    with (hist_dir / "messages.jsonl").open("w") as f:
        async for msg in app.get_chat_history(joined.id, limit=limit):
            rec = {
                "id": msg.id, "date": str(msg.date),
                "from_bot": bool(msg.from_user and msg.from_user.is_bot),
                "from_id": msg.from_user.id if msg.from_user else None,
                "text": msg.text or msg.caption or "",
            }
            if msg.document:
                rec["document"] = {
                    "file_name": msg.document.file_name,
                    "file_size": msg.document.file_size,
                }
                fname = msg.document.file_name or f"doc_{msg.id}"
                dest = hist_dir / f"{msg.id}_{fname}"
                try:
                    await app.download_media(msg, file_name=str(dest))
                    n_docs += 1
                except Exception as e:
                    rec["download_error"] = str(e)
            f.write(json.dumps(rec, ensure_ascii=False) + "\n")
            n_msgs += 1
    return {"chat_type": ctype, "title": title, "messages": n_msgs,
            "documents": n_docs, "out": str(hist_dir)}


async def amain(args):
    candidates = collect_candidates()
    if args.chat and args.token:
        candidates = [{"token": args.token, "bot_username": "manual",
                       "chat_id": args.chat, "chat_title": "",
                       "source": "cli"}]
    if not candidates:
        print("[*] no group candidates (no chat_id in TSV, no discovered_chats)")
        return
    print(f"[*] {len(candidates)} candidate(s):")
    for c in candidates:
        print(f"    @{c['bot_username'] or '?'} chat_id={c['chat_id']} "
              f"title={c['chat_title']!r} src={c['source']}")
    if not args.yes:
        print("\n[*] dry-run. OPSEC: join is visible to group admins. "
              "Re-run with --yes to proceed (operator confirmation).")
        return

    load_dotenv(str(ENV_FILE))
    api_id = int(os.environ["TG_API_ID"])
    api_hash = os.environ["TG_API_HASH"]
    build_session_file(api_id, os.environ["TG_AUTH_KEY"].strip(),
                       int(os.environ["TG_DC_ID"]), int(os.environ["TG_USER_ID"]))
    from pyrogram import Client
    app = Client(name="breach_session", api_id=api_id, api_hash=api_hash,
                 workdir=str(SESSION_DIR))
    await app.start()
    try:
        for c in candidates:
            print(f"[*] @{c['bot_username'] or '?'} → chat {c['chat_id']}")
            r = await read_group(app, c["token"], c["chat_id"], args.limit)
            if "error" in r:
                print(f"    [skip] {r['error']}")
                log_join(f"FAIL chat_id={c['chat_id']} err={r['error']}")
            else:
                print(f"    [+] {r['messages']} messages, {r['documents']} "
                      f"documents → {r['out']}")
                log_join(f"OK chat_id={c['chat_id']} msgs={r['messages']} "
                         f"docs={r['documents']}")
    finally:
        await app.stop()


def main():
    ap = argparse.ArgumentParser(description=__doc__.splitlines()[0])
    ap.add_argument("--list", action="store_true", help="list candidates, no action")
    ap.add_argument("--yes", action="store_true",
                    help="operator confirmation: actually join groups (visible to admins)")
    ap.add_argument("--chat", help="manual chat_id (with --token)")
    ap.add_argument("--token", help="manual bot token (with --chat)")
    ap.add_argument("--limit", type=int, default=1000,
                    help="max messages to read per group")
    args = ap.parse_args()
    asyncio.run(amain(args))


if __name__ == "__main__":
    main()
