#!/usr/bin/env python3
"""
Atualiza o campo 'Status Tarefa' nas oportunidades do pipeline PX3.
🟢 = tem tarefa futura (em dia)
🟡 = sem tarefa cadastrada
🔴 = tem tarefa vencida (atrasada)
Roda via cron a cada 30 minutos (com flock).

Melhorias 2026-09-15:
- Credenciais em /opt/mia/config/farois_px3.env (chmod 600)
- Distingue "sem tarefa" de "falha de leitura" (não pinta 🟡 falso)
- Cache SQLite local para pular PUTs redundantes
- requests.Session com Retry/backoff (429/5xx + Retry-After)
- Log rotacionado (10MB x 3 backups) + duração da execução
"""

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

import requests
from requests.adapters import HTTPAdapter
from urllib3.util.retry import Retry

# ---------------------------------------------------------------------------
# Configuração
# ---------------------------------------------------------------------------

ENV_FILE = "/opt/mia/config/farois_px3.env"
DB_PATH = "/opt/mia/data/farois_px3.db"
LOG_PATH = "/opt/mia/logs/farois_pipeline.log"

GHL_BASE = "https://services.leadconnectorhq.com"
CAMPO_STATUS_ID = "cvrJJ7rlAuIGgCL42NEH"  # campo customizado criado 2026-06-23


def _load_env(path: str) -> None:
    """Carrega KEY=VALUE do arquivo pra os.environ (dotenv-like, sem dependência)."""
    if not os.path.isfile(path):
        return
    with open(path, "r", encoding="utf-8") as f:
        for raw in f:
            line = raw.strip()
            if not line or line.startswith("#") or "=" not in line:
                continue
            key, _, val = line.partition("=")
            key = key.strip()
            val = val.strip().strip('"').strip("'")
            os.environ.setdefault(key, val)


_load_env(ENV_FILE)

GHL_TOKEN = os.getenv("GHL_TOKEN", "").strip()
GHL_LOCATION = os.getenv("GHL_LOCATION", "").strip()

if not GHL_TOKEN or not GHL_LOCATION:
    sys.stderr.write(
        f"[FATAL] GHL_TOKEN/GHL_LOCATION ausentes. Verifique {ENV_FILE}\n"
    )
    sys.exit(2)

HEADERS = {
    "Authorization": f"Bearer {GHL_TOKEN}",
    "Version": "2021-07-28",
    "Content-Type": "application/json",
}

# ---------------------------------------------------------------------------
# Logger com rotação
# ---------------------------------------------------------------------------

Path(LOG_PATH).parent.mkdir(parents=True, exist_ok=True)

logger = logging.getLogger("farois_px3")
logger.setLevel(logging.INFO)
logger.propagate = False
if not logger.handlers:
    handler = RotatingFileHandler(
        LOG_PATH, maxBytes=10 * 1024 * 1024, backupCount=3, encoding="utf-8"
    )
    handler.setFormatter(
        logging.Formatter("%(asctime)s %(levelname)s %(message)s", "%Y-%m-%d %H:%M:%S")
    )
    logger.addHandler(handler)
    # Também espelha no stdout (o cron faz append; útil pra debug manual).
    stream = logging.StreamHandler(sys.stdout)
    stream.setFormatter(
        logging.Formatter("%(asctime)s %(levelname)s %(message)s", "%Y-%m-%d %H:%M:%S")
    )
    logger.addHandler(stream)


def log(msg: str, level: str = "INFO") -> None:
    getattr(logger, level.lower(), logger.info)(msg)


# ---------------------------------------------------------------------------
# HTTP session com retry / backoff
# ---------------------------------------------------------------------------

def _build_session() -> requests.Session:
    s = requests.Session()
    retry = Retry(
        total=3,
        connect=3,
        read=3,
        status=3,
        backoff_factor=1.0,  # 1s, 2s, 4s
        status_forcelist=(429, 500, 502, 503, 504),
        allowed_methods=frozenset(["GET", "PUT", "POST"]),
        respect_retry_after_header=True,
        raise_on_status=False,
    )
    adapter = HTTPAdapter(max_retries=retry, pool_connections=32, pool_maxsize=32)
    s.mount("https://", adapter)
    s.mount("http://", adapter)
    s.headers.update(HEADERS)
    return s


SESSION = _build_session()


# ---------------------------------------------------------------------------
# Cache SQLite
# ---------------------------------------------------------------------------

def _cache_conn() -> sqlite3.Connection:
    Path(DB_PATH).parent.mkdir(parents=True, exist_ok=True)
    conn = sqlite3.connect(DB_PATH, timeout=10)
    conn.execute(
        """
        CREATE TABLE IF NOT EXISTS farois_cache (
            opp_id TEXT PRIMARY KEY,
            farol TEXT NOT NULL,
            updated_at TEXT NOT NULL
        )
        """
    )
    conn.execute(
        "CREATE INDEX IF NOT EXISTS idx_farois_cache_updated ON farois_cache (updated_at)"
    )
    conn.commit()
    return conn


def cache_get(conn: sqlite3.Connection, opp_id: str) -> str | None:
    row = conn.execute("SELECT farol FROM farois_cache WHERE opp_id = ?", (opp_id,)).fetchone()
    return row[0] if row else None


def cache_set(conn: sqlite3.Connection, opp_id: str, farol: str) -> None:
    conn.execute(
        """
        INSERT INTO farois_cache (opp_id, farol, updated_at)
        VALUES (?, ?, ?)
        ON CONFLICT(opp_id) DO UPDATE SET farol=excluded.farol, updated_at=excluded.updated_at
        """,
        (opp_id, farol, datetime.now(timezone.utc).isoformat()),
    )


# ---------------------------------------------------------------------------
# Chamadas GHL
# ---------------------------------------------------------------------------

def buscar_oportunidades() -> list[dict]:
    """Busca todas as oportunidades abertas da location (paginação cursor-based)."""
    opps: list[dict] = []
    params = {"location_id": GHL_LOCATION, "status": "open", "limit": 100}
    while True:
        try:
            r = SESSION.get(f"{GHL_BASE}/opportunities/search", params=params, timeout=20)
        except requests.RequestException as e:
            log(f"buscar_oportunidades: exceção {e!r}", "ERROR")
            break
        if r.status_code != 200:
            log(f"buscar_oportunidades: HTTP {r.status_code} - {r.text[:200]}", "ERROR")
            break
        data = r.json()
        batch = data.get("opportunities") or []
        opps.extend(batch)
        meta = data.get("meta", {})
        start_after = meta.get("startAfter")
        start_after_id = meta.get("startAfterId")
        if not start_after or not start_after_id or len(batch) < 100:
            break
        params = {
            "location_id": GHL_LOCATION,
            "status": "open",
            "limit": 100,
            "startAfter": start_after,
            "startAfterId": start_after_id,
        }
    return opps


def buscar_tarefas_contato(contact_id: str) -> tuple[str, list[dict]]:
    """
    Retorna (status, tarefas):
      - ("ok",   [...])  -> leitura bem-sucedida com N tarefas
      - ("empty", [])    -> leitura bem-sucedida, sem tarefas
      - ("fail", [])     -> falha de leitura (timeout, 429 pós-retry, 5xx, 4xx inesperado)
    """
    try:
        r = SESSION.get(f"{GHL_BASE}/contacts/{contact_id}/tasks", timeout=15)
    except requests.RequestException as e:
        log(f"buscar_tarefas_contato({contact_id}): exceção {e!r}", "WARNING")
        return ("fail", [])

    if r.status_code == 404:
        # contato inexistente é tratado como "sem tarefa" (produto do dado ruim, não falha)
        return ("empty", [])
    if r.status_code != 200:
        log(
            f"buscar_tarefas_contato({contact_id}): HTTP {r.status_code} - {r.text[:200]}",
            "WARNING",
        )
        return ("fail", [])

    try:
        tarefas = r.json().get("tasks") or []
    except ValueError:
        log(f"buscar_tarefas_contato({contact_id}): JSON inválido", "WARNING")
        return ("fail", [])

    return ("ok" if tarefas else "empty", tarefas)


def calcular_farol(tarefas: list[dict]) -> str:
    pendentes = [t for t in tarefas if not t.get("completed")]
    if not pendentes:
        return "🟡"
    agora = datetime.now(timezone.utc)
    for t in pendentes:
        due = t.get("dueDate")
        if not due:
            continue
        try:
            if isinstance(due, (int, float)):
                dt = datetime.fromtimestamp(due / 1000, tz=timezone.utc)
            else:
                dt = datetime.fromisoformat(str(due).replace("Z", "+00:00"))
            if dt < agora:
                return "🔴"
        except Exception:
            continue
    return "🟢"


def atualizar_oportunidade(opp_id: str, farol: str) -> bool:
    body = {"customFields": [{"id": CAMPO_STATUS_ID, "value": farol}]}
    try:
        r = SESSION.put(f"{GHL_BASE}/opportunities/{opp_id}", json=body, timeout=15)
    except requests.RequestException as e:
        log(f"atualizar_oportunidade({opp_id}): exceção {e!r}", "WARNING")
        return False
    if r.status_code in (200, 201):
        return True
    log(
        f"atualizar_oportunidade({opp_id}): HTTP {r.status_code} - {r.text[:200]}",
        "WARNING",
    )
    return False


# ---------------------------------------------------------------------------
# Entrypoint
# ---------------------------------------------------------------------------

def processar_oportunidade(conn: sqlite3.Connection, opp: dict, counters: dict) -> None:
    opp_id = opp.get("id")
    contact_id = opp.get("contactId") or (opp.get("contact") or {}).get("id")
    if not opp_id or not contact_id:
        counters["puladas_sem_id"] += 1
        return

    status, tarefas = buscar_tarefas_contato(contact_id)
    if status == "fail":
        counters["falhas_leitura"] += 1
        return

    farol = calcular_farol(tarefas)
    anterior = cache_get(conn, opp_id)
    if anterior == farol:
        counters["inalterados"] += 1
        return

    if atualizar_oportunidade(opp_id, farol):
        cache_set(conn, opp_id, farol)
        counters["atualizados"] += 1
    else:
        counters["falhas_escrita"] += 1


def main() -> int:
    start = time.monotonic()
    log("iniciando atualização de faróis PX3")

    opps = buscar_oportunidades()
    log(f"oportunidades encontradas: {len(opps)}")

    conn = _cache_conn()
    counters = {
        "atualizados": 0,
        "inalterados": 0,
        "falhas_leitura": 0,
        "falhas_escrita": 0,
        "puladas_sem_id": 0,
    }

    try:
        for opp in opps:
            processar_oportunidade(conn, opp, counters)
        conn.commit()
    finally:
        conn.close()

    duracao = time.monotonic() - start
    log(
        "resumo: "
        f"atualizados={counters['atualizados']} | "
        f"inalterados={counters['inalterados']} | "
        f"falhas_leitura={counters['falhas_leitura']} | "
        f"falhas_escrita={counters['falhas_escrita']} | "
        f"puladas_sem_id={counters['puladas_sem_id']} | "
        f"duracao_segundos={duracao:.1f}"
    )
    return 0


if __name__ == "__main__":
    sys.exit(main())
