"""
Sweet Angels SDR Webhook — Medio Porte Sorocaba — porta 5057.

Recebe eventos do Linkia (GoHighLevel) quando um lead responde WhatsApp,
processa com o agente Gi (Claude via CLI/assinatura) e devolve a resposta
pelo Linkia.

Diferencas do sdr_fesqua:
  - Porta 5057 (fesqua = 5056).
  - Rota /webhook/medio-porte.
  - System prompt do playbook medio porte (gancho no dia a dia, nao no evento).
  - contact_opp_map carregado de arquivo JSON opcional (contact_opp_map.json)
    — enquanto o Renato nao passar a lista de contatos, o SDR funciona
    respondendo mensagens sem mover cards no pipeline.
  - Pipeline stages configuraveis por env (LINKIA_STAGE_*) porque ainda
    nao existem no Linkia (TODO Renato).

Fluxo:
  1. Webhook chega → responde 200 imediatamente.
  2. Acumula mensagens do mesmo contato num buffer (debounce DEBOUNCE_SECONDS).
  3. Apos o debounce, concatena tudo, chama Claude, aguarda TYPING_DELAY_SECONDS
     (simula digitando) e envia os blocos separados por [[BREAK]] com pausa
     entre eles.
"""

from __future__ import annotations

import json
import logging
import os
import sys
import threading
import time
from logging.handlers import RotatingFileHandler
from pathlib import Path
from typing import Any

from flask import Flask, jsonify, request

BASE_DIR = Path(__file__).resolve().parent
LOG_DIR = BASE_DIR / "logs"
LOG_DIR.mkdir(parents=True, exist_ok=True)

# ---------------------------------------------------------------------------
# Logging
# ---------------------------------------------------------------------------
log_handler = RotatingFileHandler(
    LOG_DIR / "webhook.log",
    maxBytes=5 * 1024 * 1024,
    backupCount=5,
    encoding="utf-8",
)
log_handler.setFormatter(
    logging.Formatter("%(asctime)s [%(levelname)s] %(name)s: %(message)s")
)
stream_handler = logging.StreamHandler(sys.stdout)
stream_handler.setFormatter(
    logging.Formatter("%(asctime)s [%(levelname)s] %(name)s: %(message)s")
)
logging.basicConfig(level=logging.INFO, handlers=[log_handler, stream_handler])
log = logging.getLogger("sdr_medio_porte_webhook")

# ---------------------------------------------------------------------------
# .env loading
# ---------------------------------------------------------------------------
try:
    from dotenv import load_dotenv  # type: ignore
    load_dotenv("/opt/mia/.env")
    load_dotenv(BASE_DIR.parent / "linkia.env")
except ImportError:
    for envfile in ("/opt/mia/.env", str(BASE_DIR.parent / "linkia.env")):
        if os.path.exists(envfile):
            with open(envfile, 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())

# Imports proprios (depois do load_dotenv)
from gi_agent import (  # noqa: E402
    gerar_resposta,
    gerar_followup,
    is_conversa_encerrada,
)
from ghl_client import send_whatsapp_message, move_opportunity_stage  # noqa: E402
from followup_manager import FollowupManager  # noqa: E402

LINKIA_TOKEN = os.getenv("LINKIA_TOKEN", "").strip()
LINKIA_LOCATION_ID = os.getenv("LINKIA_LOCATION_ID", "").strip()

# ---------------------------------------------------------------------------
# Pipeline stages (Medio Porte Sorocaba)
# TODO Renato vai criar pipeline "PROSPECCAO ATIVA MEDIO PORTE SOROCABA" no
# Linkia e passar os stage IDs. Ate la, esses valores ficam vazios e o
# _move_stage_safe simplesmente pula a acao (nao quebra o fluxo).
# ---------------------------------------------------------------------------
STAGE_LEAD_RESPONDEU = os.getenv("LINKIA_STAGE_LEAD_RESPONDEU", "").strip()
STAGE_QUALIFICADO = os.getenv("LINKIA_STAGE_QUALIFICADO", "").strip()

# ---------------------------------------------------------------------------
# Mapping contact_id -> opportunity_id
# TODO Renato vai mandar a lista de contatos/oportunidades (formato similar
# ao fesqua_envio_linkia.json). Enquanto isso, o mapa fica vazio e o SDR
# responde mensagens sem mover cards.
# ---------------------------------------------------------------------------
_contact_opp_map: dict[str, str] = {}
CONTACT_MAP_PATH = BASE_DIR.parent / "medio_porte_envio_linkia.json"


def _load_contact_opp_map() -> None:
    if not CONTACT_MAP_PATH.exists():
        log.info(
            "contact_opp_map: arquivo %s nao existe ainda (esperado ate o Renato passar a lista)",
            CONTACT_MAP_PATH,
        )
        return
    try:
        data = json.loads(CONTACT_MAP_PATH.read_text(encoding="utf-8"))
        for s in data.get("sucessos", []):
            cid = s.get("contact_id")
            oid = s.get("opportunity_id")
            if cid and oid:
                _contact_opp_map[cid] = oid
        log.info("contact_opp_map carregado: %d entradas", len(_contact_opp_map))
    except Exception as e:
        log.warning("nao carregou contact_opp_map: %s", e)


_load_contact_opp_map()


def _move_stage_safe(
    contact_id: str,
    stage_id: str | None,
    status: str | None = None,
    motivo: str = "",
) -> None:
    """Wrapper seguro pra mover card: nao quebra o fluxo se falhar."""
    if not stage_id and not status:
        log.info("skip move stage contact=%s (nem stage_id nem status)", contact_id)
        return
    opp_id = _contact_opp_map.get(contact_id)
    if not opp_id:
        log.info(
            "skip move stage contact=%s (sem opportunity mapeada) motivo=%s",
            contact_id, motivo,
        )
        return
    if not LINKIA_TOKEN:
        log.warning("skip move stage contact=%s (LINKIA_TOKEN ausente)", contact_id)
        return
    try:
        ok, status_code, body = move_opportunity_stage(
            opportunity_id=opp_id,
            stage_id=stage_id,
            token=LINKIA_TOKEN,
            status=status,
        )
        log.info(
            "move stage contact=%s opp=%s stage=%s status=%s motivo=%s -> ok=%s http=%s",
            contact_id, opp_id, stage_id, status, motivo, ok, status_code,
        )
        if not ok:
            log.warning(
                "move stage falhou contact=%s opp=%s body=%s",
                contact_id, opp_id, str(body)[:300],
            )
    except Exception as e:
        log.exception(
            "erro movendo card contact=%s opp=%s motivo=%s: %s",
            contact_id, opp_id, motivo, e,
        )

# Tempo de espera antes de processar (acumula msgs consecutivas do mesmo contato)
DEBOUNCE_SECONDS = float(os.getenv("SDR_DEBOUNCE_SECONDS", "5.0"))
# Pausa antes de enviar o primeiro bloco (simula "digitando...")
TYPING_DELAY_SECONDS = float(os.getenv("SDR_TYPING_DELAY_SECONDS", "3.0"))
# Pausa entre blocos consecutivos
BLOCK_DELAY_SECONDS = float(os.getenv("SDR_BLOCK_DELAY_SECONDS", "4.0"))

# Modo dry-run: NAO envia mensagem real ao GHL, apenas loga o que seria enviado.
# Usado pra testar o fluxo antes de ativar de verdade.
DRY_RUN = os.getenv("SDR_DRY_RUN", "0") == "1"

if not LINKIA_TOKEN:
    log.warning("LINKIA_TOKEN nao configurado — respostas nao serao enviadas ao GHL")
if DRY_RUN:
    log.warning("DRY_RUN ativo — nenhuma mensagem sera enviada ao GHL de verdade")

app = Flask(__name__)

# ---------------------------------------------------------------------------
# Follow-up manager (thread background)
# ---------------------------------------------------------------------------

def _followup_send(contact_id: str, message: str, conversation_id: str | None):
    if DRY_RUN:
        log.info(
            "[DRY_RUN] followup contact=%s preview=%r", contact_id, message[:120]
        )
        return True, 200, "dry_run"
    return send_whatsapp_message(
        contact_id=contact_id,
        message=message,
        token=LINKIA_TOKEN,
        conversation_id=conversation_id,
    )


def _on_followup_encerrado(contact_id: str) -> None:
    """Chamado apos o followup de 24h: move card para perdido."""
    _move_stage_safe(contact_id, None, status="lost", motivo="followup_24h_sem_resposta")


followup = FollowupManager(
    send_fn=_followup_send,
    generate_followup_fn=gerar_followup,
    is_encerrada_fn=is_conversa_encerrada,
    on_conversation_closed_fn=_on_followup_encerrado,
)

# ---------------------------------------------------------------------------
# Debounce / acumulador por contato
# ---------------------------------------------------------------------------
_pending: dict[str, dict] = {}
_pending_lock = threading.Lock()


def _fire_contact(contact_id: str) -> None:
    """Chamado apos o debounce. Concatena msgs, chama Claude e envia blocos."""
    with _pending_lock:
        data = _pending.pop(contact_id, None)
    if not data:
        return

    messages = data["messages"]
    conversation_id = data["conversation_id"]
    contact_meta = data["contact_meta"]
    mensagem_combinada = "\n".join(messages)

    log.info(
        "processando contact=%s msgs=%d combinado=%r",
        contact_id, len(messages), mensagem_combinada[:300],
    )

    try:
        resposta = gerar_resposta(
            contact_id=contact_id,
            mensagem_lead=mensagem_combinada,
            contact_meta=contact_meta,
        )
    except Exception as e:
        log.exception("erro gerando resposta: %s", e)
        return

    if not LINKIA_TOKEN and not DRY_RUN:
        log.error("LINKIA_TOKEN ausente — nao enviando ao GHL")
        return

    # Movimenta card conforme sinal na resposta da Gi
    resposta_lower_pre = resposta.lower()
    if "[[qualificado]]" in resposta_lower_pre:
        _move_stage_safe(contact_id, STAGE_QUALIFICADO, motivo="qualificado")
    elif "@sweetangelsgastronomia" in resposta_lower_pre:
        # Conversa encerrada sem qualificacao -> marca oportunidade como perdida
        _move_stage_safe(contact_id, None, status="lost", motivo="encerramento_sem_qualificar")

    # Divide em blocos (separador [[BREAK]]) e envia com delay entre eles
    blocos = [b.strip() for b in resposta.split("[[BREAK]]") if b.strip()]
    if not blocos:
        log.warning("resposta vazia apos split contact=%s", contact_id)
        return

    # Pausa que simula "digitando..." antes do primeiro bloco
    time.sleep(TYPING_DELAY_SECONDS)

    algum_ok = False
    for i, bloco in enumerate(blocos):
        if i > 0:
            time.sleep(BLOCK_DELAY_SECONDS)
        if DRY_RUN:
            log.info(
                "[DRY_RUN] bloco %d/%d contact=%s preview=%r",
                i + 1, len(blocos), contact_id, bloco[:120],
            )
            ok, status = True, 200
        else:
            ok, status, _body = send_whatsapp_message(
                contact_id=contact_id,
                message=bloco,
                token=LINKIA_TOKEN,
                conversation_id=conversation_id,
            )
        if ok:
            algum_ok = True
        log.info(
            "bloco %d/%d enviado ok=%s status=%s contact=%s preview=%r",
            i + 1, len(blocos), ok, status, contact_id, bloco[:80],
        )

    # Follow-up: registra que a Gi acabou de enviar msg outbound
    if algum_ok:
        followup.registrar_mensagem_enviada(contact_id, conversation_id)

    # Se a resposta gerada ja e um encerramento, marca imediatamente
    resposta_lower = resposta.lower()
    if "@sweetangelsgastronomia" in resposta_lower or "[[qualificado]]" in resposta_lower:
        followup.marcar_encerrado(contact_id)


# ---------------------------------------------------------------------------
# Helpers
# ---------------------------------------------------------------------------

def _get_first(d: dict[str, Any], *keys: str, default: Any = None) -> Any:
    for k in keys:
        if not k:
            continue
        if "." in k:
            cur: Any = d
            ok = True
            for part in k.split("."):
                if isinstance(cur, dict) and part in cur:
                    cur = cur[part]
                else:
                    ok = False
                    break
            if ok and cur not in (None, "", []):
                return cur
        elif k in d and d[k] not in (None, "", []):
            return d[k]
    return default


def extrair_payload(payload: dict[str, Any]) -> dict[str, Any]:
    contact_id = _get_first(
        payload,
        "contactId", "contact_id", "contact.id", "id",
    )
    conversation_id = _get_first(
        payload,
        "conversationId", "conversation_id", "conversation.id",
    )
    mensagem = _get_first(
        payload,
        "body", "message", "messageBody", "text",
        "message.body", "lastMessage.body",
    )
    first_name = _get_first(payload, "firstName", "first_name", "contact.firstName")
    last_name = _get_first(payload, "lastName", "last_name", "contact.lastName")
    phone = _get_first(payload, "phone", "contact.phone")
    empresa = _get_first(payload, "companyName", "company_name", "contact.companyName")
    direction = _get_first(payload, "direction", default="inbound")
    msg_type = _get_first(payload, "type", "messageType", default="WhatsApp")

    return {
        "contact_id": contact_id,
        "conversation_id": conversation_id,
        "mensagem": mensagem,
        "first_name": first_name,
        "last_name": last_name,
        "phone": phone,
        "empresa": empresa,
        "direction": direction,
        "type": msg_type,
    }


# ---------------------------------------------------------------------------
# Rotas
# ---------------------------------------------------------------------------

@app.route("/health", methods=["GET"])
def health():
    return jsonify({
        "ok": True,
        "service": "sweet_angels_sdr_medio_porte",
        "port": int(os.getenv("PORT", "5057")),
        "linkia_token_ok": bool(LINKIA_TOKEN),
        "pipeline_configured": bool(STAGE_LEAD_RESPONDEU and STAGE_QUALIFICADO),
        "contact_opp_map_size": len(_contact_opp_map),
        "dry_run": DRY_RUN,
        "debounce_s": DEBOUNCE_SECONDS,
        "followup_loop": bool(followup._thread and followup._thread.is_alive()),
        "followup_tracked": len(followup._state),
    })


@app.route("/webhook/medio-porte", methods=["POST"])
def webhook_medio_porte():
    # 1. Parse payload
    try:
        payload = request.get_json(force=True, silent=False) or {}
    except Exception as e:
        log.warning("payload nao-JSON: %s | raw=%s", e, request.data[:400])
        return jsonify({"ok": False, "error": "invalid_json"}), 400

    log.info("webhook recebido: %s", json.dumps(payload, ensure_ascii=False)[:800])

    dados = extrair_payload(payload)

    # 2. Validacoes
    if not dados["contact_id"]:
        log.warning("webhook sem contact_id | payload=%s",
                    json.dumps(payload, ensure_ascii=False)[:400])
        return jsonify({"ok": False, "error": "missing_contact_id"}), 200

    if not dados["mensagem"]:
        log.info("webhook sem mensagem util — ignorando")
        return jsonify({"ok": True, "skipped": "no_message"}), 200

    if str(dados["direction"]).lower() == "outbound":
        log.info("mensagem outbound — ignorando (evita loop)")
        return jsonify({"ok": True, "skipped": "outbound"}), 200

    contact_id = str(dados["contact_id"])
    conversation_id = str(dados["conversation_id"]) if dados["conversation_id"] else None
    mensagem_lead = str(dados["mensagem"]).strip()
    contact_meta = {
        "first_name": dados.get("first_name"),
        "last_name": dados.get("last_name"),
        "phone": dados.get("phone"),
        "empresa": dados.get("empresa"),
        "contact_id": contact_id,
        "conversation_id": conversation_id,
    }

    log.info(
        "recebido contact=%s nome=%s empresa=%s msg=%r",
        contact_id, dados.get("first_name"), dados.get("empresa"), mensagem_lead[:200],
    )

    # Follow-up: lead respondeu -> reseta timers e persiste conversation_id
    followup.registrar_resposta_recebida(contact_id)
    if conversation_id:
        followup.set_conversation_id(contact_id, conversation_id)

    # Move card para LEAD RESPONDEU (em thread pra nao atrasar o 200 do webhook)
    if STAGE_LEAD_RESPONDEU and _contact_opp_map.get(contact_id):
        threading.Thread(
            target=_move_stage_safe,
            args=(contact_id, STAGE_LEAD_RESPONDEU, None, "inbound"),
            daemon=True,
        ).start()

    # 3. Adiciona ao buffer de debounce (reseta o timer se ja existe)
    with _pending_lock:
        if contact_id in _pending:
            existing = _pending[contact_id]
            if existing.get("timer"):
                existing["timer"].cancel()
            existing["messages"].append(mensagem_lead)
            existing["conversation_id"] = conversation_id
            existing["contact_meta"] = contact_meta
            log.info("debounce reset contact=%s total_msgs=%d", contact_id, len(existing["messages"]))
        else:
            _pending[contact_id] = {
                "messages": [mensagem_lead],
                "timer": None,
                "conversation_id": conversation_id,
                "contact_meta": contact_meta,
            }

        timer = threading.Timer(DEBOUNCE_SECONDS, _fire_contact, args=(contact_id,))
        timer.daemon = True
        _pending[contact_id]["timer"] = timer
        timer.start()

    return jsonify({"ok": True, "queued": True, "debounce_s": DEBOUNCE_SECONDS}), 200


@app.route("/", methods=["GET"])
def index():
    return jsonify({
        "service": "Sweet Angels SDR Webhook — Medio Porte Sorocaba",
        "endpoints": {
            "health": "GET /health",
            "webhook": "POST /webhook/medio-porte",
        },
    })


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

# Sobe a thread de follow-up assim que o modulo eh importado (funciona tanto
# rodando via `python app.py` quanto via gunicorn/systemd).
followup.start()

if __name__ == "__main__":
    port = int(os.getenv("PORT", "5057"))
    host = os.getenv("HOST", "0.0.0.0")
    log.info(
        "subindo Flask em %s:%s (debounce=%.1fs, dry_run=%s)",
        host, port, DEBOUNCE_SECONDS, DRY_RUN,
    )
    app.run(host=host, port=port, debug=False, threaded=True)
