#!/usr/bin/env python3
"""Local breach archive intake (no SMB, no network to fetch).

Processes a local .rar/.zip archive: test → extract → aggregate → corp_search → index.
Idempotent: re-running on an already-processed archive is safe.

Usage:
  # One archive (password auto-guessed from archive name for known channels)
  python3 intake_local.py --archive "PATH/TO/6226_05.07.2026_@PIXELCLOUD3.rar"

  # Explicit password
  python3 intake_local.py --archive "PATH/TO/NAME.rar" --password "@PIXELCLOUD3"

  # Archive already extracted to a dir — skip test/extract, just aggregate+search
  python3 intake_local.py --extracted "PATH/TO/EXTRACTED_DIR" --source SOURCE_ID

Pipeline:
  1. Test archive integrity (unrar t / 7z t)
  2. Extract to findings/breaches/<SOURCE_ID>/
  3. aggregate_breach.py → findings/data/<SOURCE_ID>.tsv
  4. corp_search.py --source <SOURCE_ID> → findings/data/corp_search_<SOURCE_ID>_results.md
  5. Append row to findings/data/index.csv (idempotent)

Source ID derivation: strip date + channel suffix from archive name.
  "6226_05.07.2026_@PIXELCLOUD3.rar"  → PIXELCLOUD3-6226
  "@beetraffic 2000 MIX 05-07-2026"   → BEETRAFFIC-2000
  "AYANKOUJI PRIVATE #528 PART 4"     → AYANKOUJI-528
"""
import argparse
import csv
import re
import subprocess
import sys
from pathlib import Path

SCRIPT_DIR = Path(__file__).parent
PROJECT_DIR = SCRIPT_DIR.parent
BREACH_DIR = PROJECT_DIR / "findings" / "breaches"
DATA_DIR = PROJECT_DIR / "findings" / "data"
INDEX_CSV = DATA_DIR / "index.csv"

KNOWN_PASSWORDS = {
    "PIXELCLOUD3": "@PIXELCLOUD3",
    "UP_DAISYCLOUD": "@UP_DAISYCLOUD",
    "kir3info": "@ScroogeUrl",
    "fatetraffic": "@fatetraffic",
    "valenciga": "@VALENCIGA",
    "watercloud": "@watercloud_info",
    "watercloud_info": "@watercloud_info",
    "beetraffic": "https://t.me/+jglATzWg88A4ODg8",
    "ghostcloud": "https://t.me/GhostCloud_info",
}


def derive_source_id(archive_name: str) -> str:
    name = archive_name
    for ext in (".rar", ".zip", ".7z", ".tar", ".gz"):
        if name.lower().endswith(ext):
            name = name[:-len(ext)]
    m = re.match(r"^(\d+)_[\d.]+_@([\w]+)$", name)
    if m:
        return f"{m.group(2).upper()}-{m.group(1)}"
    m = re.match(r"^@?(\w+)\s+(\d+)", name)
    if m:
        return f"{m.group(1).upper()}-{m.group(2)}"
    m = re.match(r"^@?(UP_DAISYCLOUD)[^_]*_.*?_(\d+)_ON_CHANNEL$", name, re.I)
    if m:
        return f"{m.group(1).upper()}-{m.group(2)}"
    m = re.match(r"^@?(\w+)\s*-\s*.*?#(\d+)", name, re.I)
    if m:
        return f"{m.group(1).upper()}-{m.group(2)}"
    m = re.match(r"^(\w+)\s+PRIVATE\s+#(\d+)", name, re.I)
    if m:
        return f"{m.group(1).upper()}-{m.group(2)}"
    m = re.match(r"^logX\s+[—-]\s+([\d.]+)", name, re.I)
    if m:
        return f"LOGX-{m.group(1).replace('.', '-')}"
    return re.sub(r"[^A-Za-z0-9]+", "-", name).strip("-").upper()


def guess_password(archive_name: str) -> str | None:
    name_lower = archive_name.lower()
    for channel, pw in KNOWN_PASSWORDS.items():
        if channel.lower() in name_lower:
            return pw
    return None


def test_archive(path: Path, password: str | None) -> bool:
    if path.suffix.lower() == ".rar":
        cmd = ["unrar", "t", "-y"]
        cmd.append(f"-p{password}" if password else "-p-")
        cmd.append(str(path))
        r = subprocess.run(cmd, capture_output=True, text=True, timeout=600, stdin=subprocess.DEVNULL)
        if "Incorrect" in r.stdout or "Incorrect" in r.stderr:
            return False
        return r.returncode == 0
    elif path.suffix.lower() == ".zip":
        cmd = ["7z", "t", str(path), "-y"]
        if password:
            cmd.insert(-2, f"-p{password}")
        r = subprocess.run(cmd, capture_output=True, text=True, timeout=600, stdin=subprocess.DEVNULL)
        return r.returncode == 0 and "Everything is Ok" in r.stdout
    return False


def extract_archive(path: Path, dest_dir: Path, password: str | None) -> bool:
    dest_dir.mkdir(parents=True, exist_ok=True)
    if path.suffix.lower() == ".rar":
        cmd = ["unrar", "x", "-y", "-o+"]
        if password:
            cmd.append(f"-p{password}")
        cmd += [str(path), f"{dest_dir}/"]
    elif path.suffix.lower() == ".zip":
        cmd = ["7z", "x", str(path), f"-o{dest_dir}", "-y"]
        if password:
            cmd.insert(-1, f"-p{password}")
    else:
        return False
    r = subprocess.run(cmd, capture_output=True, text=True, timeout=1800, stdin=subprocess.DEVNULL)
    return r.returncode == 0


def aggregate(source_id: str, breach_dir: Path) -> bool:
    out_tsv = DATA_DIR / f"{source_id}.tsv"
    if out_tsv.exists() and out_tsv.stat().st_size > 0:
        print(f"  already aggregated: {out_tsv}")
        return True
    r = subprocess.run(
        ["python3", str(SCRIPT_DIR / "aggregate_breach.py"), str(breach_dir), str(out_tsv)],
        capture_output=True, text=True, timeout=600,
    )
    print(r.stdout.strip())
    if r.returncode != 0:
        print(f"  [error] aggregate failed: {r.stderr[:300]}", file=sys.stderr)
        return False
    return out_tsv.exists()


def run_corp_search(source_id: str) -> bool:
    out_md = DATA_DIR / f"corp_search_{source_id}_results.md"
    tsv_path = DATA_DIR / f"{source_id}.tsv"
    r = subprocess.run(
        ["python3", str(SCRIPT_DIR / "corp_search.py"), "--tsv", str(tsv_path)],
        capture_output=True, text=True, timeout=300,
    )
    if r.returncode != 0:
        print(f"  [error] corp_search failed: {r.stderr[:300]}", file=sys.stderr)
        return False
    header = f"# Corporate Access Search Results — {source_id}\n**Source:** findings/data/{source_id}.tsv\n\n"
    body = r.stdout
    if body.startswith("# Corporate Access Search Results"):
        body = body.split("\n", 2)[2] if "\n" in body else body
    out_md.write_text(header + body, encoding="utf-8")
    return True


def update_index(source_id: str, breach_dir: Path, archive_name: str | None) -> bool:
    rows = []
    fieldnames = None
    if INDEX_CSV.exists():
        with open(INDEX_CSV, encoding="utf-8") as f:
            reader = csv.DictReader(f)
            rows = list(reader)
            fieldnames = reader.fieldnames
    if not fieldnames:
        fieldnames = ["dataset_id","path","kind","source","format","producer","consumers","status","notes"]
    dataset_id = f"{source_id.lower()}_tsv"
    if any(r.get("dataset_id") == dataset_id for r in rows):
        print(f"  index already has {dataset_id}")
        return True
    new_rows = [
        {"dataset_id": dataset_id,
         "path": f"findings/data/{source_id}.tsv",
         "kind": "derived", "source": "breach_archive", "format": "tsv",
         "producer": "scripts/aggregate_breach.py",
         "consumers": f"findings/data/corp_search_{source_id}_results.md",
         "status": "active",
         "notes": f"Aggregated credentials from {archive_name or breach_dir.name}"},
        {"dataset_id": f"{source_id.lower()}_corp_search",
         "path": f"findings/data/corp_search_{source_id}_results.md",
         "kind": "derived", "source": "breach_archive", "format": "md",
         "producer": "scripts/corp_search.py",
         "consumers": f"findings/data/{source_id}.tsv",
         "status": "active",
         "notes": f"Corp search output for {source_id}"},
    ]
    rows.extend(new_rows)
    DATA_DIR.mkdir(parents=True, exist_ok=True)
    with open(INDEX_CSV, "w", encoding="utf-8", newline="") as f:
        writer = csv.DictWriter(f, fieldnames=fieldnames)
        writer.writeheader()
        for r in rows:
            writer.writerow(r)
    print(f"  index +{len(new_rows)} rows")
    return True


def process_archive(archive_path: Path, password: str | None = None, skip_if_done: bool = True) -> bool:
    archive_name = archive_path.name
    print(f"\n{'='*70}\n# Processing: {archive_name}\n{'='*70}")
    source_id = derive_source_id(archive_name)
    print(f"  source_id: {source_id}")
    if not password:
        password = guess_password(archive_name)
        print(f"  guessed password: {password}" if password else "  no password guessed")

    print(f"  [1/5] testing archive integrity...")
    if not test_archive(archive_path, password):
        print(f"  [error] archive test failed (wrong password or corrupt)", file=sys.stderr)
        return False
    print("  OK")

    breach_dir = BREACH_DIR / source_id
    print(f"  [2/5] extracting to {breach_dir}...")
    if breach_dir.exists() and any(breach_dir.iterdir()) and skip_if_done:
        print(f"  already extracted ({sum(1 for _ in breach_dir.iterdir())} entries)")
    else:
        if not extract_archive(archive_path, breach_dir, password):
            print(f"  [error] extract failed", file=sys.stderr)
            return False

    print(f"  [3/5] aggregating credentials...")
    if not aggregate(source_id, breach_dir):
        return False
    print(f"  [4/5] running corp_search --source {source_id}...")
    if not run_corp_search(source_id):
        return False
    print(f"  [5/5] updating index.csv...")
    update_index(source_id, breach_dir, archive_name)
    print(f"\n  done: {source_id}")
    return True


def process_extracted(breach_dir: Path, source_id: str) -> bool:
    print(f"\n{'='*70}\n# Processing extracted dir: {breach_dir}\n{'='*70}")
    print(f"  source_id: {source_id}")
    print(f"  [1/3] aggregating credentials...")
    if not aggregate(source_id, breach_dir):
        return False
    print(f"  [2/3] running corp_search --source {source_id}...")
    if not run_corp_search(source_id):
        return False
    print(f"  [3/3] updating index.csv...")
    update_index(source_id, breach_dir, None)
    print(f"\n  done: {source_id}")
    return True


def main():
    p = argparse.ArgumentParser(description="Local breach archive intake (test → extract → aggregate → corp_search → index)")
    p.add_argument("--archive", "-a", help="Path to local .rar/.zip archive")
    p.add_argument("--password", "-p", help="Archive password (auto-guessed for known channels)")
    p.add_argument("--extracted", help="Path to already-extracted breach dir (skip test/extract)")
    p.add_argument("--source", help="Source ID (required with --extracted, derived from --archive otherwise)")
    p.add_argument("--force", action="store_true", help="Re-process even if outputs exist")
    args = p.parse_args()

    if args.extracted:
        breach_dir = Path(args.extracted)
        if not breach_dir.is_dir():
            print(f"[error] {breach_dir} is not a directory", file=sys.stderr); sys.exit(1)
        sid = args.source or derive_source_id(breach_dir.name)
        process_extracted(breach_dir, sid)
    elif args.archive:
        archive_path = Path(args.archive)
        if not archive_path.is_file():
            print(f"[error] {archive_path} not found", file=sys.stderr); sys.exit(1)
        process_archive(archive_path, args.password, skip_if_done=not args.force)
    else:
        p.error("specify --archive or --extracted")


if __name__ == "__main__":
    main()
