#!/usr/bin/env python3
"""
Segundo passe da importacao RD -> Linkia (Borrello).

Objetivo: pegar os leads que falharam no 1o passe por 'duplicated contact / phone'
e fazer UPDATE no contato existente (dono do phone), MERGE-ando tags e custom fields.

Regras:
  - Nao sobrescreve firstName/lastName/email/phone do contato existente.
  - Tags: uniao (dedup) das tags atuais + as tags que o lead-que-falhou traria.
  - Custom Fields: se o campo existente ja tem valor, MANTEM (loga conflito);
    se esta ausente, adiciona o valor do RD.
  - DRY RUN: processa os 3 primeiros contactIds sem PUT antes de rodar tudo.
  - Rate: 1 req/s (get + put por lead => 2 reqs, ~0.5s/req).
  - Retry backoff 429/timeout.

Uso:
  python3 segundo_passe_dup.py --dry-run   # so log, sem PUT
  python3 segundo_passe_dup.py             # roda de verdade
"""

import argparse
import csv
import io
import json
import logging
import os
import re
import subprocess
import sys
import time
from datetime import datetime, timezone
from pathlib import Path

import requests

# Reutiliza helpers do script original
sys.path.insert(0, str(Path(__file__).parent))
from importar_rd_linkia import (  # type: ignore
    load_env, make_session, api_request,
    build_contact_payload, build_tags, build_custom_values,
    clean, ENV_PATH, API_BASE,
)

BASE_DIR = Path("/opt/mia/workspace/clientes/borrello/importacao_rd_20260828")
CSV_PATH = BASE_DIR / "rd_export.csv"
ERRORS_JSONL = BASE_DIR / "erros.jsonl"
LOG_PATH = BASE_DIR / "segundo_passe.log"
SUMMARY_PATH = BASE_DIR / "segundo_passe_resumo.json"

DELAY_BETWEEN_REQS = 0.5   # 2 reqs por lead => ~1s efetivo por lead
MAX_BACKOFF = 60

# Reaproveita a spec de custom fields do script original para o mapa
from importar_rd_linkia import CUSTOM_FIELDS_SPEC  # type: ignore


# ---------------------------------------------------------------------------
# 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__)


# ---------------------------------------------------------------------------
# Parse dos erros: pega email + contactId + matchingField
# ---------------------------------------------------------------------------
RE_CONTACT_ID = re.compile(r'"contactId":"([^"]+)"')
RE_MATCHING = re.compile(r'"matchingField":"([^"]+)"')


def parse_errors_jsonl(path: Path):
    """
    Retorna lista de dicts {email, contactId, matchingField}
    filtrando SO os erros de duplicated contact.
    """
    out = []
    with open(path, encoding="utf-8") as f:
        for line in f:
            line = line.strip()
            if not line:
                continue
            try:
                d = json.loads(line)
            except Exception:
                continue
            action = d.get("action", "")
            if "This location does not allow duplicated" not in action:
                continue
            m = RE_CONTACT_ID.search(action)
            if not m:
                continue
            m2 = RE_MATCHING.search(action)
            out.append({
                "email": d.get("email", "").lower(),
                "contactId": m.group(1),
                "matchingField": m2.group(1) if m2 else "unknown",
            })
    return out


# ---------------------------------------------------------------------------
# CSV -> dict email->row
# ---------------------------------------------------------------------------
def load_csv_by_email(csv_path: Path) -> dict:
    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")
    idx = {}
    for row in reader:
        email = clean(row.get("Email", "")).lower()
        if email:
            idx[email] = row
    return idx


# ---------------------------------------------------------------------------
# Custom fields map: pega da location (nao recria)
# ---------------------------------------------------------------------------
def load_field_id_map(session, location_id: str) -> dict:
    """
    Retorna {rd_key -> field_id} casando pelo LABEL (name) da spec original.
    O GHL renomeia fieldKeys automaticamente (ex: rd_estagio_funil -> estagio_no_funil_rd),
    entao o label eh a chave estavel.
    """
    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]}"
        )
    data = resp.json()
    fields_list = data if isinstance(data, list) else data.get("customFields", [])
    by_name = {f.get("name", "").strip(): f.get("id") for f in fields_list}

    rd_map = {}
    for (rd_key, _dt, label, _opts) in CUSTOM_FIELDS_SPEC:
        fid = by_name.get(label.strip())
        if fid:
            rd_map[rd_key] = fid
        else:
            log.warning(f"Custom field nao encontrado por label: '{label}' (rd_key={rd_key})")
    return rd_map


# ---------------------------------------------------------------------------
# Merge helpers
# ---------------------------------------------------------------------------
def merge_tags(current: list, new: list) -> tuple:
    """Retorna (union_list, added_count). Preserva ordem, case-insensitive dedupe."""
    seen = set()
    out = []
    for t in current or []:
        k = t.lower()
        if k not in seen:
            seen.add(k)
            out.append(t)
    added = 0
    for t in new or []:
        k = t.lower()
        if k not in seen:
            seen.add(k)
            out.append(t)
            added += 1
    return out, added


def merge_custom_fields(current: list, new: list) -> tuple:
    """
    current: [{id, value}, ...] do GET
    new:     [{id, value}, ...] montado do CSV
    Retorna (union_list, added_count, conflict_count, conflicts_detail)
    Se o id ja existe no current com valor nao-vazio -> mantem current (conflito se valor diferente).
    Se o id nao existe -> adiciona new.
    """
    by_id = {}
    for cf in current or []:
        cid = cf.get("id")
        if cid:
            by_id[cid] = cf

    added = 0
    conflict = 0
    conflicts_detail = []

    for cf in new or []:
        cid = cf.get("id")
        nval = cf.get("value")
        if not cid:
            continue
        cur = by_id.get(cid)
        if cur is None:
            by_id[cid] = {"id": cid, "value": nval}
            added += 1
        else:
            cur_val = cur.get("value")
            # se atual ta vazio, preenche
            if cur_val in (None, "", []):
                by_id[cid] = {"id": cid, "value": nval}
                added += 1
            elif str(cur_val).strip() != str(nval).strip():
                conflict += 1
                if len(conflicts_detail) < 20:
                    conflicts_detail.append({
                        "id": cid,
                        "current": str(cur_val)[:80],
                        "new": str(nval)[:80],
                    })
            # se igual, nao faz nada
    return list(by_id.values()), added, conflict, conflicts_detail


# ---------------------------------------------------------------------------
# GET contato existente
# ---------------------------------------------------------------------------
def get_contact(session, contact_id: str):
    resp = api_request(session, "GET", f"{API_BASE}/contacts/{contact_id}")
    if resp.status_code == 404:
        return None, "not_found"
    if not resp.ok:
        return None, f"error:{resp.status_code}:{resp.text[:200]}"
    d = resp.json()
    return d.get("contact", d), "ok"


# ---------------------------------------------------------------------------
# PUT contato (update)
# ---------------------------------------------------------------------------
# GHL v2 PUT /contacts/{id} NAO aceita locationId no body.
UPDATE_ALLOWED_KEYS = {"tags", "customFields"}


def update_contact(session, contact_id: str, payload: dict):
    body = {k: v for k, v in payload.items() if k in UPDATE_ALLOWED_KEYS}
    resp = api_request(session, "PUT",
                       f"{API_BASE}/contacts/{contact_id}",
                       json=body)
    if resp.ok:
        return True, "ok"
    return False, f"error:{resp.status_code}:{resp.text[:300]}"


# ---------------------------------------------------------------------------
# Processa 1 lead
# ---------------------------------------------------------------------------
def process_one(session, location_id, err, csv_idx, field_id_map, dry_run):
    """
    err: {email, contactId, matchingField}
    Retorna dict com status pra logar.
    """
    email = err["email"]
    contact_id = err["contactId"]

    row = csv_idx.get(email)
    if not row:
        return {"status": "no_csv_row", "email": email, "contactId": contact_id}

    # Payload como o 1o passe teria montado
    payload = build_contact_payload(row, field_id_map, location_id)
    new_tags = payload.get("tags", [])
    new_cfs = payload.get("customFields", [])

    # GET contato existente
    contact, status = get_contact(session, contact_id)
    time.sleep(DELAY_BETWEEN_REQS)
    if status == "not_found":
        return {"status": "contact_404", "email": email, "contactId": contact_id}
    if status != "ok":
        return {"status": "get_failed", "email": email, "contactId": contact_id,
                "reason": status}

    cur_tags = contact.get("tags", []) or []
    cur_cfs = contact.get("customFields", []) or []

    merged_tags, tags_added = merge_tags(cur_tags, new_tags)
    merged_cfs, cfs_added, cfs_conflict, cfs_conflicts_detail = merge_custom_fields(
        cur_cfs, new_cfs
    )

    # Se nada a adicionar, skip
    if tags_added == 0 and cfs_added == 0:
        return {
            "status": "no_changes",
            "email": email,
            "contactId": contact_id,
            "existing_email": contact.get("email"),
            "cfs_conflict": cfs_conflict,
        }

    update_payload = {
        "tags": merged_tags,
        "customFields": merged_cfs,
    }

    result = {
        "status": "would_update" if dry_run else "pending",
        "email": email,
        "contactId": contact_id,
        "existing_email": contact.get("email"),
        "existing_phone": contact.get("phone"),
        "tags_added": tags_added,
        "cfs_added": cfs_added,
        "cfs_conflict": cfs_conflict,
        "cfs_conflicts_sample": cfs_conflicts_detail[:5],
        "new_tags_added": [t for t in new_tags if t.lower() not in {x.lower() for x in cur_tags}][:10],
    }

    if dry_run:
        return result

    ok, msg = update_contact(session, contact_id, update_payload)
    time.sleep(DELAY_BETWEEN_REQS)
    if ok:
        result["status"] = "updated"
    else:
        result["status"] = "put_failed"
        result["reason"] = msg
    return result


# ---------------------------------------------------------------------------
# Main
# ---------------------------------------------------------------------------
def main():
    parser = argparse.ArgumentParser()
    parser.add_argument("--dry-run", action="store_true",
                        help="Nao faz PUT, so loga o que faria (primeiros 3)")
    parser.add_argument("--limit", type=int, default=0,
                        help="Limite de leads a processar (0 = todos)")
    args = parser.parse_args()

    start = time.time()
    log.info("=" * 70)
    log.info("SEGUNDO PASSE - Duplicated contacts (MERGE tags + customFields)")
    log.info(f"DRY_RUN={args.dry_run}  LIMIT={args.limit}")
    log.info("=" * 70)

    # 1. Env + session
    env = load_env(ENV_PATH)
    location_id = env["LINKIA_BORRELLO_LOCATION_ID"]
    token = env["LINKIA_BORRELLO_TOKEN"]
    session = make_session(token)
    log.info(f"Location: {location_id}")

    # 2. Custom fields map
    log.info("Carregando custom fields da location...")
    field_id_map = load_field_id_map(session, location_id)
    log.info(f"Custom fields RD encontrados: {len(field_id_map)}")

    # 3. Parse erros
    log.info(f"Lendo erros de {ERRORS_JSONL}...")
    errors = parse_errors_jsonl(ERRORS_JSONL)
    log.info(f"Erros de duplicated contact com contactId: {len(errors)}")

    # 4. Carrega CSV em memoria
    log.info(f"Carregando CSV {CSV_PATH} em memoria...")
    csv_idx = load_csv_by_email(CSV_PATH)
    log.info(f"Emails indexados do CSV: {len(csv_idx)}")

    # 5. Verifica cobertura
    faltando = [e["email"] for e in errors if e["email"] not in csv_idx]
    log.info(f"Erros sem match no CSV: {len(faltando)}")
    if faltando[:5]:
        log.info(f"  Exemplos: {faltando[:5]}")

    # 6. DRY RUN dos 3 primeiros SEMPRE, mesmo se nao for --dry-run
    #    (regra da tarefa: confirmar payload antes de PUT em prod)
    log.info("--- DRY RUN dos 3 primeiros para inspecao ---")
    for i, err in enumerate(errors[:3]):
        r = process_one(session, location_id, err, csv_idx, field_id_map, dry_run=True)
        log.info(f"  [{i}] {json.dumps(r, ensure_ascii=False)}")

    if args.dry_run:
        log.info("DRY RUN completo. Nao vai fazer PUT. Sai aqui.")
        return

    # 7. Roda de verdade
    log.info("--- EXECUCAO REAL ---")
    results = []
    counts = {
        "updated": 0,
        "no_changes": 0,
        "no_csv_row": 0,
        "contact_404": 0,
        "get_failed": 0,
        "put_failed": 0,
        "would_update": 0,
        "pending": 0,
    }
    tags_added_total = 0
    cfs_added_total = 0
    cfs_conflict_total = 0

    total = len(errors) if not args.limit else min(args.limit, len(errors))
    log.info(f"Total a processar: {total}")

    for i, err in enumerate(errors[:total], 1):
        try:
            r = process_one(session, location_id, err, csv_idx, field_id_map,
                            dry_run=False)
        except Exception as e:
            r = {"status": "exception", "email": err["email"],
                 "contactId": err["contactId"], "reason": str(e)}
            log.exception(f"exception no lead {err['email']}")

        results.append(r)
        counts[r["status"]] = counts.get(r["status"], 0) + 1
        tags_added_total += r.get("tags_added", 0) or 0
        cfs_added_total += r.get("cfs_added", 0) or 0
        cfs_conflict_total += r.get("cfs_conflict", 0) or 0

        if i % 25 == 0 or i == total:
            elapsed = time.time() - start
            rate = i / elapsed if elapsed else 0
            eta_min = (total - i) / rate / 60 if rate else 0
            log.info(
                f"progresso {i}/{total}  updated={counts['updated']}  "
                f"no_changes={counts['no_changes']}  put_failed={counts['put_failed']}  "
                f"404={counts['contact_404']}  ETA={eta_min:.1f}min"
            )

    # 8. Resumo
    duracao = time.time() - start
    put_failed_reasons = {}
    for r in results:
        if r["status"] == "put_failed":
            reason = r.get("reason", "")[:80]
            put_failed_reasons[reason] = put_failed_reasons.get(reason, 0) + 1

    summary = {
        "timestamp": datetime.now(timezone.utc).isoformat(),
        "duracao_minutos": round(duracao / 60, 2),
        "total_erros_processados": total,
        "updates_ok": counts["updated"],
        "updates_falhou": counts["put_failed"],
        "updates_falhou_motivos": put_failed_reasons,
        "sem_alteracao": counts["no_changes"],
        "sem_linha_no_csv": counts["no_csv_row"],
        "contato_deletado_404": counts["contact_404"],
        "get_falhou": counts["get_failed"],
        "tags_novas_adicionadas_total": tags_added_total,
        "customfields_novos_adicionados_total": cfs_added_total,
        "customfields_conflito_total": cfs_conflict_total,
    }

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

    # detalhe completo pra auditoria
    detail_path = BASE_DIR / "segundo_passe_detalhes.jsonl"
    with open(detail_path, "w", encoding="utf-8") as f:
        for r in results:
            f.write(json.dumps(r, ensure_ascii=False) + "\n")

    log.info("=" * 70)
    log.info("SEGUNDO PASSE CONCLUIDO")
    log.info(f"Total:            {total}")
    log.info(f"Updates OK:       {counts['updated']}")
    log.info(f"Sem alteracao:    {counts['no_changes']}")
    log.info(f"PUT falhou:       {counts['put_failed']}")
    log.info(f"Contato 404:      {counts['contact_404']}")
    log.info(f"Sem linha CSV:    {counts['no_csv_row']}")
    log.info(f"Tags adicionadas: {tags_added_total}")
    log.info(f"CFs adicionados:  {cfs_added_total}")
    log.info(f"CFs conflito:     {cfs_conflict_total}")
    log.info(f"Duracao:          {duracao/60:.2f} min")
    log.info(f"Resumo:           {SUMMARY_PATH}")
    log.info(f"Detalhes:         {detail_path}")
    log.info("=" * 70)


if __name__ == "__main__":
    main()
