#!/usr/bin/env python3
"""
Backfill de enriquecimento UTM - PX3 Lab
=========================================

Percorre contatos criados nos últimos N dias no Linkia PX3 e, pra cada um
que não tem custom field utm_source populado, aplica enrich_contact()
usando o mesmo mapeamento do serviço em produção.

Uso:
    # Dry-run (só simula, não escreve nada) — RECOMENDADO na 1ª execução
    python3 backfill_enrichment.py --days 90 --dry-run

    # Execução real (limita a 100 contatos p/ segurança)
    python3 backfill_enrichment.py --days 90 --limit 100

    # Sem limite (backfill completo)
    python3 backfill_enrichment.py --days 90 --limit 0

    # Só analisa/reporta distribuição (não escreve mesmo sem --dry-run)
    python3 backfill_enrichment.py --days 90 --analyze

Saída: log em /opt/mia/logs/ghl_utm_enricher_px3_backfill.log e resumo no stdout.
"""

from __future__ import annotations

import argparse
import json
import logging
import os
import sys
import time
from collections import Counter
from datetime import datetime, timedelta, timezone
from logging.handlers import RotatingFileHandler
from pathlib import Path

import requests

# Importa o core do serviço
sys.path.insert(0, str(Path(__file__).resolve().parent))
import app as service  # noqa: E402

# ---------------------------------------------------------------------------
LOG_FILE = Path("/opt/mia/logs/ghl_utm_enricher_px3_backfill.log")
LOG_FILE.parent.mkdir(parents=True, exist_ok=True)
_fmt = logging.Formatter("%(asctime)s [%(levelname)s] %(message)s")
_fh = RotatingFileHandler(LOG_FILE, maxBytes=10 * 1024 * 1024, backupCount=3, encoding="utf-8")
_fh.setFormatter(_fmt)
_sh = logging.StreamHandler(sys.stdout)
_sh.setFormatter(_fmt)
logging.basicConfig(level=logging.INFO, handlers=[_fh, _sh])
log = logging.getLogger("backfill")


def list_contacts_since(days: int, page_limit: int = 100):
    """
    Gera contatos criados nos últimos N dias, paginado.
    Usa POST /contacts/search com filtro dateAdded.
    """
    since = datetime.now(timezone.utc) - timedelta(days=days)
    since_iso = since.strftime("%Y-%m-%dT%H:%M:%S.000Z")

    url = f"{service.GHL_API_BASE}/contacts/search"
    payload = {
        "locationId": service.GHL_LOCATION_ID,
        "pageLimit": page_limit,
        "filters": [
            {"field": "dateAdded", "operator": "range", "value": {"gte": since_iso}},
        ],
        "sort": [{"field": "dateAdded", "direction": "desc"}],
    }

    total_seen = 0
    search_after = None
    while True:
        body = dict(payload)
        if search_after:
            body["searchAfter"] = search_after
        try:
            r = requests.post(url, headers=service._ghl_headers(), json=body, timeout=15)
        except requests.RequestException as e:
            log.error("search excecao: %s", e)
            return
        if r.status_code >= 300:
            log.error("search %s: %s", r.status_code, r.text[:400])
            return
        data = r.json() or {}
        contacts = data.get("contacts", []) or []
        if not contacts:
            return
        for c in contacts:
            total_seen += 1
            yield c
        last = contacts[-1]
        # cursor de paginação (GHL usa searchAfter no último item)
        search_after = last.get("searchAfter")
        if not search_after:
            # fallback: paginação por total
            if total_seen >= (data.get("total") or 0):
                return
            # se não veio searchAfter, evita loop infinito
            return
        time.sleep(0.3)  # rate limit friendly


def run_backfill(days: int, limit: int, dry_run: bool, analyze_only: bool):
    if dry_run:
        # Força DRY_RUN no módulo em runtime
        service.DRY_RUN = True
        log.info(">>> DRY_RUN ativado: nenhum PUT sera feito na GHL")

    log.info("backfill iniciando: days=%d limit=%d analyze_only=%s dry_run=%s",
             days, limit, analyze_only, dry_run)

    stats = Counter()
    processed = 0
    written = 0
    skipped_has = 0
    skipped_no_map = 0
    errored = 0

    written_examples = []
    skipped_examples = []

    for c in list_contacts_since(days=days):
        contact_id = c.get("id")
        if not contact_id:
            continue

        # A listagem NÃO retorna attributionSource — só o GET individual.
        # Então enrich_contact() faz o GET completo internamente.
        # Isso significa 1 request extra por contato — respeitar rate limit.

        if analyze_only:
            # Só busca detalhes e categoriza (não escreve mesmo se dry_run=False)
            full = service.fetch_contact(contact_id)
            if not full:
                stats["fetch_error"] += 1
                continue
            has_utm = service._get_cf_value(full, service.UTM_FIELD_SOURCE)
            att = full.get("attributionSource") or {}
            mapped = service._map_attribution_to_utm(att)
            if has_utm:
                stats["has_utm_source"] += 1
            elif mapped:
                stats[f"would_write:{mapped.get('utm_source','?')}"] += 1
            else:
                sess = (att.get("sessionSource") or "(none)").lower()
                stats[f"no_mapping:{sess}"] += 1
            processed += 1
        else:
            result = service.enrich_contact(contact_id, force=True)
            processed += 1
            if not result.get("ok"):
                errored += 1
                stats[f"error:{result.get('reason', 'unknown')}"] += 1
            elif result.get("skip") == "already_has_utm_source":
                skipped_has += 1
                stats["skip_already_has"] += 1
            elif result.get("skip"):
                skipped_no_map += 1
                stats[f"skip:{result.get('skip')}"] += 1
                if len(skipped_examples) < 3:
                    skipped_examples.append(result)
            else:
                written += 1
                stats["written"] += 1
                stats[f"written_source:{result.get('utm_written', {}).get('utm_source','?')}"] += 1
                if len(written_examples) < 5:
                    written_examples.append(result)

        if limit and processed >= limit:
            log.info("limite atingido (%d), parando", limit)
            break

        # Rate limit: pausa a cada 20 requests
        if processed % 20 == 0:
            log.info(
                "progresso: %d processados / writes=%d skip_has=%d skip_no_map=%d err=%d",
                processed, written, skipped_has, skipped_no_map, errored,
            )
            time.sleep(1.0)
        else:
            time.sleep(0.15)

    log.info("=" * 60)
    log.info("RESUMO backfill")
    log.info("=" * 60)
    log.info("processados: %d", processed)
    log.info("escritos:    %d", written)
    log.info("skip (já tem utm_source): %d", skipped_has)
    log.info("skip (sem mapeamento):    %d", skipped_no_map)
    log.info("erros: %d", errored)
    log.info("distribuição detalhada:")
    for k, v in sorted(stats.items(), key=lambda kv: -kv[1]):
        log.info("  %s: %d", k, v)
    if written_examples:
        log.info("EXEMPLOS de escrita:")
        for ex in written_examples:
            log.info("  %s", json.dumps(ex, ensure_ascii=False)[:300])
    if skipped_examples:
        log.info("EXEMPLOS de skip (sem mapping):")
        for ex in skipped_examples:
            log.info("  %s", json.dumps(ex, ensure_ascii=False)[:300])


def main():
    parser = argparse.ArgumentParser()
    parser.add_argument("--days", type=int, default=90)
    parser.add_argument("--limit", type=int, default=100, help="0 = sem limite")
    parser.add_argument("--dry-run", action="store_true", help="simula, não escreve")
    parser.add_argument("--analyze", action="store_true",
                        help="só lê e categoriza — nunca escreve")
    args = parser.parse_args()
    run_backfill(days=args.days, limit=args.limit, dry_run=args.dry_run, analyze_only=args.analyze)


if __name__ == "__main__":
    main()
