#!/usr/bin/env python3
"""
Importacao RD Station -> CRM Linkia (GoHighLevel white-label)
Sub-account: Francisco Borrello
Volume: 25.426 leads
Data: 2026-08-28

Regras:
  - Cria custom fields na location antes de importar
  - Cria SOMENTE contatos (zero oportunidade, zero pipeline)
  - Tag global: rd-import-2026-08-28
  - Tag por estagio: rd-lead | rd-cliente | rd-lead-qualificado
  - Dedupe por email (CREATE se novo, UPDATE se existe)
  - Checkpoint: resume de onde parou se interrompido
  - Backoff exponencial em 429
  - Log a cada 500 processados com ETA
"""

import csv
import io
import json
import logging
import os
import random
import subprocess
import sys
import time
from collections import defaultdict
from datetime import datetime, timezone
from pathlib import Path

import requests

# ---------------------------------------------------------------------------
# Configuracao
# ---------------------------------------------------------------------------
BASE_DIR = Path("/opt/mia/workspace/clientes/borrello/importacao_rd_20260828")
CSV_PATH = BASE_DIR / "rd_export.csv"
LOG_PATH = BASE_DIR / "import.log"
CHECKPOINT_PATH = BASE_DIR / "checkpoint.json"
ERRORS_PATH = BASE_DIR / "erros.jsonl"
SUMMARY_PATH = BASE_DIR / "resumo_final.json"
ENV_PATH = Path("/opt/mia/config/linkia_borrello.env")

API_BASE = "https://services.leadconnectorhq.com"
API_VERSION = "2021-07-28"
USER_AGENT = (
    "Mozilla/5.0 (X11; Linux x86_64) AppleWebKit/537.36 "
    "(KHTML, like Gecko) Chrome/124.0.0.0 Safari/537.36"
)

TAG_GLOBAL = "rd-import-2026-08-28"
TAG_ESTAGIO = {
    "lead": "rd-lead",
    "cliente": "rd-cliente",
    "lead qualificado": "rd-lead-qualificado",
}

LOG_INTERVAL = 500
# ~5 req/lead (search + upsert + margem) => 10 req/s limite => 0.12s delay conservador
DELAY_BETWEEN_REQS = 0.12
MAX_BACKOFF = 60

# ---------------------------------------------------------------------------
# Logging
# ---------------------------------------------------------------------------
logging.basicConfig(
    level=logging.INFO,
    format="%(asctime)s [%(levelname)s] %(message)s",
    handlers=[
        logging.FileHandler(LOG_PATH, encoding="utf-8"),
        logging.StreamHandler(sys.stdout),
    ],
)
log = logging.getLogger(__name__)


# ---------------------------------------------------------------------------
# Carregar env
# ---------------------------------------------------------------------------
def load_env(path: Path) -> dict:
    env = {}
    with open(path) as f:
        for line in f:
            line = line.strip()
            if not line or line.startswith("#"):
                continue
            if "=" in line:
                k, v = line.split("=", 1)
                env[k.strip()] = v.strip()
    return env


# ---------------------------------------------------------------------------
# Sessao HTTP com headers padrao
# ---------------------------------------------------------------------------
def make_session(token: str) -> requests.Session:
    s = requests.Session()
    s.headers.update({
        "Authorization": f"Bearer {token}",
        "Version": API_VERSION,
        "Content-Type": "application/json",
        "User-Agent": USER_AGENT,
        "Accept": "application/json",
    })
    return s


# ---------------------------------------------------------------------------
# Request com retry/backoff
# ---------------------------------------------------------------------------
def api_request(session: requests.Session, method: str, url: str, **kwargs):
    backoff = 1
    for attempt in range(10):
        try:
            resp = session.request(method, url, timeout=30, **kwargs)
            if resp.status_code == 429:
                wait = min(backoff, MAX_BACKOFF)
                log.warning(f"429 rate-limit - aguardando {wait}s (tentativa {attempt+1})")
                time.sleep(wait)
                backoff = min(backoff * 2, MAX_BACKOFF)
                continue
            return resp
        except requests.exceptions.RequestException as e:
            wait = min(backoff, MAX_BACKOFF)
            log.warning(f"Erro de rede: {e} - aguardando {wait}s")
            time.sleep(wait)
            backoff = min(backoff * 2, MAX_BACKOFF)
    raise RuntimeError(f"Falhou apos 10 tentativas: {method} {url}")


# ---------------------------------------------------------------------------
# Custom Fields: criar na location se nao existir
# ---------------------------------------------------------------------------
CUSTOM_FIELDS_SPEC = [
    ("rd_estagio_funil",                "DROPDOWN",   "Estagio no Funil (RD)",
     ["Lead", "Cliente", "Lead Qualificado"]),
    ("rd_cargo",                        "TEXT",        "Cargo (RD)", []),
    ("rd_biografia",                    "LARGE_TEXT",  "Biografia (RD)", []),
    ("rd_dono_lead",                    "TEXT",        "Dono do Lead (RD)", []),
    ("rd_data_ultima_oportunidade",     "DATE",        "Data Ultima Oportunidade (RD)", []),
    ("rd_data_ultima_venda",            "DATE",        "Data Ultima Venda (RD)", []),
    ("rd_valor_ultima_venda",           "MONETARY",    "Valor Ultima Venda (RD)", []),
    ("rd_lead_scoring_perfil",          "TEXT",        "Lead Scoring Perfil (RD)", []),
    ("rd_lead_scoring_interesse",       "TEXT",        "Lead Scoring Interesse (RD)", []),
    ("rd_status_comunicacao_email",     "TEXT",        "Status Comunicacao Email (RD)", []),
    ("rd_url_publica",                  "TEXT",        "URL Publica (RD)", []),
    ("rd_base_legal_comunicacao",       "TEXT",        "Base Legal Comunicacao (RD)", []),
    ("rd_total_conversoes",             "NUMERICAL",   "Total de Conversoes (RD)", []),
    ("rd_data_primeira_conversao",      "DATE",        "Data Primeira Conversao (RD)", []),
    ("rd_origem_primeira_conversao",    "TEXT",        "Origem Primeira Conversao (RD)", []),
    ("rd_data_ultima_conversao",        "DATE",        "Data Ultima Conversao (RD)", []),
    ("rd_origem_ultima_conversao",      "TEXT",        "Origem Ultima Conversao (RD)", []),
    ("rd_eventos_ultimos_100",          "LARGE_TEXT",  "Eventos Ultimos 100 (RD)", []),
    ("rd_ddd",                          "TEXT",        "DDD (RD)", []),
    ("rd_areas_perdeu_oportunidade",    "LARGE_TEXT",  "Areas que Perdeu Oportunidade (RD)", []),
    ("rd_escolha_recebimento",          "LARGE_TEXT",  "Escolha Recebimento (RD)", []),
    ("rd_link_codigo_barras",           "TEXT",        "Link Codigo de Barras (RD)", []),
    ("rd_curso_interesse",              "TEXT",        "Curso de Interesse (RD)", []),
    ("rd_utm_campaign",                 "TEXT",        "UTM Campaign (RD)", []),
    ("rd_utm_content",                  "TEXT",        "UTM Content (RD)", []),
    ("rd_utm_medium",                   "TEXT",        "UTM Medium (RD)", []),
    ("rd_utm_source",                   "TEXT",        "UTM Source (RD)", []),
    ("rd_utm_term",                     "TEXT",        "UTM Term (RD)", []),
    ("rd_confirma_participacao_evento", "TEXT",        "Confirma Participacao Evento (RD)", []),
    ("rd_facebook",                     "TEXT",        "Facebook (RD)", []),
    ("rd_twitter",                      "TEXT",        "Twitter (RD)", []),
    ("rd_linkedin",                     "TEXT",        "LinkedIn (RD)", []),
]


def ensure_custom_fields(session: requests.Session, location_id: str) -> dict:
    """
    Garante que todos os custom fields existam na location.
    Retorna dict {rd_key: field_id}.
    """
    log.info("Carregando custom fields existentes na location...")
    resp = api_request(session, "GET",
                       f"{API_BASE}/locations/{location_id}/customFields")
    if not resp.ok:
        raise RuntimeError(
            f"Erro ao listar custom fields: {resp.status_code} {resp.text[:300]}"
        )

    existing = {}
    data = resp.json()
    fields_list = data if isinstance(data, list) else data.get("customFields", [])
    for f in fields_list:
        key = f.get("fieldKey", "")
        # GHL prefixa com "contact." - normaliza
        short_key = key.replace("contact.", "")
        existing[short_key] = f["id"]

    log.info(f"Custom fields ja existentes: {len(existing)}")

    field_id_map = {}
    for (rd_key, data_type, label, options) in CUSTOM_FIELDS_SPEC:
        if rd_key in existing:
            field_id_map[rd_key] = existing[rd_key]
            log.info(f"  [ok] {rd_key} ja existe ({existing[rd_key]})")
            continue

        payload = {
            "name": label,
            "dataType": data_type,
        }
        if options:
            payload["options"] = [{"label": o, "value": o.lower()} for o in options]

        log.info(f"  Criando custom field: {rd_key} ({data_type})...")
        resp = api_request(
            session, "POST",
            f"{API_BASE}/locations/{location_id}/customFields",
            json=payload,
        )
        time.sleep(0.4)

        if resp.ok:
            created = resp.json()
            cf = created.get("customField", created)
            fid = cf.get("id", "")
            field_id_map[rd_key] = fid
            log.info(f"  [criado] {rd_key} -> {fid}")
        else:
            log.error(f"  [ERRO] {rd_key}: {resp.status_code} {resp.text[:200]}")
            # continua sem o campo; nao bloqueia a importacao

    return field_id_map


# ---------------------------------------------------------------------------
# Helpers de parsing
# ---------------------------------------------------------------------------

def clean(val) -> str:
    if not val:
        return ""
    v = str(val).strip().strip('"').strip()
    return v if v.lower() not in ("", '""', "none") else ""


def parse_name(nome: str):
    parts = clean(nome).split(" ", 1)
    first = parts[0] if parts else ""
    last = parts[1] if len(parts) > 1 else ""
    return first, last


def parse_phone(celular: str, telefone: str, whatsapp: str) -> str:
    for val in [celular, telefone, whatsapp]:
        v = clean(val)
        if v:
            return v
    return ""


def parse_date(val: str) -> str:
    """Converte '2026-08-23 19:04:20 -0300' -> 'YYYY-MM-DD'"""
    v = clean(val)
    if not v:
        return ""
    try:
        return v.split(" ")[0]
    except Exception:
        return ""


def parse_monetary(val: str):
    v = clean(val)
    if not v or v.lower() in ("d", "0", ""):
        return None
    try:
        return float(v.replace(",", ".").replace("R$", "").strip())
    except Exception:
        return None


def parse_int(val: str):
    v = clean(val)
    if not v:
        return None
    try:
        return int(float(v))
    except Exception:
        return None


def build_tags(row: dict) -> list:
    tags = [TAG_GLOBAL]

    estagio = clean(row.get("Estagio no funil", row.get("Estágio no funil", ""))).lower()
    tag_e = TAG_ESTAGIO.get(estagio)
    if tag_e:
        tags.append(tag_e)

    rd_tags_raw = clean(row.get("Tags", ""))
    if rd_tags_raw:
        for t in rd_tags_raw.split(","):
            t = t.strip()
            if t:
                t_clean = t.lower().replace(" ", "-")
                # GHL nao aceita tags muito longas; trunca em 100 chars
                if 0 < len(t_clean) <= 100:
                    tags.append(t_clean)

    return list(dict.fromkeys(tags))  # dedupe preservando ordem


def build_custom_values(row: dict, field_id_map: dict) -> list:
    """Monta lista de {id, value} para customFields no payload do contato."""

    # Chaves PT-BR exatas do CSV
    def col(name):
        # tenta chave exata primeiro, depois variantes sem acento
        return row.get(name, "")

    mapping = [
        ("rd_estagio_funil",
         clean(col("Estágio no funil"))),
        ("rd_cargo",
         clean(col("Cargo"))),
        ("rd_biografia",
         clean(col("Biografia"))),
        ("rd_dono_lead",
         clean(col("Dono do Lead"))),
        ("rd_data_ultima_oportunidade",
         parse_date(col("Data da última oportunidade"))),
        ("rd_data_ultima_venda",
         parse_date(col("Data da última venda"))),
        ("rd_valor_ultima_venda",
         parse_monetary(col("Valor da última venda"))),
        ("rd_lead_scoring_perfil",
         clean(col("Lead Scoring - Perfil"))),
        ("rd_lead_scoring_interesse",
         clean(col("Lead Scoring - Interesse"))),
        ("rd_status_comunicacao_email",
         clean(col("Status para comunicação por email"))),
        ("rd_url_publica",
         clean(col("URL pública"))),
        ("rd_base_legal_comunicacao",
         clean(col("Base legal para comunicação"))),
        ("rd_total_conversoes",
         parse_int(col("Total de conversões"))),
        ("rd_data_primeira_conversao",
         parse_date(col("Data da primeira conversão"))),
        ("rd_origem_primeira_conversao",
         clean(col("Origem da primeira conversão"))),
        ("rd_data_ultima_conversao",
         parse_date(col("Data da última conversão"))),
        ("rd_origem_ultima_conversao",
         clean(col("Origem da última conversão"))),
        ("rd_eventos_ultimos_100",
         clean(col("Eventos (Últimos 100)"))[:4000]),
        ("rd_ddd",
         clean(col("DDD"))),
        ("rd_areas_perdeu_oportunidade",
         clean(col("Em que áreas você perdeu alguma oportunidade?"))[:4000]),
        ("rd_escolha_recebimento",
         clean(col("Escolha o que você quer receber (pode marcar mais que uma opção)"))[:2000]),
        ("rd_link_codigo_barras",
         clean(col("Link Código de Barras"))),
        ("rd_curso_interesse",
         clean(col("Qual curso de interesse"))),
        ("rd_utm_campaign",
         clean(col("utm_campaign"))),
        ("rd_utm_content",
         clean(col("utm_content"))),
        ("rd_utm_medium",
         clean(col("utm_medium"))),
        ("rd_utm_source",
         clean(col("utm_source"))),
        ("rd_utm_term",
         clean(col("utm_term"))),
        ("rd_confirma_participacao_evento",
         clean(col("Você confirma sua participação no evento?"))),
        ("rd_facebook",
         clean(col("Facebook"))),
        ("rd_twitter",
         clean(col("Twitter"))),
        ("rd_linkedin",
         clean(col("Linkedin"))),
    ]

    result = []
    for key, value in mapping:
        fid = field_id_map.get(key)
        if not fid:
            continue
        if value is None or value == "":
            continue
        # MONETARY aceita float; o resto vai como string
        if isinstance(value, float):
            result.append({"id": fid, "value": value})
        else:
            result.append({"id": fid, "value": str(value)})

    return result


def build_contact_payload(row: dict, field_id_map: dict, location_id: str) -> dict:
    first, last = parse_name(row.get("Nome", ""))
    phone = parse_phone(
        row.get("Celular", ""),
        row.get("Telefone", ""),
        row.get("Whatsapp", ""),
    )
    tags = build_tags(row)
    custom_values = build_custom_values(row, field_id_map)

    payload = {
        "locationId": location_id,
        "firstName": first,
        "lastName": last,
        "email": clean(row.get("Email", "")),
        "tags": tags,
        "customFields": custom_values,
    }

    if phone:
        payload["phone"] = phone

    country = clean(row.get("País", row.get("Pais", "")))
    if country:
        payload["country"] = country

    state = clean(row.get("Estado", ""))
    if state:
        payload["state"] = state

    city = clean(row.get("Cidade", ""))
    if city:
        payload["city"] = city

    dob = parse_date(row.get("Data de aniversário", row.get("Data de aniversario", "")))
    if dob:
        payload["dateOfBirth"] = dob

    website = clean(row.get("Website", ""))
    if website:
        payload["website"] = website

    company = clean(row.get("Empresa", ""))
    if company:
        payload["companyName"] = company

    return payload


# ---------------------------------------------------------------------------
# Buscar contato por email
# Endpoint confirmado: GET /contacts/?locationId=X&query=EMAIL&limit=1
# ---------------------------------------------------------------------------
def find_contact_by_email(
    session: requests.Session, location_id: str, email: str
) -> str | None:
    resp = api_request(
        session, "GET",
        f"{API_BASE}/contacts/",
        params={"locationId": location_id, "query": email, "limit": 5},
    )
    if not resp.ok:
        return None
    data = resp.json()
    contacts = data.get("contacts", [])
    for c in contacts:
        if c.get("email", "").lower() == email.lower():
            return c["id"]
    return None


# ---------------------------------------------------------------------------
# Criar ou atualizar contato
# ---------------------------------------------------------------------------
def upsert_contact(
    session: requests.Session,
    location_id: str,
    payload: dict,
    existing_id: str | None,
) -> tuple:
    """Returns (contact_id, action) where action in 'created'|'updated'|'error:...'"""
    if existing_id:
        resp = api_request(
            session, "PUT",
            f"{API_BASE}/contacts/{existing_id}",
            json=payload,
        )
        if resp.ok:
            return existing_id, "updated"
        else:
            return existing_id, f"error:{resp.status_code}:{resp.text[:200]}"
    else:
        resp = api_request(
            session, "POST",
            f"{API_BASE}/contacts/",
            json=payload,
        )
        if resp.ok:
            data = resp.json()
            contact = data.get("contact", data)
            return contact.get("id", ""), "created"
        elif resp.status_code == 400:
            # GHL as vezes detecta duplicata e devolve contactId no corpo
            try:
                err_data = resp.json()
                cid = err_data.get("contactId") or err_data.get("id", "")
                if cid:
                    resp2 = api_request(
                        session, "PUT",
                        f"{API_BASE}/contacts/{cid}",
                        json=payload,
                    )
                    if resp2.ok:
                        return cid, "updated"
            except Exception:
                pass
            return "", f"error:400:{resp.text[:200]}"
        else:
            return "", f"error:{resp.status_code}:{resp.text[:200]}"


# ---------------------------------------------------------------------------
# Checkpoint
# ---------------------------------------------------------------------------
def load_checkpoint() -> dict:
    if CHECKPOINT_PATH.exists():
        try:
            with open(CHECKPOINT_PATH) as f:
                return json.load(f)
        except Exception:
            pass
    return {}


def save_checkpoint(ckpt: dict):
    tmp = str(CHECKPOINT_PATH) + ".tmp"
    with open(tmp, "w") as f:
        json.dump(ckpt, f)
    os.replace(tmp, CHECKPOINT_PATH)


# ---------------------------------------------------------------------------
# Leitura do CSV UTF-16 LE com BOM
# ---------------------------------------------------------------------------
def iter_rows(csv_path: Path):
    proc = subprocess.run(
        ["iconv", "-f", "UTF-16", "-t", "UTF-8", str(csv_path)],
        capture_output=True,
    )
    content = proc.stdout.decode("utf-8", errors="replace")
    reader = csv.DictReader(io.StringIO(content), delimiter="\t")
    for row in reader:
        yield row


# ---------------------------------------------------------------------------
# Main
# ---------------------------------------------------------------------------
def main():
    start_time = time.time()

    log.info("=" * 70)
    log.info("Iniciando importacao RD -> CRM Linkia (Borrello)")
    log.info(f"CSV: {CSV_PATH}")
    log.info(f"Inicio: {datetime.now(timezone.utc).isoformat()}")
    log.info("=" * 70)

    env = load_env(ENV_PATH)
    location_id = env["LINKIA_BORRELLO_LOCATION_ID"]
    token = env["LINKIA_BORRELLO_TOKEN"]
    log.info(f"Location ID: {location_id}")

    session = make_session(token)

    # Etapa 1: Garantir custom fields
    log.info("--- ETAPA 1: Custom Fields ---")
    field_id_map = ensure_custom_fields(session, location_id)
    log.info(f"Custom fields prontos: {len(field_id_map)}")

    # Carregar checkpoint
    checkpoint = load_checkpoint()
    already_done = set(checkpoint.keys())
    log.info(f"Checkpoint carregado: {len(already_done)} emails ja processados")

    # Contadores
    n_processed = 0
    n_created = 0
    n_updated = 0
    n_errors = 0
    n_skipped = 0
    tag_counts = defaultdict(int)

    erros_file = open(ERRORS_PATH, "a", encoding="utf-8")

    log.info("--- ETAPA 2: Importacao de contatos ---")

    try:
        for row in iter_rows(CSV_PATH):
            email = clean(row.get("Email", "")).lower()

            if not email:
                n_skipped += 1
                if n_skipped <= 10:
                    log.warning("Linha sem email ignorada")
                continue

            if email in already_done:
                n_skipped += 1
                continue

            n_processed += 1

            # Dedupe
            existing_id = find_contact_by_email(session, location_id, email)
            time.sleep(DELAY_BETWEEN_REQS)

            # Montar payload
            try:
                payload = build_contact_payload(row, field_id_map, location_id)
            except Exception as e:
                log.error(f"Erro payload {email}: {e}")
                erros_file.write(
                    json.dumps({"email": email, "error": str(e)}, ensure_ascii=False) + "\n"
                )
                n_errors += 1
                checkpoint[email] = "ERROR_PAYLOAD"
                save_checkpoint(checkpoint)
                time.sleep(DELAY_BETWEEN_REQS)
                continue

            # Contabiliza tags de estagio
            estagio = clean(
                row.get("Estágio no funil", "")
            ).lower()
            tag_e = TAG_ESTAGIO.get(estagio)
            if tag_e:
                tag_counts[tag_e] += 1
            tag_counts[TAG_GLOBAL] += 1

            # Upsert
            contact_id, action = upsert_contact(
                session, location_id, payload, existing_id
            )
            time.sleep(DELAY_BETWEEN_REQS)

            if "error" in action:
                n_errors += 1
                erros_file.write(
                    json.dumps({
                        "email": email,
                        "action": action,
                    }, ensure_ascii=False) + "\n"
                )
                checkpoint[email] = "ERROR"
                log.error(f"ERRO [{n_processed}] {email}: {action}")
            elif action == "created":
                n_created += 1
                checkpoint[email] = contact_id
            elif action == "updated":
                n_updated += 1
                checkpoint[email] = contact_id

            if n_processed % 100 == 0:
                save_checkpoint(checkpoint)

            if n_processed % LOG_INTERVAL == 0:
                elapsed = time.time() - start_time
                total_done = n_processed + len(already_done)
                total_estimate = 25426
                remaining = max(0, total_estimate - total_done)
                rate = n_processed / elapsed if elapsed > 0 else 1
                eta_min = (remaining / rate / 60) if rate > 0 else 0
                log.info(
                    f"processados={total_done}, criados={n_created}, "
                    f"updated={n_updated}, erros={n_errors}, "
                    f"ETA={eta_min:.0f}min"
                )

    finally:
        save_checkpoint(checkpoint)
        erros_file.close()

    # ---------------------------------------------------------------------------
    # Etapa 3: Validacao final
    # ---------------------------------------------------------------------------
    log.info("--- ETAPA 3: Validacao ---")
    log.info("Aguardando 30s para indexacao GHL...")
    time.sleep(30)

    total_ghl = None
    resp = api_request(
        session, "GET",
        f"{API_BASE}/contacts/",
        params={"locationId": location_id, "limit": 1},
    )
    if resp.ok:
        meta = resp.json().get("meta", {})
        total_ghl = meta.get("total")
        log.info(f"Total contatos na location (GHL): {total_ghl}")

    # 5 amostras aleatorias
    samples = []
    valid_ckpt = [
        (e, cid) for e, cid in checkpoint.items()
        if cid not in ("ERROR", "ERROR_PAYLOAD") and len(cid) > 10
    ]
    sample_pairs = random.sample(valid_ckpt, min(5, len(valid_ckpt)))

    for sample_email, contact_id in sample_pairs:
        resp = api_request(session, "GET", f"{API_BASE}/contacts/{contact_id}")
        if resp.ok:
            c = resp.json().get("contact", resp.json())
            samples.append({
                "email": c.get("email"),
                "firstName": c.get("firstName"),
                "lastName": c.get("lastName"),
                "tags": c.get("tags", [])[:10],
                "customFields_count": len(c.get("customFields", [])),
                "customFields_sample": {
                    cf.get("fieldKey", ""): cf.get("value", "")
                    for cf in (c.get("customFields", []) or [])[:3]
                },
            })
        time.sleep(DELAY_BETWEEN_REQS)

    # ---------------------------------------------------------------------------
    # Sumario
    # ---------------------------------------------------------------------------
    elapsed_total = time.time() - start_time
    total_done_final = n_processed + len(already_done)

    summary = {
        "timestamp": datetime.now(timezone.utc).isoformat(),
        "duracao_minutos": round(elapsed_total / 60, 1),
        "total_processados": total_done_final,
        "criados": n_created,
        "atualizados": n_updated,
        "erros": n_errors,
        "ignorados_sem_email": n_skipped,
        "total_ghl_location": total_ghl,
        "custom_fields_prontos": len(field_id_map),
        "tag_global": TAG_GLOBAL,
        "tags_estagio": TAG_ESTAGIO,
        "tag_distribuicao": dict(tag_counts),
        "amostras": samples,
    }

    with open(SUMMARY_PATH, "w", encoding="utf-8") as f:
        json.dump(summary, f, ensure_ascii=False, indent=2)

    log.info("=" * 70)
    log.info("IMPORTACAO CONCLUIDA")
    log.info(f"Processados:   {total_done_final}")
    log.info(f"Criados:       {n_created}")
    log.info(f"Atualizados:   {n_updated}")
    log.info(f"Erros:         {n_errors}")
    log.info(f"Sem email:     {n_skipped}")
    log.info(f"Total GHL:     {total_ghl}")
    log.info(f"Tempo total:   {elapsed_total/60:.1f} min")
    log.info(f"Resumo salvo:  {SUMMARY_PATH}")
    log.info("=" * 70)

    print("\n===RESULTADO_IMPORTACAO===")
    print(json.dumps(summary, ensure_ascii=False, indent=2))
    print("===FIM===\n")

    return summary


if __name__ == "__main__":
    main()
