#!/usr/bin/env python3
"""
GHL UTM Enricher - PX3 Lab
==========================

Micro-serviço webhook que enriquece contatos do Linkia PX3 com UTM baseado
no attributionSource NATIVO do GoHighLevel.

Fluxo:
  1. Recebe webhook (POST /webhook) do painel Linkia (evento Contact Create/Update).
  2. Extrai contactId do payload (aceita vários formatos).
  3. GET /contacts/{id} na GHL API -> pega attributionSource + customFields.
  4. Se contato JÁ tem custom field utm_source preenchido -> skip.
  5. Se attributionSource sinaliza origem paga/orgânica identificável ->
     mapeia pra utm_source/utm_medium/... e faz PUT /contacts/{id}.
  6. O workflow existente do Linkia ("Aplicar tag de origem") dispara
     naturalmente com o campo agora preenchido e aplica origem_meta/google/etc.

Dedupe:
  - SQLite local (dedupe.db) marca contact_id processado com TTL de 6h.
  - Evita reprocessar quando o Linkia dispara múltiplos eventos (Create + Update)
    do mesmo contato em curto período.

Endpoints:
  GET  /health              -> status + config resumida
  POST /webhook             -> recebe evento do Linkia
  POST /enrich/{contact_id} -> força enriquecimento manual de um contato
  GET  /enrich/{contact_id} -> mesmo mas útil pra teste em browser

Config: /opt/mia/config/ghl_utm_enricher_px3.env
Log:    /opt/mia/logs/ghl_utm_enricher_px3.log
DB:     ./dedupe.db

Handler modular: `enrich_contact()` é chamado tanto pelo webhook quanto pelo
backfill_enrichment.py. No futuro (Meta CAPI Purchase), plugar outro handler
em `run_handlers(contact, changes)`.
"""

from __future__ import annotations

import json
import logging
import os
import sqlite3
import sys
import threading
import time
from datetime import datetime, timezone
from logging.handlers import RotatingFileHandler
from pathlib import Path
from typing import Any

import requests
from flask import Flask, jsonify, request

BASE_DIR = Path(__file__).resolve().parent

# ---------------------------------------------------------------------------
# Logging
# ---------------------------------------------------------------------------
LOG_FILE = Path("/opt/mia/logs/ghl_utm_enricher_px3.log")
LOG_FILE.parent.mkdir(parents=True, exist_ok=True)

_fmt = logging.Formatter("%(asctime)s [%(levelname)s] %(name)s: %(message)s")
_fh = RotatingFileHandler(
    LOG_FILE, maxBytes=10 * 1024 * 1024, backupCount=5, 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("utm_enricher_px3")

# ---------------------------------------------------------------------------
# Env loading (dotenv com fallback manual)
# ---------------------------------------------------------------------------
ENV_FILE = "/opt/mia/config/ghl_utm_enricher_px3.env"
try:
    from dotenv import load_dotenv

    load_dotenv(ENV_FILE)
except ImportError:
    if os.path.exists(ENV_FILE):
        with open(ENV_FILE, encoding="utf-8") as f:
            for raw in f:
                line = raw.strip()
                if not line or line.startswith("#") or "=" not in line:
                    continue
                k, v = line.split("=", 1)
                os.environ.setdefault(k.strip(), v.strip())

GHL_TOKEN = os.environ.get("GHL_TOKEN", "").strip()
GHL_LOCATION_ID = os.environ.get("GHL_LOCATION_ID", "").strip()
GHL_API_BASE = os.environ.get("GHL_API_BASE", "https://services.leadconnectorhq.com").rstrip("/")
GHL_API_VERSION = os.environ.get("GHL_API_VERSION", "2021-07-28").strip()

UTM_FIELD_SOURCE = os.environ.get("UTM_FIELD_SOURCE", "").strip()
UTM_FIELD_MEDIUM = os.environ.get("UTM_FIELD_MEDIUM", "").strip()
UTM_FIELD_CAMPAIGN = os.environ.get("UTM_FIELD_CAMPAIGN", "").strip()
UTM_FIELD_CONTENT = os.environ.get("UTM_FIELD_CONTENT", "").strip()
UTM_FIELD_TERM = os.environ.get("UTM_FIELD_TERM", "").strip()

WEBHOOK_SECRET = os.environ.get("WEBHOOK_SECRET", "").strip()
DRY_RUN = os.environ.get("DRY_RUN", "0").strip() in ("1", "true", "yes")

if not GHL_TOKEN or not GHL_LOCATION_ID or not UTM_FIELD_SOURCE:
    log.critical(
        "Config incompleta: GHL_TOKEN=%s LOC=%s UTM_FIELD_SOURCE=%s",
        bool(GHL_TOKEN), bool(GHL_LOCATION_ID), bool(UTM_FIELD_SOURCE),
    )
    sys.exit(1)

# Mapa dos IDs de custom field -> "chave lógica"
CF_MAP = {
    "utm_source": UTM_FIELD_SOURCE,
    "utm_medium": UTM_FIELD_MEDIUM,
    "utm_campaign": UTM_FIELD_CAMPAIGN,
    "utm_content": UTM_FIELD_CONTENT,
    "utm_term": UTM_FIELD_TERM,
}
# só campos com ID configurado
CF_MAP = {k: v for k, v in CF_MAP.items() if v}


# ---------------------------------------------------------------------------
# Dedupe SQLite
# ---------------------------------------------------------------------------
DB_PATH = BASE_DIR / "dedupe.db"
DEDUPE_TTL_SECONDS = 6 * 60 * 60  # 6h — evita reprocessar em Create + Update
CLEANUP_TTL_SECONDS = 30 * 24 * 60 * 60  # retém 30 dias no DB

_db_lock = threading.Lock()


def _db():
    conn = sqlite3.connect(str(DB_PATH), timeout=5.0)
    conn.execute("PRAGMA journal_mode=WAL")
    return conn


def _init_db() -> None:
    with _db_lock, _db() as conn:
        conn.execute(
            """
            CREATE TABLE IF NOT EXISTS processed (
                contact_id TEXT PRIMARY KEY,
                ts         INTEGER NOT NULL,
                result     TEXT
            )
            """
        )
        conn.execute("CREATE INDEX IF NOT EXISTS idx_proc_ts ON processed(ts)")
        cutoff = int(time.time()) - CLEANUP_TTL_SECONDS
        n = conn.execute("DELETE FROM processed WHERE ts < ?", (cutoff,)).rowcount
        conn.commit()
        if n:
            log.info("dedupe cleanup: %d registros expirados removidos", n)


def _is_recently_processed(contact_id: str) -> bool:
    if not contact_id:
        return False
    cutoff = int(time.time()) - DEDUPE_TTL_SECONDS
    with _db_lock, _db() as conn:
        row = conn.execute(
            "SELECT ts FROM processed WHERE contact_id = ? AND ts >= ?",
            (contact_id, cutoff),
        ).fetchone()
    return row is not None


def _mark_processed(contact_id: str, result: str) -> None:
    if not contact_id:
        return
    now = int(time.time())
    with _db_lock, _db() as conn:
        conn.execute(
            "INSERT OR REPLACE INTO processed (contact_id, ts, result) VALUES (?, ?, ?)",
            (contact_id, now, result),
        )
        conn.commit()


# ---------------------------------------------------------------------------
# GHL API helpers
# ---------------------------------------------------------------------------
def _ghl_headers() -> dict[str, str]:
    return {
        "Authorization": f"Bearer {GHL_TOKEN}",
        "Version": GHL_API_VERSION,
        "Accept": "application/json",
        "Content-Type": "application/json",
    }


def fetch_contact(contact_id: str) -> dict | None:
    try:
        r = requests.get(
            f"{GHL_API_BASE}/contacts/{contact_id}",
            headers=_ghl_headers(),
            timeout=8,
        )
        if r.status_code == 404:
            log.warning("contato %s nao encontrado (404)", contact_id)
            return None
        if r.status_code >= 300:
            log.warning(
                "GET /contacts/%s falhou: %s %s",
                contact_id, r.status_code, r.text[:200],
            )
            return None
        return (r.json() or {}).get("contact") or {}
    except requests.RequestException as e:
        log.error("GET /contacts/%s excecao: %s", contact_id, e)
        return None


def update_contact_custom_fields(contact_id: str, fields: list[dict]) -> tuple[bool, str]:
    """PUT /contacts/{id} com customFields. Retorna (ok, resp_body_or_err)."""
    if DRY_RUN:
        log.info("DRY_RUN: PUT /contacts/%s custom_fields=%s", contact_id, fields)
        return True, "dry_run"
    try:
        r = requests.put(
            f"{GHL_API_BASE}/contacts/{contact_id}",
            headers=_ghl_headers(),
            json={"customFields": fields},
            timeout=8,
        )
        ok = 200 <= r.status_code < 300
        return ok, f"{r.status_code} {r.text[:300]}"
    except requests.RequestException as e:
        return False, f"REQ_ERROR: {e}"


# ---------------------------------------------------------------------------
# Attribution -> UTM mapping (o coração do serviço)
# ---------------------------------------------------------------------------
def _get_cf_value(contact: dict, field_id: str) -> str | None:
    for cf in contact.get("customFields", []) or []:
        if cf.get("id") == field_id:
            v = cf.get("value")
            if v is None:
                return None
            s = str(v).strip()
            return s or None
    return None


def _map_attribution_to_utm(att: dict) -> dict[str, str]:
    """
    Regras de mapeamento (atribuição nativa GHL -> UTMs lógicos).

    Prioridade:
    1. Se attribution já traz utmSource, usa direto (com utmMedium/Campaign/etc.).
    2. sessionSource == "Paid Social" -> utm_source=meta, utm_medium=whatsapp_ads
       (se medium=whatsapp) ou paid_social; extra: adName->utm_campaign, adId->utm_content.
    3. sessionSource == "Social media" / "Social Media" (orgânico) -> utm_source=meta,
       utm_medium=whatsapp_organic (se medium=whatsapp) senão social.
    4. sessionSource contém "Paid Search" / "cpc" -> utm_source=google, utm_medium=cpc.
    5. sessionSource == "Organic Search" -> utm_source=google, utm_medium=organic.
    6. sessionSource == "Referral" -> utm_source=referral, utm_medium=referral.
    7. sessionSource == "Email" -> utm_source=email.
    8. Direct / CRM Workflows / CRM UI / Manual / vazio -> {} (deixa workflow
       classificar como origem_desconhecida).
    """
    if not att or not isinstance(att, dict):
        return {}

    out: dict[str, str] = {}
    utm_src = (att.get("utmSource") or "").strip()
    utm_med = (att.get("utmMedium") or "").strip()
    utm_camp = (att.get("utmCampaign") or "").strip()
    utm_cont = (att.get("utmContent") or "").strip()
    utm_term = (att.get("utmTerm") or "").strip()
    sess_src = (att.get("sessionSource") or "").strip()
    medium = (att.get("medium") or "").strip().lower()
    ad_name = (att.get("adName") or "").strip()
    ad_id = (att.get("adId") or "").strip()

    # Caso 1: attribution já tem UTMs -> usa direto
    if utm_src:
        out["utm_source"] = utm_src.lower()
        if utm_med:
            out["utm_medium"] = utm_med.lower()
        if utm_camp:
            out["utm_campaign"] = utm_camp
        if utm_cont:
            out["utm_content"] = utm_cont
        if utm_term:
            out["utm_term"] = utm_term
        return out

    # Caso 2+: infere pela sessionSource (case-insensitive)
    sess_lower = sess_src.lower()

    if "paid social" in sess_lower or "paid_social" in sess_lower:
        out["utm_source"] = "meta"
        out["utm_medium"] = "whatsapp_ads" if "whats" in medium else "paid_social"
        # Bonus: aproveita adName/adId nativos se existirem
        if ad_name:
            out["utm_campaign"] = ad_name[:200]
        if ad_id:
            out["utm_content"] = ad_id
    elif "social media" in sess_lower or "organic social" in sess_lower:
        out["utm_source"] = "meta"
        out["utm_medium"] = "whatsapp_organic" if "whats" in medium else "social"
    elif "paid search" in sess_lower or "cpc" in sess_lower:
        out["utm_source"] = "google"
        out["utm_medium"] = "cpc"
    elif "organic search" in sess_lower:
        out["utm_source"] = "google"
        out["utm_medium"] = "organic"
    elif "referral" in sess_lower:
        out["utm_source"] = "referral"
        out["utm_medium"] = "referral"
    elif "email" in sess_lower:
        out["utm_source"] = "email"
        out["utm_medium"] = "email"
    # Direct / CRM Workflows / CRM UI / Manual / vazio -> {} (workflow classifica como desconhecida)

    return out


def _build_custom_fields_payload(utm: dict[str, str]) -> list[dict]:
    """Converte {utm_source: 'meta', ...} em [{id: <cfid>, value: 'meta'}, ...]"""
    out = []
    for key, val in utm.items():
        cf_id = CF_MAP.get(key)
        if not cf_id or not val:
            continue
        out.append({"id": cf_id, "value": val})
    return out


# ---------------------------------------------------------------------------
# Core: enrich_contact (chamado pelo webhook E pelo backfill)
# ---------------------------------------------------------------------------
def enrich_contact(contact_id: str, force: bool = False) -> dict:
    """
    Aplica enriquecimento UTM a um contato.

    Args:
        contact_id: ID GHL
        force: se True, ignora dedupe (mas ainda respeita 'já tem utm_source').

    Retorna dict com resultado (ok/skip/erro) — sempre serializável.
    """
    if not contact_id:
        return {"ok": False, "reason": "no_contact_id"}

    if not force and _is_recently_processed(contact_id):
        return {"ok": True, "skip": "dedupe", "contact_id": contact_id}

    contact = fetch_contact(contact_id)
    if not contact:
        return {"ok": False, "reason": "contact_not_found", "contact_id": contact_id}

    # Guard: não sobrescreve utm_source existente
    current_utm_source = _get_cf_value(contact, UTM_FIELD_SOURCE)
    if current_utm_source:
        _mark_processed(contact_id, "skip_has_utm_source")
        return {
            "ok": True, "skip": "already_has_utm_source",
            "contact_id": contact_id, "utm_source": current_utm_source,
        }

    att = contact.get("attributionSource") or {}
    last_att = contact.get("lastAttributionSource") or {}

    # Tenta attributionSource primeiro; se vazio, tenta lastAttributionSource
    utm = _map_attribution_to_utm(att)
    used = "attributionSource"
    if not utm and last_att:
        utm = _map_attribution_to_utm(last_att)
        used = "lastAttributionSource"

    if not utm:
        _mark_processed(contact_id, "skip_no_mapping")
        return {
            "ok": True, "skip": "no_attribution_mapping",
            "contact_id": contact_id,
            "attribution": att, "last_attribution": last_att,
        }

    cf_payload = _build_custom_fields_payload(utm)
    if not cf_payload:
        _mark_processed(contact_id, "skip_no_payload")
        return {
            "ok": True, "skip": "empty_payload",
            "contact_id": contact_id, "utm": utm,
        }

    ok, resp = update_contact_custom_fields(contact_id, cf_payload)
    result = "written" if ok else "write_error"
    _mark_processed(contact_id, result)

    return {
        "ok": ok,
        "contact_id": contact_id,
        "name": contact.get("firstName"),
        "used_attribution": used,
        "utm_written": utm,
        "cf_payload": cf_payload,
        "resp": resp if not ok else "ok",
        "dry_run": DRY_RUN,
    }


# ---------------------------------------------------------------------------
# Modular handler chain (futuro: Meta CAPI Purchase, etc.)
# ---------------------------------------------------------------------------
def run_handlers(contact_id: str, event_type: str) -> list[dict]:
    """
    Executa todos os handlers pra este contact_id/evento.
    Por ora só o enricher; no futuro plug outros (ex: Meta CAPI).
    """
    results = []
    # 1) UTM enricher
    results.append({"handler": "utm_enricher", **enrich_contact(contact_id)})
    return results


# ---------------------------------------------------------------------------
# Webhook payload parser
# ---------------------------------------------------------------------------
def _extract_contact_id(payload: dict) -> tuple[str, str]:
    """
    Retorna (contact_id, event_type). O Linkia/GHL manda formatos variados:
    - {"type": "ContactCreate", "contact_id": "..."}
    - {"type": "ContactCreate", "id": "..."}
    - {"contactId": "...", "type": "..."}
    - {"contact": {"id": "..."}, "type": "..."}
    """
    if not isinstance(payload, dict):
        return "", ""
    et = str(payload.get("type") or payload.get("event") or "").strip()
    cid = ""
    for key in ("contact_id", "contactId", "id"):
        v = payload.get(key)
        if v and isinstance(v, str) and len(v) >= 10:
            cid = v.strip()
            break
    if not cid:
        c = payload.get("contact")
        if isinstance(c, dict):
            cid = str(c.get("id") or "").strip()
    return cid, et


# ---------------------------------------------------------------------------
# Flask app
# ---------------------------------------------------------------------------
app = Flask(__name__)


@app.get("/health")
def health():
    return jsonify(
        {
            "ok": True,
            "service": "ghl_utm_enricher_px3",
            "location": GHL_LOCATION_ID,
            "cf_fields": list(CF_MAP.keys()),
            "dry_run": DRY_RUN,
            "webhook_secret_required": bool(WEBHOOK_SECRET),
        }
    )


@app.post("/webhook")
def webhook():
    # Auth opcional
    if WEBHOOK_SECRET:
        got = request.headers.get("X-Webhook-Secret", "").strip()
        if got != WEBHOOK_SECRET:
            log.warning("webhook rejeitado: secret invalido (ip=%s)", request.remote_addr)
            return jsonify({"ok": False, "error": "unauthorized"}), 401

    payload = request.get_json(silent=True) or {}
    contact_id, event_type = _extract_contact_id(payload)

    log.info(
        "webhook recebido: event=%s contact_id=%s ip=%s",
        event_type or "?",
        contact_id or "(vazio)",
        request.headers.get("X-Forwarded-For", request.remote_addr or "?"),
    )

    if not contact_id:
        log.warning("webhook sem contact_id parseável: %s", json.dumps(payload)[:400])
        return jsonify({"ok": False, "error": "no_contact_id", "payload_keys": list(payload.keys())}), 400

    results = run_handlers(contact_id, event_type)
    log.info(
        "webhook processado: contact=%s results=%s",
        contact_id, json.dumps(results, ensure_ascii=False)[:500],
    )
    return jsonify({"ok": True, "contact_id": contact_id, "results": results})


@app.route("/enrich/<contact_id>", methods=["GET", "POST"])
def enrich_endpoint(contact_id: str):
    """Manual: força enriquecimento de um contato (ignora dedupe)."""
    force = request.args.get("force", "1") in ("1", "true", "yes")
    result = enrich_contact(contact_id, force=force)
    log.info("manual enrich contact=%s result=%s", contact_id, json.dumps(result, ensure_ascii=False)[:400])
    return jsonify(result)


# ---------------------------------------------------------------------------
# Boot
# ---------------------------------------------------------------------------
_init_db()

if __name__ == "__main__":
    port = int(os.environ.get("PORT", "8919"))
    host = os.environ.get("HOST", "0.0.0.0")
    log.info(
        "ghl_utm_enricher_px3 subindo: host=%s port=%s dry_run=%s location=%s",
        host, port, DRY_RUN, GHL_LOCATION_ID,
    )
    app.run(host=host, port=port, debug=False, threaded=True)
