#!/usr/bin/env python3
"""
Interface Web - Prospecção Ativa via Google Maps
Flask app rodando na porta 5050
"""

import sys
import os
# path[0] = diretório deste próprio app.py (permite staging/prod isolados).
# Antes era hardcoded em prod, o que fazia staging importar auth/config do
# diretório de produção e escrever no DB errado (bug detectado no Sprint 1 F2).
from pathlib import Path as _Path
sys.path.insert(0, str(_Path(__file__).resolve().parent))

import uuid
import time
import threading
import requests
import re
import json
import csv
import io
from datetime import datetime, timezone, timedelta
from pathlib import Path
from flask import Flask, jsonify, render_template, request, redirect, url_for, flash
from flask_login import LoginManager, login_user, logout_user, login_required, current_user
from auth import (init_db, buscar_por_email, criar_usuario, listar_usuarios,
                  atualizar_usuario, deletar_usuario, verificar_senha, registrar_acesso, User,
                  get_ghl_config, save_ghl_config,
                  criar_conta, listar_contas, get_conta, deletar_conta, renomear_conta,
                  get_quota, incrementar_quota, resetar_quota_se_novo_mes,
                  salvar_job_historico, listar_jobs_historico,
                  listar_perfis_crm, get_perfil_crm, salvar_perfil_crm,
                  deletar_perfil_crm, contar_perfis_crm, migrar_ghl_config_para_perfis,
                  buscar_por_token, gerar_api_token, set_ip_whitelist)
from dashboard_saude import get_dashboard_saude
from functools import wraps
import json as _json_api

app = Flask(__name__)
# Chave de sessão vem do ambiente (.env). Sem ela, cai num fallback só para não travar,
# mas em produção FLASK_SECRET_KEY deve estar sempre definida (systemd EnvironmentFile).
app.secret_key = os.environ.get("FLASK_SECRET_KEY") or "prospeccao-climb-2026-secret"

# Segredo opcional para autenticar webhooks externos. Enquanto vazio, os webhooks
# seguem abertos (comportamento atual — não quebra Dinastia/GHL). Ao definir
# WEBHOOK_SECRET no .env, todo webhook passa a exigir ?s=<segredo> ou header X-Webhook-Secret.
WEBHOOK_SECRET = os.environ.get("WEBHOOK_SECRET", "").strip()


def _webhook_autorizado() -> bool:
    """True se o webhook pode prosseguir. Se WEBHOOK_SECRET estiver vazio, sempre libera."""
    if not WEBHOOK_SECRET:
        return True
    fornecido = (request.args.get("s") or request.headers.get("X-Webhook-Secret") or "").strip()
    return fornecido == WEBHOOK_SECRET


# ── Rate-limit simples de login (em memória, expira sozinho — sem dependência extra) ──
from collections import defaultdict
import threading as _threading_rl

_LOGIN_FAILS = defaultdict(list)  # ip -> [timestamps de falha]
_LOGIN_FAILS_LOCK = _threading_rl.Lock()
_LOGIN_MAX_FAILS = 10           # falhas permitidas por janela
_LOGIN_WINDOW_SEG = 15 * 60     # janela de 15 minutos


def _login_ip() -> str:
    fwd = (request.headers.get("X-Forwarded-For") or "").split(",")[0].strip()
    return fwd or (request.remote_addr or "desconhecido")


def _login_bloqueado(ip: str) -> bool:
    agora = time.time()
    with _LOGIN_FAILS_LOCK:
        recentes = [t for t in _LOGIN_FAILS[ip] if agora - t < _LOGIN_WINDOW_SEG]
        _LOGIN_FAILS[ip] = recentes
        return len(recentes) >= _LOGIN_MAX_FAILS


def _login_registrar_falha(ip: str):
    with _LOGIN_FAILS_LOCK:
        _LOGIN_FAILS[ip].append(time.time())


def _login_limpar(ip: str):
    with _LOGIN_FAILS_LOCK:
        _LOGIN_FAILS.pop(ip, None)

login_manager = LoginManager(app)
login_manager.login_view = "login"
login_manager.login_message = "Faça login para acessar a ferramenta."
login_manager.login_message_category = "error"

@login_manager.user_loader
def load_user(user_id):
    return User.get(user_id)

# ═════════════════════════════════════════════════════════════════════════
# Sprint 1 Fase 2 ClimbLeads — signup público, Asaas, cota, LGPD
# Módulo isolado em climbleads_signup.py. Falha silenciosa se módulo ausente
# (não quebra deploy de emergência que rode com app.py antigo).
# ═════════════════════════════════════════════════════════════════════════
try:
    from climbleads_signup import register_climbleads_routes
    register_climbleads_routes(app)
except Exception as _e_cl:
    print(f"[climbleads] AVISO: módulo não carregado ({_e_cl})")

# Armazenamento em memória dos jobs
jobs = {}
MAX_JOBS_HISTORY = 10

# Registro persistente de empresas já prospectadas — MULTI-TENANT via SQLite
# (users.db tabela `prospectados` com FK conta_id). O JSON legado NÃO é mais
# lido nem escrito: substituído pela tabela em 2026-08-12 (Fase 1 ClimbLeads).
_USERS_DB_PATH = Path(__file__).resolve().parent / "users.db"

def _users_db_conn():
    """Abre conexão nova no users.db. Cada thread abre a própria."""
    import sqlite3
    return sqlite3.connect(str(_USERS_DB_PATH), timeout=15)

def _carregar_registro(conta_id) -> set:
    """Carrega chaves de empresas já prospectadas PELA CONTA especificada.
    conta_id None ou inválido → retorna set vazio (fail-safe)."""
    if not conta_id:
        return set()
    try:
        conn = _users_db_conn()
        cur = conn.cursor()
        rows = cur.execute(
            "SELECT chave FROM prospectados WHERE conta_id = ?",
            (int(conta_id),)
        ).fetchall()
        conn.close()
        return {r[0] for r in rows}
    except Exception as e:
        print(f"[prospectados] erro ao carregar conta_id={conta_id}: {e}")
        return set()

def _salvar_registro(conta_id, chave: str, empresa: str = "", cidade: str = "",
                     telefone: str = "", job_id: str = None):
    """Insere UMA chave no registro da conta. Idempotente (UNIQUE conta_id+chave).
    Substituiu a versão antiga que gravava JSON global."""
    if not conta_id or not chave:
        return
    try:
        conn = _users_db_conn()
        cur = conn.cursor()
        cur.execute(
            "INSERT OR IGNORE INTO prospectados (conta_id, chave, empresa, cidade, telefone, criado_em, job_id) "
            "VALUES (?,?,?,?,?,?,?)",
            (int(conta_id), chave, empresa or "", cidade or "", telefone or "",
             datetime.now().isoformat(), job_id)
        )
        conn.commit()
        conn.close()
    except Exception as e:
        print(f"[prospectados] erro ao salvar conta_id={conta_id} chave={chave!r}: {e}")

def _chave_empresa(nome: str, cidade: str) -> str:
    return f"{nome.lower().strip()}|{cidade.lower().strip()}"

GHL_BASE = "https://services.leadconnectorhq.com"
UNNICHAT_BASE = "https://unnichat.com.br"


def _ghl_headers(token: str) -> dict:
    return {
        "Authorization": f"Bearer {token}",
        "Content-Type": "application/json",
        "Version": "2021-07-28"
    }


# ── UNNICHAT ──────────────────────────────────────────

def unnichat_adicionar_pipeline(contact_id: str, token: str, connection_id: str,
                                pipeline_id: str, column_id: str, business_name: str) -> tuple:
    """Adiciona contato a uma coluna do pipeline no Unnichat."""
    headers = {"Authorization": f"Bearer {token}", "Content-Type": "application/json"}
    if connection_id:
        headers["x-connection-id"] = connection_id
    payload = {
        "crm_id": pipeline_id,
        "column_id": column_id,
        "business_name": business_name,
    }
    try:
        r = requests.post(
            f"{UNNICHAT_BASE}/api/contact/{contact_id}/crm",
            headers=headers, json=payload, timeout=15
        )
        if r.status_code in (200, 201):
            return True, None
        try:
            msg = r.json().get("message") or r.text[:120]
        except Exception:
            msg = r.text[:120]
        return False, f"HTTP {r.status_code}: {msg}"
    except Exception as e:
        return False, str(e)


def unnichat_criar_contato(lead: dict, token: str, tags: list, connection_id: str = "") -> tuple:
    """Cria contato no Unnichat. Retorna (contact_id, erro)."""
    telefone = re.sub(r"[^\d+]", "", lead.get("telefone", ""))
    payload = {
        "name": lead["nome"],
        "phone": telefone or None,
        "email": lead.get("email") or None,
        "tags": tags or ["prospeccao-ativa"],
    }
    payload = {k: v for k, v in payload.items() if v is not None}
    headers = {"Authorization": f"Bearer {token}", "Content-Type": "application/json"}
    if connection_id:
        headers["x-connection-id"] = connection_id
    try:
        r = requests.post(
            f"{UNNICHAT_BASE}/api/contact",
            headers=headers,
            json=payload,
            timeout=15
        )
        if r.status_code in (200, 201):
            data = r.json().get("data", {})
            return data.get("id"), None
        try:
            msg = r.json().get("message") or r.text[:120]
        except Exception:
            msg = r.text[:120]
        return None, f"HTTP {r.status_code}: {msg}"
    except Exception as e:
        return None, str(e)


def buscar_contato_existente_ghl(telefone: str, nome: str, token: str, location_id: str) -> str | None:
    """Busca contato existente por telefone ou nome."""
    headers = _ghl_headers(token)

    # Tenta por telefone primeiro (mais preciso)
    if telefone:
        try:
            r = requests.get(
                f"{GHL_BASE}/contacts/",
                headers=headers,
                params={"locationId": location_id, "query": telefone, "limit": 1},
                timeout=10
            )
            if r.status_code == 200:
                contacts = r.json().get("contacts", [])
                if contacts:
                    return contacts[0].get("id")
        except Exception:
            pass

    # Tenta por nome
    try:
        r = requests.get(
            f"{GHL_BASE}/contacts/",
            headers=headers,
            params={"locationId": location_id, "query": nome[:50], "limit": 5},
            timeout=10
        )
        if r.status_code == 200:
            contacts = r.json().get("contacts", [])
            # Busca correspondência exata por companyName
            for c in contacts:
                if (c.get("companyName") or "").lower() == nome.lower():
                    return c.get("id")
            # Se não encontrou exato, retorna o primeiro
            if contacts:
                return contacts[0].get("id")
    except Exception:
        pass

    return None


def _carregar_campos_crm(location_id: str) -> dict:
    """Carrega mapeamento de campos customizados por location_id."""
    mapa = {
        "4z6Fjpgj2soAQZpu8Ddk": "/opt/mia/workspace/clientes/sweet_angels/campos_crm.json",
    }
    path = mapa.get(location_id)
    if path and os.path.exists(path):
        try:
            with open(path) as f:
                data = json.load(f)
            # Remove _meta e retorna só campo -> id
            return {k: v for k, v in data.items() if k != "_meta" and isinstance(v, str)}
        except Exception:
            pass
    return {}


def criar_contato_ghl(lead: dict, token: str, location_id: str, source: str = "Google Maps", custom_fields_map: dict = None) -> tuple:
    """Cria contato no CRM e retorna (contact_id, erro). Se duplicado, retorna o ID existente."""
    nome_split = lead["nome"].split(" ", 1)
    primeiro = nome_split[0]
    sobrenome = nome_split[1] if len(nome_split) > 1 else ""

    telefone = re.sub(r"[^\d+]", "", lead.get("telefone", ""))
    tags = lead.get("tags") or ["prospeccao-ativa", lead.get("nicho", ""), lead.get("cidade", "")]
    tags = [t for t in tags if t]

    payload = {
        "locationId": location_id,
        "firstName": primeiro,
        "lastName": sobrenome,
        "companyName": lead.get("empresa") or lead["nome"],
        "phone": telefone or None,
        "email": lead.get("email") or None,
        "website": lead.get("site") or None,
        "address1": lead.get("endereco") or None,
        "tags": tags,
        "source": source,
    }
    payload = {k: v for k, v in payload.items() if v is not None}

    # Campos customizados por CRM (ex.: Sweet Angels)
    custom_fields = []
    if custom_fields_map:
        for campo, field_id in custom_fields_map.items():
            if campo == "fontes_enriquecimento":
                fontes = lead.get("fontes") or []
                if fontes:
                    custom_fields.append({"id": field_id, "value": ", ".join(fontes)})
                continue
            valor = lead.get(campo)
            if valor:
                custom_fields.append({"id": field_id, "value": str(valor)})
    if custom_fields:
        payload["customFields"] = custom_fields

    r = requests.post(
        f"{GHL_BASE}/contacts/",
        headers=_ghl_headers(token),
        json=payload,
        timeout=15
    )

    if r.status_code in (200, 201):
        return r.json().get("contact", {}).get("id"), None

    # Se duplicado, busca o contato existente e usa o ID dele
    try:
        msg = r.json().get("message") or r.json().get("msg") or ""
    except Exception:
        msg = r.text[:200]

    if r.status_code == 400 and "duplicate" in msg.lower():
        contact_id = buscar_contato_existente_ghl(telefone, lead["nome"], token, location_id)
        if contact_id:
            return contact_id, None
        return None, f"Duplicado mas não encontrado na busca"

    return None, f"HTTP {r.status_code}: {msg[:120]}"


def criar_oportunidade_ghl(lead: dict, contact_id: str, token: str, location_id: str,
                           pipeline_id: str, stage_id: str, source: str = "Google Maps") -> tuple:
    """Cria oportunidade no pipeline. Retorna (ok, opportunity_id, erro)."""
    payload = {
        "pipelineId": pipeline_id,
        "locationId": location_id,
        "name": lead.get("empresa") or lead["nome"],
        "pipelineStageId": stage_id,
        "status": "open",
        "contactId": contact_id,
        "source": source,
    }

    r = requests.post(
        f"{GHL_BASE}/opportunities/",
        headers=_ghl_headers(token),
        json=payload,
        timeout=15
    )
    if r.status_code in (200, 201):
        opp_id = r.json().get("opportunity", {}).get("id") or r.json().get("id")
        return True, opp_id, None
    try:
        erro = r.json().get("message") or r.json().get("msg") or str(r.json())
    except Exception:
        erro = r.text[:120]
    return False, None, f"HTTP {r.status_code}: {erro}"


def calcular_score_gmb(lead: dict) -> int:
    """Calcula score 0-100 do perfil GMB com base em rating, avaliações, descrição, categoria e fotos."""
    score = 0
    # Avaliação média (estrelas) — peso 40pts
    try:
        rating = float(lead.get("avaliacao") or 0)
        score += int((rating / 5) * 40)
    except Exception:
        pass
    # Volume de avaliações — peso 30pts
    try:
        n = int(lead.get("num_avaliacoes") or 0)
    except Exception:
        n = 0
    if n >= 200:
        score += 30
    elif n >= 50:
        score += 20
    elif n >= 10:
        score += 10
    # Tem descrição/editorial — peso 15pts
    if lead.get("descricao"):
        score += 15
    # Tem categoria/segmento — peso 10pts
    if lead.get("categoria"):
        score += 10
    # Tem foto de perfil — peso 5pts
    if lead.get("tem_fotos"):
        score += 5
    return min(score, 100)


def run_prospector(job_id: str, nicho: str, cidade: str, limite: int,
                   token: str, location_id: str, pipeline_id: str, stage_id: str,
                   crm_type: str = "ghl", connection_id: str = "",
                   pais: str = "BR", enriquecimento_auto: bool = False,
                   conta_id: int = None, filtro_sem_site: bool = False,
                   filtro_nota: str = "", filtro_site_gmb: str = "",
                   score_minimo: int = 0,
                   lat: float = None, lng: float = None, raio_km: int = 0):
    """Executa o scraper em thread separada com callbacks de progresso.

    Nota: os parâmetros filtro_nota, filtro_site_gmb e score_minimo permanecem
    na assinatura por compatibilidade, mas NÃO são mais aplicados durante a coleta.
    O score continua sendo calculado e salvo em cada lead; a filtragem por
    nota/site/score foi movida pra tabela de resultados (pós-coleta).
    """
    job = jobs[job_id]
    job["status"] = "em_andamento"
    job["pais"] = pais
    job["conta_id"] = conta_id
    preposicao = "in" if pais == "US" else "em"
    job["logs"].append(f"[{_ts()}] Iniciando busca: {nicho} {preposicao} {cidade} (limite: {limite})")

    try:
        leads_coletados = []

        def scrape_com_progresso():
            try:
                from config import GOOGLE_PLACES_KEY
            except ImportError:
                GOOGLE_PLACES_KEY = ""

            if GOOGLE_PLACES_KEY:
                return _buscar_via_places_api(GOOGLE_PLACES_KEY)
            else:
                return _buscar_via_playwright()

        def _buscar_via_places_api(api_key: str) -> list:
            """Busca empresas via Google Places API (New) — oficial, sem risco de bloqueio."""
            PLACES_URL = "https://places.googleapis.com/v1/places:searchText"
            FIELD_MASK = "places.displayName,places.formattedAddress,places.nationalPhoneNumber,places.websiteUri,places.rating,places.userRatingCount,places.id,places.editorialSummary,places.primaryTypeDisplayName,places.photos,places.businessStatus"

            preposicao_q = "in" if pais == "US" else "em"
            # Quando o usuário fornece área no mapa, o "cidade" vira contexto do
            # texto (ex: "Maringá, PR") mas quem manda mesmo é o círculo. Ainda
            # incluímos no textQuery pra melhorar relevância semântica.
            query = f"{nicho} {preposicao_q} {cidade}"
            if lat and lng and raio_km:
                job["logs"].append(
                    f"[{_ts()}] Google Places API: buscando '{query}' | círculo {raio_km}km em ({lat:.4f}, {lng:.4f})"
                )
            else:
                job["logs"].append(f"[{_ts()}] Google Places API: buscando '{query}'...")

            leads = []
            registro = _carregar_registro(conta_id)
            pulados = 0
            page_token = None
            coletados = 0

            while coletados < limite:
                if job["status"] in ("cancelado", "pausado"):
                    break

                body = {"textQuery": query, "maxResultCount": 20}
                if page_token:
                    body["pageToken"] = page_token

                # locationBias.circle: Places API (New) v1 aceita circle apenas
                # em locationBias (bias, nao hard-restrict). Raio maximo: 50000m.
                if lat and lng and raio_km:
                    raio_m = min(float(raio_km) * 1000.0, 50000.0)
                    body["locationBias"] = {
                        "circle": {
                            "center": {"latitude": float(lat), "longitude": float(lng)},
                            "radius": raio_m,
                        }
                    }

                try:
                    r = requests.post(
                        PLACES_URL,
                        headers={
                            "Content-Type": "application/json",
                            "X-Goog-Api-Key": api_key,
                            "X-Goog-FieldMask": FIELD_MASK,
                        },
                        json=body,
                        timeout=15
                    )
                    if r.status_code != 200:
                        job["logs"].append(f"[{_ts()}] Places API erro: {r.status_code} {r.text[:100]}")
                        break

                    data = r.json()
                    places = data.get("places", [])
                    if not places:
                        break

                    job["total"] = min(limite, job["total"] + len(places))

                    for place in places:
                        if job["status"] in ("cancelado", "pausado") or coletados >= limite:
                            break

                        nome = place.get("displayName", {}).get("text", "").strip()
                        if not nome or len(nome) < 2:
                            continue

                        chave = _chave_empresa(nome, cidade)
                        if chave in registro:
                            pulados += 1
                            continue

                        rating_val = place.get("rating", 0) or 0
                        try:
                            rating = float(rating_val)
                        except Exception:
                            rating = 0.0

                        # Heurística de perfil verificado no Google:
                        # A API do Places não expõe diretamente o campo "verified".
                        # Usamos negócio OPERATIONAL + tem fotos + tem avaliações como proxy
                        # de um perfil ativo/completo (provável selo verificado).
                        business_status = place.get("businessStatus", "") or ""
                        tem_fotos_gmb = bool(place.get("photos"))
                        num_aval_gmb = place.get("userRatingCount", 0) or 0
                        verificado_gmb = (
                            business_status == "OPERATIONAL"
                            and tem_fotos_gmb
                            and num_aval_gmb > 0
                        )

                        lead = {
                            "nome": nome,
                            "telefone": place.get("nationalPhoneNumber", "").strip(),
                            "site": place.get("websiteUri", "").strip(),
                            "endereco": place.get("formattedAddress", "").strip(),
                            "avaliacao": str(place.get("rating", "")),
                            "num_avaliacoes": num_aval_gmb,
                            "descricao": (place.get("editorialSummary") or {}).get("text", "") or "",
                            "categoria": (place.get("primaryTypeDisplayName") or {}).get("text", "") or "",
                            "tem_fotos": tem_fotos_gmb,
                            "business_status": business_status,
                            "verificado_gmb": verificado_gmb,
                            "nicho": nicho,
                            "cidade": cidade,
                            "enviado_crm": False,
                            "contact_id": None,
                        }
                        # Score continua sendo calculado e salvo em cada lead —
                        # o descarte por score mínimo durante a coleta foi removido.
                        # Usuário filtra depois na tabela.
                        lead["score_gmb"] = calcular_score_gmb(lead)

                        registro.add(chave)
                        _salvar_registro(conta_id, chave, empresa=nome, cidade=cidade,
                                         telefone=lead.get("telefone", ""), job_id=job_id)
                        leads.append(lead)
                        coletados += 1
                        job["progresso"] = coletados
                        job["leads"].append(lead)
                        job["logs"].append(
                            f"[{_ts()}] Coletado ({coletados}): {nome} | {lead['telefone'] or 'sem tel'} | score {lead['score_gmb']}"
                        )
                        if len(job["logs"]) > 50:
                            job["logs"] = job["logs"][-50:]

                    page_token = data.get("nextPageToken")
                    if not page_token:
                        break

                    time.sleep(2)

                except Exception as e:
                    job["logs"].append(f"[{_ts()}] Places API exceção: {str(e)[:100]}")
                    break

            if pulados > 0:
                job["logs"].append(f"[{_ts()}] {pulados} empresa(s) puladas (já prospectadas antes)")

            return leads

        def _buscar_via_playwright() -> list:
            """Fallback: busca via scraping do Google Maps com Playwright."""
            from playwright.sync_api import sync_playwright

            leads = []
            preposicao_q = "in" if pais == "US" else "em"
            query = f"{nicho} {preposicao_q} {cidade}"
            lang_param = "hl=en&gl=us" if pais == "US" else "hl=pt-BR&gl=br"

            with sync_playwright() as p:
                browser = p.chromium.launch(headless=True)
                page = browser.new_page()

                job["logs"].append(f"[{_ts()}] Buscando (scraping): {query}")
                page.goto(f"https://www.google.com/maps/search/{query.replace(' ', '+')}?{lang_param}")
                page.wait_for_timeout(3000)

                if job["status"] == "cancelado":
                    browser.close()
                    return leads

                scrolls = min(limite // 5, 20)
                for i in range(scrolls):
                    if job["status"] == "cancelado":
                        break
                    try:
                        panel = page.locator('div[role="feed"]').first
                        panel.evaluate("el => el.scrollTop += 1500")
                        page.wait_for_timeout(1500)
                    except Exception:
                        break

                cards = page.locator('a[href*="/maps/place/"]').all()
                job["total"] = len(cards)
                job["logs"].append(f"[{_ts()}] {len(cards)} resultados encontrados, coletando até {limite} novos...")

                registro = _carregar_registro(conta_id)
                pulados = 0

                for idx, card in enumerate(cards):
                    if job["status"] in ("cancelado", "pausado") or len(leads) >= limite:
                        break
                    try:
                        lead = {}
                        try:
                            aria = card.get_attribute("aria-label", timeout=1000) or ""
                            lead["nome"] = aria.strip().split('\n')[0] if aria else card.inner_text(timeout=1000).strip().split('\n')[0]
                        except Exception:
                            continue

                        if not lead["nome"] or len(lead["nome"]) < 2:
                            continue

                        chave = _chave_empresa(lead["nome"], cidade)
                        if chave in registro:
                            pulados += 1
                            continue

                        card.click()
                        page.wait_for_timeout(2000)

                        for seletor, campo, transform in [
                            ('[data-item-id^="phone"]', "telefone", lambda v: v.replace("phone:tel:", "").strip()),
                            ('[data-item-id="authority"]', "site", lambda v: v or ""),
                            ('[data-item-id^="address"]', "endereco", lambda v: v),
                        ]:
                            try:
                                el = page.locator(seletor).first
                                attr = el.get_attribute("data-item-id" if "phone" in seletor or "authority" in seletor else "data-item-id", timeout=2000)
                                if campo == "site":
                                    lead[campo] = el.get_attribute("href", timeout=2000) or ""
                                elif campo == "endereco":
                                    lead[campo] = el.inner_text(timeout=2000).strip()
                                else:
                                    lead[campo] = transform(attr or "")
                            except Exception:
                                lead[campo] = ""

                        lead.setdefault("avaliacao", "")

                        # Tenta capturar número de avaliações (texto próximo ao rating)
                        num_aval = 0
                        try:
                            aval_txt = page.locator('button[jsaction*="pane.rating.moreReviews"]').first.inner_text(timeout=1500)
                            m = re.search(r"([\d\.,]+)", aval_txt or "")
                            if m:
                                num_aval = int(m.group(1).replace(".", "").replace(",", ""))
                        except Exception:
                            pass
                        lead["num_avaliacoes"] = num_aval

                        # Tenta capturar rating do card (aria-label costuma trazer estrelas)
                        try:
                            if not lead.get("avaliacao"):
                                m_star = re.search(r"(\d[\.,]\d)\s*(?:estrelas|stars|/5)?", aria or "")
                                if m_star:
                                    lead["avaliacao"] = m_star.group(1).replace(",", ".")
                        except Exception:
                            pass

                        # Descrição / editorial summary (quando aparece no painel)
                        descricao = ""
                        for sel in [
                            '[data-attrid="kc:/local:summary"]',
                            'div[jsaction*="pane.editorialSummary"]',
                        ]:
                            try:
                                descricao = page.locator(sel).first.inner_text(timeout=1000).strip()
                                if descricao:
                                    break
                            except Exception:
                                continue
                        lead["descricao"] = descricao

                        # Categoria (botão sob o nome do lugar)
                        try:
                            lead["categoria"] = page.locator('button[jsaction*="pane.rating.category"]').first.inner_text(timeout=1000).strip()
                        except Exception:
                            lead["categoria"] = ""

                        # Foto: se existe o botão de galeria
                        try:
                            lead["tem_fotos"] = page.locator('button[jsaction*="pane.heroHeaderImage"]').count() > 0
                        except Exception:
                            lead["tem_fotos"] = False

                        # Detecção de perfil verificado no Google (checkmark)
                        # O GMN marca perfis verificados com um ícone/tooltip.
                        # Tentamos localizar por aria-label ou tooltip que contenha "verificado"/"verified".
                        try:
                            verif_el = page.locator(
                                '[aria-label*="verificado" i], [aria-label*="verified" i], '
                                '[data-tooltip*="verificado" i], [data-tooltip*="verified" i]'
                            ).first
                            lead["verificado_gmb"] = verif_el.is_visible(timeout=1000)
                        except Exception:
                            lead["verificado_gmb"] = False

                        lead["nicho"] = nicho
                        lead["cidade"] = cidade
                        lead["enviado_crm"] = False
                        lead["contact_id"] = None
                        # Score continua sendo calculado e salvo em cada lead —
                        # o descarte por score mínimo durante a coleta foi removido.
                        lead["score_gmb"] = calcular_score_gmb(lead)

                        registro.add(chave)
                        _salvar_registro(conta_id, chave, empresa=lead["nome"], cidade=cidade,
                                         telefone=lead.get("telefone", ""), job_id=job_id)
                        leads.append(lead)
                        job["progresso"] = len(leads)
                        job["leads"].append(lead)
                        job["logs"].append(
                            f"[{_ts()}] Coletado ({idx+1}): {lead['nome']} | {lead.get('telefone') or 'sem tel'} | score {lead['score_gmb']}"
                        )
                        if len(job["logs"]) > 50:
                            job["logs"] = job["logs"][-50:]

                    except Exception as e:
                        job["logs"].append(f"[{_ts()}] Aviso card {idx+1}: {str(e)[:80]}")
                        continue

                if pulados > 0:
                    job["logs"].append(f"[{_ts()}] {pulados} empresa(s) puladas (já prospectadas antes)")
                browser.close()

            return leads

        leads_coletados = scrape_com_progresso()

        if job["status"] == "cancelado":
            job["logs"].append(f"[{_ts()}] Job cancelado pelo usuario")
            return

        if job["status"] == "pausado":
            job["logs"].append(f"[{_ts()}] Pausado pelo usuario - enviando {len(leads_coletados)} leads coletados ao CRM...")

        # Todos os leads coletados — para exibição na tabela (sem filtros de CRM)
        todos_leads_para_exibir = list(leads_coletados)

        # ── Filtro "somente leads sem site" — aplica APENAS no envio ao CRM ──
        if filtro_sem_site and leads_coletados:
            antes = len(leads_coletados)
            leads_coletados = [
                l for l in leads_coletados
                if not (l.get("site") or "").strip()
            ]
            pulados_site = antes - len(leads_coletados)
            job["logs"].append(
                f"[{_ts()}] Filtro 'sem site': {pulados_site} com site ignorados, "
                f"{len(leads_coletados)} vão ao CRM"
            )

        if not todos_leads_para_exibir:
            job["status"] = "concluido"
            job["concluido_em"] = datetime.now().isoformat()
            job["logs"].append(f"[{_ts()}] Nenhum lead encontrado")
            try:
                salvar_job_historico(job)
            except Exception as e:
                print(f"[jobs_historico] erro ao salvar conclusao {job.get('id')}: {e}")
            return

        if enriquecimento_auto:
            job["logs"].append(f"[{_ts()}] {len(todos_leads_para_exibir)} leads coletados. Enviando {len(leads_coletados)} ao CRM...")
        else:
            job["logs"].append(f"[{_ts()}] {len(todos_leads_para_exibir)} leads coletados. Modo Análise: leads ficam na tabela para análise lead a lead.")

        # ── Enriquecimento automático (BR: CNPJ/Receita | US: site + DuckDuckGo) ──
        if enriquecimento_auto:
            job["logs"].append(f"[{_ts()}] Enriquecendo leads ({pais}) antes de enviar ao CRM...")
            for i, lead in enumerate(leads_coletados):
                try:
                    if pais == "US":
                        dados = enriquecer_empresa_us(
                            lead["nome"],
                            website=lead.get("site", ""),
                            cidade=cidade
                        )
                    else:
                        dados = enriquecer_empresa(lead["nome"])

                    if dados.get("telefone") and not lead.get("telefone"):
                        lead["telefone"] = dados["telefone"]
                    if dados.get("email") and not lead.get("email"):
                        lead["email"] = dados["email"]
                    if dados.get("cnpj"):
                        lead["cnpj"] = dados["cnpj"]
                    if dados.get("decisor_nome"):
                        lead["decisor_nome"] = dados["decisor_nome"]
                        lead["decisor_cargo"] = dados.get("decisor_cargo", "")
                    if dados.get("linkedin_url"):
                        lead["linkedin_url"] = dados["linkedin_url"]
                    if dados.get("site") and not lead.get("site"):
                        lead["site"] = dados["site"]
                    if dados.get("instagram"):
                        lead["instagram"] = dados["instagram"]
                    if dados.get("facebook"):
                        lead["facebook"] = dados["facebook"]
                    fontes = ", ".join(dados.get("fontes", []))
                    job["logs"].append(f"[{_ts()}] Enriq ({i+1}/{len(leads_coletados)}): {lead['nome']} [{fontes or 'sem dados extras'}]")
                except Exception as e:
                    job["logs"].append(f"[{_ts()}] Enriq erro para {lead['nome']}: {e}")
            job["logs"].append(f"[{_ts()}] Enriquecimento concluído. Enviando ao CRM...")

        enviados = 0
        job["status"] = "em_andamento"

        # Só envia ao CRM se o modo foi marcado. Caso contrário (modo Análise),
        # os leads ficam salvos no job/tabela pra análise lead a lead.
        if enriquecimento_auto:
            for i, lead in enumerate(leads_coletados):
                if job["status"] == "cancelado":
                    break

                if crm_type == "unnichat":
                    # ── Unnichat: cria contato + adiciona ao pipeline ──
                    contact_id, err_contato = unnichat_criar_contato(
                        lead, token, ["prospeccao-ativa", lead.get("nicho", ""), lead.get("cidade", "")],
                        connection_id=connection_id
                    )
                    if contact_id:
                        # Adiciona ao pipeline AGENTE SDR / FILA DE ATENDIMENTO
                        pip_ok, pip_err = unnichat_adicionar_pipeline(
                            contact_id, token, connection_id,
                            pipeline_id=pipeline_id or "8TgKxh0VK37bGNkbN8yO",
                            column_id=stage_id or "HThNTr24Yy2fN0vmeLm2",
                            business_name=lead["nome"]
                        )
                        lead["enviado_crm"] = True
                        lead["contact_id"] = contact_id
                        enviados += 1
                        job["leads_enviados"] = enviados
                        status_pip = "+ pipeline" if pip_ok else f"(pipeline: {pip_err})"
                        job["logs"].append(
                            f"[{_ts()}] Unnichat ({i+1}/{len(leads_coletados)}): {lead['nome']} enviado {status_pip}"
                        )
                    else:
                        job["logs"].append(
                            f"[{_ts()}] Unnichat: falhou para {lead['nome']} | {err_contato}"
                        )
                else:
                    # ── CRM: cria contato + oportunidade no pipeline ──
                    contact_id, err_contato = criar_contato_ghl(lead, token, location_id)
                    if contact_id:
                        ok, opportunity_id, err_oport = criar_oportunidade_ghl(lead, contact_id, token, location_id, pipeline_id, stage_id)
                        is_duplicate = err_oport and "duplicate" in err_oport.lower()
                        if ok or is_duplicate:
                            lead["enviado_crm"] = True
                            lead["contact_id"] = contact_id
                            enviados += 1
                            job["leads_enviados"] = enviados
                            status_msg = "já existe no pipeline" if is_duplicate else "enviado"
                            job["logs"].append(
                                f"[{_ts()}] CRM ({i+1}/{len(leads_coletados)}): {lead['nome']} {status_msg}"
                            )
                            # ── Popula campos enriquecidos no CRM ──
                            if enriquecimento_auto and any(lead.get(k) for k in ("email", "cnpj", "decisor_nome", "linkedin_url")):
                                try:
                                    _adicionar_nota_ghl(contact_id, token, location_id, lead, opportunity_id=opportunity_id)
                                    job["logs"].append(
                                        f"[{_ts()}] Campos enriquecidos salvos: {lead['nome']}"
                                    )
                                except Exception as e:
                                    job["logs"].append(f"[{_ts()}] Aviso: erro ao salvar campos de {lead['nome']}: {e}")
                        else:
                            job["logs"].append(
                                f"[{_ts()}] CRM: oportunidade falhou para {lead['nome']} | {err_oport}"
                            )
                    else:
                        job["logs"].append(
                            f"[{_ts()}] CRM: contato falhou para {lead['nome']} | {err_contato}"
                        )

                if len(job["logs"]) > 50:
                    job["logs"] = job["logs"][-50:]

                time.sleep(0.5)

        job["leads_enviados"] = enviados
        # Guarda TODOS os leads coletados (sem filtro_sem_site) — exibição na tabela
        job["leads"] = todos_leads_para_exibir
        job["status"] = "concluido"
        if enriquecimento_auto:
            job["logs"].append(
                f"[{_ts()}] Concluido: {enviados}/{len(leads_coletados)} leads enviados ao CRM ({len(todos_leads_para_exibir)} total coletados)"
            )
        else:
            job["logs"].append(
                f"[{_ts()}] Concluido: {len(todos_leads_para_exibir)} leads na tabela (modo Análise). Use 'Analisar Lead' pra estudar cada um."
            )
        job["concluido_em"] = datetime.now().isoformat()

        # Incrementa quota mensal da conta (total de leads coletados pela ferramenta)
        try:
            if conta_id and len(todos_leads_para_exibir) > 0:
                resetar_quota_se_novo_mes(conta_id)
                quota = incrementar_quota(conta_id, len(todos_leads_para_exibir))
                job["logs"].append(
                    f"[{_ts()}] Quota: {quota['usados']}/{quota['limite']} ({quota['disponivel']} disponiveis)"
                )
        except Exception as e:
            job["logs"].append(f"[{_ts()}] Aviso quota: {str(e)[:80]}")

        try:
            salvar_job_historico(job)
        except Exception as e:
            print(f"[jobs_historico] erro ao salvar conclusao {job.get('id')}: {e}")

    except Exception as e:
        job["status"] = "erro"
        job["erro"] = str(e)
        job["concluido_em"] = datetime.now().isoformat()
        job["logs"].append(f"[{_ts()}] ERRO: {str(e)}")
        try:
            salvar_job_historico(job)
        except Exception as ee:
            print(f"[jobs_historico] erro ao salvar erro {job.get('id')}: {ee}")


def _ts():
    return datetime.now().strftime("%H:%M:%S")


def _dono_atual():
    """Retorna (conta_id, user_id) do usuário logado, ou (None, None) fora de contexto."""
    try:
        if current_user and current_user.is_authenticated:
            return current_user.conta_id, int(current_user.id)
    except Exception:
        pass
    return None, None


def _novo_job(nicho, cidade, limite):
    job_id = str(uuid.uuid4())[:8]
    _cid, _uid = _dono_atual()
    jobs[job_id] = {
        "id": job_id,
        "tipo": "prospeccao",
        "nicho": nicho,
        "cidade": cidade,
        "limite": limite,
        "conta_id": _cid,
        "user_id": _uid,
        "status": "aguardando",
        "progresso": 0,
        "total": limite,
        "leads_enviados": 0,
        "leads": [],
        "logs": [],
        "erro": None,
        "criado_em": datetime.now().isoformat(),
        "concluido_em": None,
    }
    if len(jobs) > MAX_JOBS_HISTORY:
        ids_ordenados = sorted(jobs.keys(), key=lambda k: jobs[k]["criado_em"])
        for old_id in ids_ordenados[:-MAX_JOBS_HISTORY]:
            del jobs[old_id]
    try:
        salvar_job_historico(jobs[job_id])
    except Exception as e:
        print(f"[jobs_historico] erro ao salvar criação {job_id}: {e}")
    return job_id


ENRIQ_LEADS_DIR = Path("/opt/mia/workspace/prospeccao_ativa/enriq_leads")
ENRIQ_LEADS_DIR.mkdir(exist_ok=True)


def _persistir_leads_enriq(job_id: str, leads: list):
    (ENRIQ_LEADS_DIR / f"{job_id}.json").write_text(json.dumps(leads, ensure_ascii=False))


def _carregar_leads_enriq(job_id: str) -> list:
    f = ENRIQ_LEADS_DIR / f"{job_id}.json"
    if f.exists():
        return json.loads(f.read_text())
    return []


def _novo_job_enriq(empresas: list) -> str:
    job_id = str(uuid.uuid4())[:8]
    _cid, _uid = _dono_atual()
    jobs[job_id] = {
        "id": job_id,
        "tipo": "enriquecimento",
        "empresas": empresas,
        "conta_id": _cid,
        "user_id": _uid,
        "status": "aguardando",
        "progresso": 0,
        "total": len(empresas),
        "leads_enviados": 0,
        "logs": [],
        "leads": [],
        "erro": None,
        "criado_em": datetime.now().isoformat(),
        "concluido_em": None,
    }
    if len(jobs) > MAX_JOBS_HISTORY:
        ids_ordenados = sorted(jobs.keys(), key=lambda k: jobs[k]["criado_em"])
        for old_id in ids_ordenados[:-MAX_JOBS_HISTORY]:
            del jobs[old_id]
    try:
        salvar_job_historico(jobs[job_id])
    except Exception as e:
        print(f"[jobs_historico] erro ao salvar criação {job_id}: {e}")
    return job_id


# ── ENRIQUECIMENTO ────────────────────────────────────────────────────────────

def enriquecer_empresa(nome_empresa: str) -> dict:
    """
    Dado o nome da empresa:
    1. Busca CNPJ via DuckDuckGo (extrai de sites de consulta)
    2. Consulta ReceitaWS com o CNPJ -> telefone, email, endereço, sócios (Receita Federal)
    3. Busca LinkedIn do decisor via DuckDuckGo
    """
    from ddgs import DDGS
    from bs4 import BeautifulSoup
    from urllib.parse import urljoin, urlparse

    dados: dict = {
        "nome": nome_empresa,
        "cnpj": "",
        "telefone": "",
        "whatsapp": "",
        "email": "",
        "site": "",
        "instagram": "",
        "facebook": "",
        "decisor_nome": "",
        "decisor_cargo": "",
        "linkedin_url": "",
        "linkedin_empresa": "",
        "endereco": "",
        "fontes": [],
    }

    # --- 1. Busca CNPJ via DuckDuckGo (região Brasil) ---
    cnpj_raw = ""
    try:
        with DDGS() as ddgs:
            resultados = list(ddgs.text(f"{nome_empresa} CNPJ", max_results=5, region="br-pt"))
        for r in resultados:
            texto = r.get("body", "") + " " + r.get("href", "") + " " + r.get("title", "")
            matches = re.findall(r"\d{2}\.?\d{3}\.?\d{3}\/?\d{4}-?\d{2}", texto)
            for match in matches:
                limpo = re.sub(r"[^\d]", "", match)
                if len(limpo) == 14 and len(set(limpo)) > 3:
                    cnpj_raw = match
                    break
            if cnpj_raw:
                break
        time.sleep(1)
    except Exception:
        pass

    # --- 2. Consulta ReceitaWS (dados da Receita Federal) ---
    if cnpj_raw:
        cnpj_limpo = re.sub(r"[^\d]", "", cnpj_raw)
        if len(cnpj_limpo) == 14:
            try:
                r_api = requests.get(
                    f"https://receitaws.com.br/v1/cnpj/{cnpj_limpo}",
                    headers={"User-Agent": "Mozilla/5.0"},
                    timeout=12
                )
                if r_api.status_code == 200:
                    api_data = r_api.json()
                    if api_data.get("status") != "ERROR":
                        tel = api_data.get("telefone", "").strip()
                        if tel:
                            dados["telefone"] = tel
                            dados["whatsapp"] = tel
                        dados["email"] = (api_data.get("email", "") or "").strip()
                        # Endereço
                        partes = [p for p in [
                            api_data.get("logradouro", ""),
                            api_data.get("numero", ""),
                            api_data.get("bairro", ""),
                            api_data.get("municipio", ""),
                            api_data.get("uf", ""),
                        ] if p]
                        dados["endereco"] = ", ".join(partes)
                        # Formata CNPJ
                        c = cnpj_limpo
                        dados["cnpj"] = f"{c[:2]}.{c[2:5]}.{c[5:8]}/{c[8:12]}-{c[12:14]}"
                        # Sócio principal (decisor)
                        qsa = api_data.get("qsa", [])
                        if qsa:
                            dados["decisor_nome"] = qsa[0].get("nome", "").strip()
                            dados["decisor_cargo"] = qsa[0].get("qual", "").strip()
                        dados["fontes"].append("receita_federal")
                time.sleep(1)
            except Exception:
                pass

    # --- 3. Busca LinkedIn via DuckDuckGo ---
    try:
        with DDGS() as ddgs:
            li_results = list(ddgs.text(
                f'site:linkedin.com/in "{nome_empresa}"',
                max_results=3,
                region="br-pt"
            ))
        for r in li_results:
            url = r.get("href", "")
            if "linkedin.com/in/" in url:
                dados["linkedin_url"] = url.split("?")[0]
                dados["fontes"].append("linkedin")
                break
        time.sleep(1)
    except Exception:
        pass

    # --- 3.5. Busca LinkedIn da empresa (company page) ---
    if not dados.get("linkedin_empresa"):
        try:
            with DDGS() as ddgs:
                li_company = list(ddgs.text(
                    f'site:linkedin.com/company "{nome_empresa}"',
                    max_results=3,
                    region="br-pt"
                ))
            for r in li_company:
                url = r.get("href", "")
                if "linkedin.com/company/" in url:
                    dados["linkedin_empresa"] = url.split("?")[0]
                    dados["fontes"].append("linkedin_company")
                    break
            time.sleep(1)
        except Exception:
            pass

    # --- 4. Busca site oficial via DuckDuckGo ---
    if not dados["site"]:
        try:
            with DDGS() as ddgs:
                site_results = list(ddgs.text(
                    f'"{nome_empresa}" site oficial',
                    max_results=5,
                    region="br-pt"
                ))
            for r in site_results:
                url = r.get("href", "")
                # Ignora agregadores, busca direta no domínio da empresa
                ignorar = ["facebook.com", "instagram.com", "linkedin.com", "cnpj.biz",
                           "receita.", "infocnpj", "cnpja", "jusbrasil", "google.com",
                           "youtube.com", "wikipedia", "receitaws"]
                if url and not any(x in url for x in ignorar):
                    dados["site"] = url.split("?")[0]
                    dados["fontes"].append("site_oficial")
                    break
            time.sleep(1)
        except Exception:
            pass

    # --- 4.5. Scraping profundo do site oficial (redes sociais, emails, telefones) ---
    if dados["site"]:
        headers_req = {
            "User-Agent": "Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 Chrome/120 Safari/537.36",
            "Accept-Language": "pt-BR,pt;q=0.9,en;q=0.8",
        }
        base_site = dados["site"].rstrip("/")
        # Normaliza base pra origem (protocolo+host), caso o "site" tenha vindo com path
        try:
            _parsed = urlparse(base_site)
            if _parsed.scheme and _parsed.netloc:
                origem = f"{_parsed.scheme}://{_parsed.netloc}"
            else:
                origem = base_site
        except Exception:
            origem = base_site

        paginas_tentar = [
            base_site,
            f"{origem}/contato",
            f"{origem}/sobre",
            f"{origem}/quem-somos",
            f"{origem}/equipe",
            f"{origem}/fale-conosco",
        ]
        # Remove duplicatas mantendo ordem
        seen_urls = set()
        paginas_unicas = []
        for u in paginas_tentar:
            if u and u not in seen_urls:
                seen_urls.add(u)
                paginas_unicas.append(u)

        achou_algo_site = False
        for pag_url in paginas_unicas:
            try:
                resp = requests.get(pag_url, headers=headers_req, timeout=8, allow_redirects=True)
                if resp.status_code != 200:
                    continue
                html = resp.text
                soup = BeautifulSoup(html, "html.parser")

                # Extrai links de redes sociais
                for a_tag in soup.find_all("a", href=True):
                    href = a_tag["href"]
                    # Instagram
                    if "instagram.com/" in href and not dados.get("instagram"):
                        m = re.search(r'instagram\.com/([^/?#"\']+)', href)
                        if m and m.group(1) not in ("p", "reel", "explore", "stories", "tv", "accounts", "sharedfiles"):
                            dados["instagram"] = f"https://instagram.com/{m.group(1).rstrip('/')}"
                            achou_algo_site = True
                    # Facebook
                    if "facebook.com/" in href and not dados.get("facebook"):
                        m = re.search(r'facebook\.com/([^/?#"\']+)', href)
                        if m and m.group(1) not in ("sharer", "share", "login", "groups", "pages", "events", "watch", "dialog", "tr"):
                            dados["facebook"] = f"https://facebook.com/{m.group(1).rstrip('/')}"
                            achou_algo_site = True
                    # LinkedIn pessoa
                    if "linkedin.com/in/" in href and not dados.get("linkedin_url"):
                        dados["linkedin_url"] = href.split("?")[0]
                        achou_algo_site = True
                    # LinkedIn empresa
                    if "linkedin.com/company/" in href and not dados.get("linkedin_empresa"):
                        dados["linkedin_empresa"] = href.split("?")[0]
                        achou_algo_site = True

                # Emails no HTML todo (inclui mailto)
                if not dados.get("email"):
                    emails_pag = re.findall(r"[a-zA-Z0-9._%+\-]+@[a-zA-Z0-9.\-]+\.[a-zA-Z]{2,}", html)
                    for em in emails_pag:
                        em_low = em.lower()
                        if not any(x in em_low for x in ["example", "test", "noreply", "no-reply", "sentry", "wixpress", "sentry.io", "@2x", "@3x"]):
                            dados["email"] = em
                            achou_algo_site = True
                            break

                # Telefones BR
                if not dados.get("telefone"):
                    tels_pag = re.findall(r"(?:\+55[-.\s]?)?(?:\(?\d{2}\)?)[-.\s]?\d{4,5}[-.\s]?\d{4}", html)
                    for tel in tels_pag:
                        # Filtra ruído: precisa ter pelo menos 10 dígitos
                        digitos = re.sub(r"[^\d]", "", tel)
                        if 10 <= len(digitos) <= 13:
                            dados["telefone"] = tel.strip()
                            if not dados.get("whatsapp"):
                                dados["whatsapp"] = tel.strip()
                            achou_algo_site = True
                            break

            except Exception:
                continue

        if achou_algo_site and "website_scraping" not in dados["fontes"]:
            dados["fontes"].append("website_scraping")

    # --- 5. Busca Instagram via DuckDuckGo (só se scraping não achou) ---
    if not dados["instagram"]:
        try:
            with DDGS() as ddgs:
                ig_results = list(ddgs.text(
                    f'site:instagram.com "{nome_empresa}"',
                    max_results=3,
                    region="br-pt"
                ))
            for r in ig_results:
                url = r.get("href", "")
                if "instagram.com/" in url:
                    # Extrai handle: instagram.com/handle/
                    m = re.search(r"instagram\.com/([^/?#]+)", url)
                    if m:
                        handle = m.group(1)
                        if handle not in ("p", "reel", "explore", "stories", "tv"):
                            dados["instagram"] = f"https://instagram.com/{handle}"
                            dados["fontes"].append("instagram")
                            break
            time.sleep(1)
        except Exception:
            pass

    # --- 6. Busca Facebook via DuckDuckGo (só se scraping não achou) ---
    if not dados["facebook"]:
        try:
            with DDGS() as ddgs:
                fb_results = list(ddgs.text(
                    f'site:facebook.com "{nome_empresa}"',
                    max_results=3,
                    region="br-pt"
                ))
            for r in fb_results:
                url = r.get("href", "")
                if "facebook.com/" in url:
                    m = re.search(r"facebook\.com/([^/?#]+)", url)
                    if m:
                        handle = m.group(1)
                        if handle not in ("pages", "groups", "events", "login", "sharer", "share", "watch"):
                            dados["facebook"] = f"https://facebook.com/{handle}"
                            dados["fontes"].append("facebook")
                            break
            time.sleep(1)
        except Exception:
            pass

    # --- 7. Busca email do decisor via DuckDuckGo ---
    if dados["decisor_nome"] and not dados["email"]:
        try:
            with DDGS() as ddgs:
                em_results = list(ddgs.text(
                    f'"{dados["decisor_nome"]}" "{nome_empresa}" email contato',
                    max_results=5,
                    region="br-pt"
                ))
            for r in em_results:
                texto = r.get("body", "") + " " + r.get("title", "")
                emails = re.findall(r"[a-zA-Z0-9._%+\-]+@[a-zA-Z0-9.\-]+\.[a-zA-Z]{2,}", texto)
                for em in emails:
                    if not any(x in em.lower() for x in ["example", "test", "noreply", "no-reply", "sentry"]):
                        dados["email"] = em
                        dados["fontes"].append("email_decisor")
                        break
                if dados["email"]:
                    break
            time.sleep(1)
        except Exception:
            pass

    return dados


def enriquecer_empresa_us(nome_empresa: str, website: str = "", cidade: str = "") -> dict:
    """
    Enriquecimento para empresas nos EUA (sem CNPJ):
    1. Scraping do site da empresa: emails, telefones, nome do responsável
    2. DuckDuckGo: decisor (owner/CEO/founder) + LinkedIn
    """
    from ddgs import DDGS
    from bs4 import BeautifulSoup

    dados: dict = {
        "nome": nome_empresa,
        "cnpj": "",
        "telefone": "",
        "email": "",
        "site": website or "",
        "decisor_nome": "",
        "decisor_cargo": "",
        "linkedin_url": "",
        "endereco": "",
        "fontes": [],
    }

    headers_req = {
        "User-Agent": "Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 Chrome/120 Safari/537.36",
        "Accept-Language": "en-US,en;q=0.9",
    }

    # --- 1. Scraping do site ---
    site_url = website.strip() if website else ""
    if site_url and not site_url.startswith("http"):
        site_url = "https://" + site_url

    if site_url:
        pages_to_try = [site_url]
        for suffix in ["/contact", "/contact-us", "/about", "/about-us", "/team"]:
            pages_to_try.append(site_url.rstrip("/") + suffix)

        found_email = ""
        found_phone = ""
        found_person = ""

        for page_url in pages_to_try[:4]:
            try:
                resp = requests.get(page_url, headers=headers_req, timeout=8, allow_redirects=True)
                if resp.status_code != 200:
                    continue
                soup = BeautifulSoup(resp.text, "html.parser")
                text = soup.get_text(" ", strip=True)

                # Extrai email
                if not found_email:
                    emails = re.findall(r"[a-zA-Z0-9._%+\-]+@[a-zA-Z0-9.\-]+\.[a-zA-Z]{2,}", text)
                    for em in emails:
                        if not any(x in em.lower() for x in ["example", "youremail", "test", "noreply", "info@", "hello@"]):
                            found_email = em
                            break
                    if not found_email and emails:
                        found_email = emails[0]

                # Extrai telefone
                if not found_phone:
                    phones = re.findall(r"(?:\+1[-.\s]?)?\(?\d{3}\)?[-.\s]\d{3}[-.\s]\d{4}", text)
                    if phones:
                        found_phone = phones[0]

                # Extrai nome de pessoa (Owner/CEO/Founder/President)
                if not found_person:
                    match = re.search(
                        r"([A-Z][a-z]+ [A-Z][a-z]+(?:\s[A-Z][a-z]+)?)\s*[,\-–|]?\s*"
                        r"(Owner|Co-Owner|CEO|Founder|Co-Founder|President|Director|Manager|Principal)",
                        text
                    )
                    if match:
                        found_person = match.group(1).strip()
                        dados["decisor_cargo"] = match.group(2).strip()

                if found_email and found_phone:
                    break
                time.sleep(0.5)

            except Exception:
                continue

        if found_email:
            dados["email"] = found_email
            dados["fontes"].append("website_scraping")
        if found_phone and not dados.get("telefone"):
            dados["telefone"] = found_phone
            if "website_scraping" not in dados["fontes"]:
                dados["fontes"].append("website_scraping")
        if found_person:
            dados["decisor_nome"] = found_person

    # --- 2. DuckDuckGo: decisor ---
    if not dados["decisor_nome"]:
        try:
            query = f'"{nome_empresa}" {cidade} owner OR CEO OR founder'.strip()
            with DDGS() as ddgs:
                results = list(ddgs.text(query, max_results=5, region="us-en"))
            for r in results:
                snippet = r.get("body", "") + " " + r.get("title", "")
                m = re.search(
                    r"([A-Z][a-z]+ [A-Z][a-z]+)\s*[,\-–]?\s*(owner|CEO|founder|president|director|manager)",
                    snippet, re.IGNORECASE
                )
                if m:
                    dados["decisor_nome"] = m.group(1).strip()
                    dados["decisor_cargo"] = m.group(2).strip().title()
                    dados["fontes"].append("duckduckgo")
                    break
            time.sleep(1)
        except Exception:
            pass

    # --- 3. DuckDuckGo: LinkedIn ---
    try:
        li_query = f'site:linkedin.com/in "{nome_empresa}" owner OR CEO OR founder'
        with DDGS() as ddgs:
            li_results = list(ddgs.text(li_query, max_results=3, region="us-en"))
        for r in li_results:
            url = r.get("href", "")
            if "linkedin.com/in/" in url:
                dados["linkedin_url"] = url.split("?")[0]
                if "linkedin" not in dados["fontes"]:
                    dados["fontes"].append("linkedin")
                # Tenta extrair nome do título do resultado
                if not dados["decisor_nome"]:
                    title = r.get("title", "")
                    m = re.match(r"^([A-Z][a-z]+ [A-Z][a-z]+)", title)
                    if m:
                        dados["decisor_nome"] = m.group(1)
                break
        time.sleep(1)
    except Exception:
        pass

    # --- 4. DuckDuckGo: email (fallback) ---
    if not dados["email"]:
        try:
            eq = f'"{nome_empresa}" {cidade} email contact'.strip()
            with DDGS() as ddgs:
                eq_results = list(ddgs.text(eq, max_results=5, region="us-en"))
            for r in eq_results:
                text = r.get("body", "") + " " + r.get("title", "")
                emails = re.findall(r"[a-zA-Z0-9._%+\-]+@[a-zA-Z0-9.\-]+\.[a-zA-Z]{2,}", text)
                for em in emails:
                    if not any(x in em.lower() for x in ["example", "youremail", "test"]):
                        dados["email"] = em
                        dados["fontes"].append("duckduckgo_email")
                        break
                if dados["email"]:
                    break
            time.sleep(1)
        except Exception:
            pass

    return dados


def _adicionar_nota_ghl(contact_id: str, token: str, location_id: str, dados: dict, opportunity_id: str = None):
    """Popula campos personalizados na oportunidade + adiciona nota."""

    # 1. Popula campos personalizados na OPORTUNIDADE (aparece no card do pipeline)
    campos_ids = _carregar_campos_ghl()
    custom_fields = []
    mapa = {
        "CNPJ":             dados.get("cnpj", ""),
        "Email":            dados.get("email", ""),
        "LinkedIn":         dados.get("linkedin_url", ""),
        "Decisor":          dados.get("decisor_nome", ""),
        "Cargo do Decisor": dados.get("decisor_cargo", ""),
        "Site":             dados.get("site", ""),
        "Instagram":        dados.get("instagram", ""),
        "Facebook":         dados.get("facebook", ""),
    }
    for nome_campo, valor in mapa.items():
        field_id = campos_ids.get(nome_campo)
        if field_id and valor:
            custom_fields.append({"id": field_id, "field_value": valor})

    # Salva na OPORTUNIDADE (campos são opportunity-level no GHL)
    if custom_fields and opportunity_id:
        try:
            requests.put(
                f"{GHL_BASE}/opportunities/{opportunity_id}",
                headers=_ghl_headers(token),
                json={"customFields": custom_fields},
                timeout=10
            )
        except Exception:
            pass

    # 2. Adiciona nota com resumo completo
    linhas = []
    if dados.get("cnpj"):         linhas.append(f"CNPJ: {dados['cnpj']}")
    if dados.get("email"):        linhas.append(f"Email: {dados['email']}")
    if dados.get("decisor_nome"): linhas.append(f"Decisor: {dados['decisor_nome']} ({dados.get('decisor_cargo', '')})")
    if dados.get("linkedin_url"): linhas.append(f"LinkedIn: {dados['linkedin_url']}")
    if dados.get("site"):         linhas.append(f"Site: {dados['site']}")
    if dados.get("instagram"):    linhas.append(f"Instagram: {dados['instagram']}")
    if dados.get("facebook"):     linhas.append(f"Facebook: {dados['facebook']}")
    if dados.get("endereco"):     linhas.append(f"Endereço: {dados['endereco']}")

    if linhas:
        try:
            requests.post(
                f"{GHL_BASE}/contacts/{contact_id}/notes",
                headers=_ghl_headers(token),
                json={"body": "\n".join(linhas), "userId": location_id},
                timeout=10
            )
        except Exception:
            pass


def run_enriquecedor(job_id: str, empresas: list, token: str,
                     location_id: str, pipeline_id: str, stage_id: str,
                     enviar_crm: bool = True):
    """Enriquece lista de empresas e (opcionalmente) cria contatos no CRM."""
    job = jobs[job_id]
    job["status"] = "em_andamento"
    job["total"] = len(empresas)
    # conta_id já é atribuído ao job antes desta thread iniciar (via /enriquecer,
    # /enriquecer-preview, etc). Fase 1: sem conta_id, o registro NÃO persiste
    # (fail-safe pra evitar vazamento entre tenants).
    conta_id = job.get("conta_id")

    try:
        for i, nome_empresa in enumerate(empresas):
            if job["status"] == "cancelado":
                break

            job["logs"].append(f"[{_ts()}] Enriquecendo ({i+1}/{len(empresas)}): {nome_empresa}")
            if len(job["logs"]) > 100:
                job["logs"] = job["logs"][-100:]

            dados = enriquecer_empresa(nome_empresa)

            # Monta lead no formato que criar_contato_ghl espera
            lead = {
                "nome": dados.get("decisor_nome") or nome_empresa,
                "telefone": dados.get("telefone", ""),
                "email": dados.get("email", ""),
                "empresa": nome_empresa,
                "endereco": dados.get("endereco", ""),
                "site": dados.get("site", ""),
                # Campos que alimentam custom fields no CRM (Sweet Angels etc.)
                "cnpj": dados.get("cnpj", ""),
                "decisor_nome": dados.get("decisor_nome", ""),
                "decisor_cargo": dados.get("decisor_cargo", ""),
                "linkedin_url": dados.get("linkedin_url", ""),
                "instagram": dados.get("instagram", ""),
                "facebook": dados.get("facebook", ""),
                "fontes": dados.get("fontes") or [],
            }

            tags = ["lista-enriquecida", "enriquecedor"]
            if dados.get("linkedin_url"):
                tags.append("linkedin-encontrado")
            if dados.get("cnpj"):
                tags.append("cnpj-encontrado")
            lead["tags"] = tags

            if enviar_crm and token and location_id:
                contact_id, err = criar_contato_ghl(
                    lead, token, location_id, source="Enriquecimento",
                    custom_fields_map=_carregar_campos_crm(location_id) or None,
                )
                if contact_id:
                    opportunity_id = None
                    if pipeline_id and stage_id:
                        ok, opportunity_id, _ = criar_oportunidade_ghl(
                            lead, contact_id, token, location_id, pipeline_id, stage_id, source="Enriquecimento"
                        )
                    _adicionar_nota_ghl(contact_id, token, location_id, dados, opportunity_id=opportunity_id)
                    job["leads_enviados"] += 1
                    job["logs"].append(
                        f"[{_ts()}] ✓ {nome_empresa} → CRM "
                        f"(fone: {dados.get('telefone') or '—'} | "
                        f"LinkedIn: {'sim' if dados.get('linkedin_url') else 'não'})"
                    )
                else:
                    job["logs"].append(f"[{_ts()}] ✗ {nome_empresa} → erro CRM: {err}")
            else:
                job["logs"].append(
                    f"[{_ts()}] ✓ {nome_empresa} → enriquecido "
                    f"(fone: {dados.get('telefone') or '—'} | "
                    f"LinkedIn: {'sim' if dados.get('linkedin_url') else 'não'}) — aguardando envio manual"
                )

            job["progresso"] = i + 1
            job["leads"].append(dados)
            time.sleep(1)

        job["status"] = "concluido"
        job["concluido_em"] = datetime.now().isoformat()
        if enviar_crm:
            job["logs"].append(
                f"[{_ts()}] Concluído: {job['leads_enviados']}/{len(empresas)} enviados ao CRM"
            )
        else:
            job["logs"].append(
                f"[{_ts()}] Concluído: {len(job['leads'])}/{len(empresas)} enriquecidos — pronto para envio ao CRM"
            )
        try:
            salvar_job_historico(job)
        except Exception as e:
            print(f"[jobs_historico] erro ao salvar conclusao {job.get('id')}: {e}")
        try:
            _persistir_leads_enriq(job_id, job.get("leads", []))
        except Exception as e:
            print(f"[enriq] erro ao persistir leads {job_id}: {e}")

    except Exception as e:
        job["status"] = "erro"
        job["erro"] = str(e)
        job["concluido_em"] = datetime.now().isoformat()
        job["logs"].append(f"[{_ts()}] ERRO: {str(e)}")
        try:
            salvar_job_historico(job)
        except Exception as ee:
            print(f"[jobs_historico] erro ao salvar erro {job.get('id')}: {ee}")


# ---------- ROUTES ----------

@app.route("/login", methods=["GET", "POST"])
def login():
    if current_user.is_authenticated:
        return redirect(url_for("index"))
    if request.method == "POST":
        ip = _login_ip()
        if _login_bloqueado(ip):
            flash("Muitas tentativas de login. Aguarde alguns minutos e tente de novo.", "error")
            return render_template("login.html"), 429
        email = (request.form.get("email") or "").strip().lower()
        senha = (request.form.get("senha") or "").strip()
        user = buscar_por_email(email)
        if user and verificar_senha(user, senha):
            if not user.ativo:
                flash("Conta inativa. Entre em contato com o administrador.", "error")
                return render_template("login.html")
            _login_limpar(ip)
            login_user(user, remember=True)
            registrar_acesso(int(user.id))
            return redirect(url_for("index"))
        _login_registrar_falha(ip)
        flash("E-mail ou senha incorretos.", "error")
    return render_template("login.html")


@app.route("/logout")
@login_required
def logout():
    logout_user()
    return redirect(url_for("login"))


@app.route("/admin")
@login_required
def admin():
    if not current_user.is_super_admin:
        return redirect(url_for("index"))
    usuarios = listar_usuarios()
    contas   = listar_contas()
    # Monta mapa conta_id -> ghl_config para exibir status no painel
    ghl_status = {}
    for c in contas:
        cfg = get_ghl_config(c["id"])
        ghl_status[c["id"]] = cfg["location_name"] if cfg else None
    # Dashboard de saúde do consumo de leads / free tier Google Places
    try:
        dashboard_saude = get_dashboard_saude()
    except Exception as _e:
        app.logger.exception("Falha ao montar dashboard_saude: %s", _e)
        dashboard_saude = None
    return render_template(
        "admin.html",
        usuarios=usuarios,
        contas=contas,
        ghl_status=ghl_status,
        dashboard_saude=dashboard_saude,
    )


@app.route("/admin/criar-usuario", methods=["POST"])
@login_required
def admin_criar_usuario():
    if not current_user.is_super_admin:
        return redirect(url_for("index"))
    nome     = (request.form.get("nome") or "").strip()
    email    = (request.form.get("email") or "").strip()
    senha    = (request.form.get("senha") or "").strip()
    plano    = (request.form.get("plano") or "beta").strip()
    conta_id = request.form.get("conta_id") or None
    if conta_id:
        conta_id = int(conta_id)
    user, erro = criar_usuario(nome, email, senha, plano, conta_id)
    if erro:
        flash(f"Erro ao criar usuário: {erro}", "error")
    else:
        flash(f"Usuário {nome} criado.", "success")
    return redirect(url_for("admin"))


@app.route("/admin/usuario/<int:uid>/conta", methods=["POST"])
@login_required
def admin_atribuir_conta(uid):
    if not current_user.is_super_admin:
        return redirect(url_for("index"))
    conta_id = request.form.get("conta_id") or None
    plano    = request.form.get("plano") or None
    updates  = {}
    if conta_id is not None:
        updates["conta_id"] = int(conta_id) if conta_id else None
    if plano:
        updates["plano"] = plano
    if updates:
        atualizar_usuario(uid, **updates)
    return redirect(url_for("admin"))


@app.route("/admin/toggle-usuario/<int:uid>", methods=["POST"])
@login_required
def admin_toggle_usuario(uid):
    if not current_user.is_super_admin:
        return redirect(url_for("index"))
    user = User.get(uid)
    if user:
        novo_status = 0 if user.ativo else 1
        atualizar_usuario(uid, ativo=novo_status)
    return redirect(url_for("admin"))


@app.route("/admin/deletar-usuario/<int:uid>", methods=["POST"])
@login_required
def admin_deletar_usuario(uid):
    if not current_user.is_super_admin:
        return redirect(url_for("index"))
    if str(uid) != current_user.id:
        deletar_usuario(uid)
    return redirect(url_for("admin"))


@app.route("/")
@login_required
def index():
    return render_template("index.html",
                           is_admin=current_user.is_super_admin,
                           is_conta_admin=current_user.is_conta_admin,
                           pode_analisar=pode_analisar_leads(current_user))


@app.route("/api/conectar", methods=["POST"])
@login_required
def api_conectar():
    """Valida credenciais do CRM e retorna configuração."""
    data = request.get_json()
    crm_type = (data.get("crm_type") or "ghl").strip().lower()
    token = (data.get("token") or "").strip()
    # Remove prefixo "Bearer " se o usuário colou com ele
    if token.lower().startswith("bearer "):
        token = token[7:].strip()

    if not token:
        return jsonify({"ok": False, "erro": "Token obrigatório"}), 400

    # ── UNNICHAT ──
    if crm_type == "unnichat":
        connection_id = (data.get("connection_id") or "").strip()
        try:
            headers = {"Authorization": f"Bearer {token}"}
            if connection_id:
                headers["x-connection-id"] = connection_id
            r = requests.get(f"{UNNICHAT_BASE}/api/tags", headers=headers, timeout=10)
            if r.status_code == 200:
                return jsonify({"ok": True, "crm_type": "unnichat", "location_name": "Unnichat", "connection_id": connection_id})
            return jsonify({"ok": False, "erro": f"Token inválido (HTTP {r.status_code})"}), 200
        except requests.exceptions.Timeout:
            return jsonify({"ok": False, "erro": "Timeout ao conectar com o Unnichat"}), 200
        except Exception as e:
            return jsonify({"ok": False, "erro": str(e)}), 200

    # ── CRM (padrão) ──
    location_id = (data.get("location_id") or "").strip()
    if not location_id:
        return jsonify({"ok": False, "erro": "token e location_id sao obrigatorios"}), 400

    try:
        r = requests.get(
            f"{GHL_BASE}/opportunities/pipelines",
            headers=_ghl_headers(token),
            params={"locationId": location_id},
            timeout=10
        )
        if r.status_code != 200:
            return jsonify({"ok": False, "erro": f"Credenciais invalidas (HTTP {r.status_code})"}), 200

        pipelines_raw = r.json().get("pipelines", [])
        location_name = location_id
        try:
            r_loc = requests.get(
                f"{GHL_BASE}/locations/{location_id}",
                headers=_ghl_headers(token),
                timeout=8
            )
            if r_loc.status_code == 200:
                loc_data = r_loc.json()
                location_name = (
                    loc_data.get("location", {}).get("name")
                    or loc_data.get("name")
                    or location_id
                )
        except Exception:
            pass

        pipelines = []
        for p in pipelines_raw:
            stages = [{"id": s.get("id"), "name": s.get("name")} for s in p.get("stages", [])]
            pipelines.append({"id": p.get("id"), "name": p.get("name"), "stages": stages})

        return jsonify({"ok": True, "crm_type": "ghl", "location_name": location_name, "pipelines": pipelines})

    except requests.exceptions.Timeout:
        return jsonify({"ok": False, "erro": "Timeout ao conectar com o CRM"}), 200
    except Exception as e:
        return jsonify({"ok": False, "erro": str(e)}), 200


@app.route("/api/ghl-config", methods=["GET"])
@login_required
def api_get_ghl_config():
    """Retorna config GHL da conta do usuário para auto-conectar."""
    conta_id = current_user.conta_id
    # super_admin sem conta: aceita ?conta_id= na query
    if current_user.is_super_admin and not conta_id:
        conta_id = request.args.get("conta_id", type=int)
    if not conta_id:
        return jsonify({"ok": False, "configurado": False})
    cfg = get_ghl_config(conta_id)
    if not cfg or not cfg.get("token"):
        return jsonify({"ok": False, "configurado": False, "conta_id": conta_id})
    return jsonify({
        "ok": True, "configurado": True,
        "crm_type": cfg["crm_type"], "token": cfg["token"],
        "location_id": cfg["location_id"], "pipeline_id": cfg["pipeline_id"],
        "stage_id": cfg["stage_id"], "location_name": cfg["location_name"],
        "pipeline_name": cfg["pipeline_name"], "stage_name": cfg["stage_name"],
    })


@app.route("/api/ghl-pipelines", methods=["GET"])
@login_required
def api_ghl_pipelines():
    """Lista pipelines e stages GHL da conta do usuário logado.
    Usado pelos dropdowns inline da tela principal para trocar pipeline/stage
    sem precisar entrar na tela de configuração.
    """
    conta_id = current_user.conta_id
    if current_user.is_super_admin and not conta_id:
        conta_id = request.args.get("conta_id", type=int)
    if not conta_id:
        return jsonify({"ok": False, "erro": "Conta não definida."}), 400

    cfg = get_ghl_config(conta_id)
    if not cfg or not cfg.get("token"):
        return jsonify({"ok": False, "erro": "CRM não configurado."}), 200

    # Só suporta pipelines GHL — Unnichat não tem esse conceito
    if (cfg.get("crm_type") or "ghl") != "ghl":
        return jsonify({"ok": False, "erro": "Pipelines só disponíveis para GHL."}), 200

    token = cfg.get("token") or ""
    location_id = cfg.get("location_id") or ""
    if not (token and location_id):
        return jsonify({"ok": False, "erro": "Token ou location_id ausentes na config."}), 200

    try:
        r = requests.get(
            f"{GHL_BASE}/opportunities/pipelines",
            headers=_ghl_headers(token),
            params={"locationId": location_id},
            timeout=10,
        )
        if r.status_code != 200:
            return jsonify({"ok": False, "erro": f"GHL retornou HTTP {r.status_code}"}), 200

        pipelines_raw = r.json().get("pipelines", [])
        pipelines = []
        for p in pipelines_raw:
            stages = [
                {"id": s.get("id"), "name": s.get("name")}
                for s in p.get("stages", [])
            ]
            pipelines.append({
                "id": p.get("id"),
                "name": p.get("name"),
                "stages": stages,
            })

        return jsonify({
            "ok": True,
            "pipelines": pipelines,
            "current_pipeline_id": cfg.get("pipeline_id") or "",
            "current_stage_id": cfg.get("stage_id") or "",
        })

    except requests.exceptions.Timeout:
        return jsonify({"ok": False, "erro": "Timeout ao consultar pipelines na GHL."}), 200
    except Exception as e:
        return jsonify({"ok": False, "erro": str(e)}), 200


@app.route("/api/ghl-config", methods=["POST"])
@login_required
def api_save_ghl_config():
    """Salva config GHL da conta (conta_admin ou super_admin)."""
    if not current_user.is_conta_admin:
        return jsonify({"ok": False, "erro": "Permissão negada."}), 403
    data = request.get_json()
    token = (data.get("token") or "").strip()
    if not token:
        return jsonify({"ok": False, "erro": "Token obrigatório."}), 400
    # super_admin pode salvar para qualquer conta via campo conta_id
    conta_id = current_user.conta_id
    if current_user.is_super_admin:
        conta_id = data.get("conta_id") or conta_id
    if not conta_id:
        return jsonify({"ok": False, "erro": "Conta não definida."}), 400
    save_ghl_config(
        conta_id=int(conta_id),
        crm_type=data.get("crm_type", "ghl"),
        token=token,
        location_id=data.get("location_id", ""),
        pipeline_id=data.get("pipeline_id", ""),
        stage_id=data.get("stage_id", ""),
        location_name=data.get("location_name", ""),
        pipeline_name=data.get("pipeline_name", ""),
        stage_name=data.get("stage_name", ""),
    )
    return jsonify({"ok": True})


# ── PERFIS DE CRM (múltiplos CRMs por conta) ─────────────

def _resolver_conta_id_admin():
    """Resolve o conta_id do request: usa o do usuário, ou aceita ?conta_id= se for super_admin."""
    conta_id = current_user.conta_id
    if current_user.is_super_admin:
        via_query = request.args.get("conta_id", type=int)
        via_json = None
        if request.is_json:
            body = request.get_json(silent=True) or {}
            via_json = body.get("conta_id")
        conta_id = via_query or via_json or conta_id
    return int(conta_id) if conta_id else None


@app.route("/api/crm-perfis", methods=["GET"])
@login_required
def api_listar_crm_perfis():
    """Lista os perfis de CRM da conta do usuário logado.
    Super_admin pode passar ?conta_id= para consultar de outra conta."""
    conta_id = _resolver_conta_id_admin()
    if not conta_id:
        return jsonify({"ok": True, "perfis": []})
    perfis = listar_perfis_crm(conta_id)
    # Não expõe o token completo na lista — só para admins da conta
    if not current_user.is_conta_admin:
        for p in perfis:
            if p.get("token"):
                p["token_preview"] = p["token"][:8] + "..." if len(p["token"]) > 12 else "***"
            # Ainda envia token pro frontend, já que precisamos usar na prospecção;
            # é o mesmo comportamento do /api/ghl-config atual.
    return jsonify({"ok": True, "perfis": perfis, "conta_id": conta_id})


@app.route("/api/crm-perfis", methods=["POST"])
@login_required
def api_criar_crm_perfil():
    """Cria um novo perfil de CRM para a conta."""
    if not current_user.is_conta_admin:
        return jsonify({"ok": False, "erro": "Permissão negada."}), 403
    data = request.get_json() or {}
    nome = (data.get("nome") or "").strip()
    token = (data.get("token") or "").strip()
    if not nome:
        return jsonify({"ok": False, "erro": "Nome do perfil é obrigatório."}), 400
    if not token:
        return jsonify({"ok": False, "erro": "Token é obrigatório."}), 400

    conta_id = current_user.conta_id
    if current_user.is_super_admin:
        conta_id = data.get("conta_id") or conta_id
    if not conta_id:
        return jsonify({"ok": False, "erro": "Conta não definida."}), 400

    perfil = salvar_perfil_crm(
        conta_id=int(conta_id),
        nome=nome,
        crm_type=data.get("crm_type", "ghl"),
        token=token,
        location_id=data.get("location_id", ""),
        pipeline_id=data.get("pipeline_id", ""),
        stage_id=data.get("stage_id", ""),
        location_name=data.get("location_name", ""),
        pipeline_name=data.get("pipeline_name", ""),
        stage_name=data.get("stage_name", ""),
        connection_id=data.get("connection_id", ""),
    )
    return jsonify({"ok": True, "perfil": perfil})


@app.route("/api/crm-perfis/<int:perfil_id>", methods=["PUT"])
@login_required
def api_atualizar_crm_perfil(perfil_id):
    """Atualiza um perfil existente."""
    if not current_user.is_conta_admin:
        return jsonify({"ok": False, "erro": "Permissão negada."}), 403
    data = request.get_json() or {}
    perfil = get_perfil_crm(perfil_id)
    if not perfil:
        return jsonify({"ok": False, "erro": "Perfil não encontrado."}), 404
    # conta_admin só pode editar perfil da própria conta
    if not current_user.is_super_admin and perfil["conta_id"] != current_user.conta_id:
        return jsonify({"ok": False, "erro": "Permissão negada."}), 403

    nome = (data.get("nome") or perfil["nome"]).strip()
    token = (data.get("token") or perfil["token"]).strip()
    if not nome or not token:
        return jsonify({"ok": False, "erro": "Nome e token são obrigatórios."}), 400

    atualizado = salvar_perfil_crm(
        conta_id=perfil["conta_id"],
        nome=nome,
        crm_type=data.get("crm_type", perfil["crm_type"]),
        token=token,
        location_id=data.get("location_id", perfil["location_id"]),
        pipeline_id=data.get("pipeline_id", perfil["pipeline_id"]),
        stage_id=data.get("stage_id", perfil["stage_id"]),
        location_name=data.get("location_name", perfil["location_name"]),
        pipeline_name=data.get("pipeline_name", perfil["pipeline_name"]),
        stage_name=data.get("stage_name", perfil["stage_name"]),
        connection_id=data.get("connection_id", perfil.get("connection_id", "")),
        perfil_id=perfil_id,
    )
    return jsonify({"ok": True, "perfil": atualizado})


@app.route("/api/crm-perfis/<int:perfil_id>", methods=["DELETE"])
@login_required
def api_deletar_crm_perfil(perfil_id):
    """Deleta um perfil, desde que não seja o único da conta."""
    if not current_user.is_conta_admin:
        return jsonify({"ok": False, "erro": "Permissão negada."}), 403
    perfil = get_perfil_crm(perfil_id)
    if not perfil:
        return jsonify({"ok": False, "erro": "Perfil não encontrado."}), 404
    if not current_user.is_super_admin and perfil["conta_id"] != current_user.conta_id:
        return jsonify({"ok": False, "erro": "Permissão negada."}), 403

    total = contar_perfis_crm(perfil["conta_id"])
    if total <= 1:
        return jsonify({"ok": False, "erro": "Não é possível deletar o único perfil da conta."}), 400
    deletar_perfil_crm(perfil_id)
    return jsonify({"ok": True})


# ── QUOTA DE LEADS ────────────────────────────────────────

@app.route("/api/quota", methods=["GET"])
@login_required
def api_quota():
    """Retorna a quota mensal de leads da conta do usuário logado."""
    conta_id = current_user.conta_id
    if not conta_id:
        # Sem conta vinculada: retorna limite zerado (não pode prospectar)
        return jsonify({
            "ok": True, "limite": 0, "usados": 0,
            "disponivel": 0, "reset_em": None,
            "sem_conta": True
        })
    try:
        resetar_quota_se_novo_mes(conta_id)
        q = get_quota(conta_id)
        return jsonify({
            "ok": True,
            "limite": q["limite"],
            "usados": q["usados"],
            "disponivel": q["disponivel"],
            "reset_em": q["reset_em"],
        })
    except Exception as e:
        return jsonify({"ok": False, "erro": str(e)}), 500


# ── ROTAS DE CONTAS (super_admin) ─────────────────────────

@app.route("/admin/conta/criar", methods=["POST"])
@login_required
def admin_criar_conta():
    if not current_user.is_super_admin:
        return jsonify({"ok": False, "erro": "Permissão negada."}), 403
    nome = (request.form.get("nome") or "").strip()
    if not nome:
        flash("Nome da conta é obrigatório.", "error")
        return redirect(url_for("admin"))
    criar_conta(nome)
    flash(f"Conta '{nome}' criada.", "success")
    return redirect(url_for("admin"))


@app.route("/admin/conta/<int:cid>/deletar", methods=["POST"])
@login_required
def admin_deletar_conta(cid):
    if not current_user.is_super_admin:
        return jsonify({"ok": False, "erro": "Permissão negada."}), 403
    deletar_conta(cid)
    flash("Conta removida.", "success")
    return redirect(url_for("admin"))


@app.route("/admin/conta/<int:cid>/renomear", methods=["POST"])
@login_required
def admin_renomear_conta(cid):
    if not current_user.is_super_admin:
        return jsonify({"ok": False, "erro": "Permissão negada."}), 403
    nome = (request.form.get("nome") or "").strip()
    if nome:
        renomear_conta(cid, nome)
    return redirect(url_for("admin"))


@app.route("/admin/entrar-como-conta/<int:cid>")
@login_required
def admin_entrar_como_conta(cid):
    from flask import session
    if not current_user.is_super_admin:
        return redirect(url_for("index"))
    membros = listar_usuarios(conta_id=cid)
    if not membros:
        flash("Esta conta não tem usuários vinculados. Crie um usuário primeiro.", "error")
        return redirect(url_for("admin"))
    session["_super_admin_id"] = current_user.id
    login_user(membros[0])
    return redirect(url_for("index"))


@app.route("/admin/sair-impersonacao")
@login_required
def admin_sair_impersonacao():
    from flask import session
    super_id = session.pop("_super_admin_id", None)
    if super_id:
        user = User.get(super_id)
        if user:
            login_user(user)
            return redirect(url_for("admin"))
    return redirect(url_for("index"))


# ────────────────────────────────────────────────────────────
# AUTOMACAO AGENDADA (systemd timer)
# ────────────────────────────────────────────────────────────

import subprocess

AUTOMACAO_CONFIG_PATH = Path("/opt/mia/workspace/prospeccao_ativa/prospeccao_schedule.json")
AUTOMACAO_SCRIPT_PATH = Path("/opt/mia/workspace/prospeccao_ativa/prospeccao_auto.py")
AUTOMACAO_LOG_PATH = Path("/opt/mia/logs/prospeccao_auto.log")
AUTOMACAO_PROGRESS_PATH = Path("/opt/mia/workspace/prospeccao_ativa/prospeccao_auto_progress.json")

DIAS_SEMANA_MAP = {
    "segunda": "Mon",
    "terca": "Tue",
    "quarta": "Wed",
    "quinta": "Thu",
    "sexta": "Fri",
    "sabado": "Sat",
    "domingo": "Sun",
}

DIAS_SEMANA_VALIDOS = list(DIAS_SEMANA_MAP.keys())


def _automacao_carregar_config() -> dict:
    """Le a config do JSON, garantindo campos padrao."""
    default = {
        "ativo": False,
        "limite_por_busca": 20,
        "palavras_chave": [],
        "nichos": [],           # legado: mantido pra compat
        "cidades": [],
        "horario": "08:00",
        "dias_semana": ["segunda", "quarta", "sexta"],
        "pais": "BR",
        "enriquecer": False,
        "filtro_sem_site": False,
        "crm": {},
    }
    try:
        if AUTOMACAO_CONFIG_PATH.exists():
            data = json.loads(AUTOMACAO_CONFIG_PATH.read_text(encoding="utf-8"))
            for k, v in default.items():
                data.setdefault(k, v)
            # Migração leve: se ainda tem só "nichos" e "palavras_chave" vazio, espelha
            if not data.get("palavras_chave") and data.get("nichos"):
                data["palavras_chave"] = list(data["nichos"])
            return data
    except Exception as e:
        print(f"[automacao] erro ao ler config: {e}", flush=True)
    return default


def _automacao_salvar_config(config: dict) -> None:
    tmp = AUTOMACAO_CONFIG_PATH.with_suffix(".json.tmp")
    tmp.write_text(json.dumps(config, ensure_ascii=False, indent=2), encoding="utf-8")
    tmp.replace(AUTOMACAO_CONFIG_PATH)


def _automacao_timer_status() -> dict:
    """Retorna se o timer systemd esta enabled/active."""
    try:
        enabled = subprocess.run(
            ["/bin/systemctl", "is-enabled", "prospeccao-auto.timer"],
            capture_output=True, text=True, timeout=5,
        ).stdout.strip()
    except Exception as e:
        enabled = f"erro: {e}"
    try:
        active = subprocess.run(
            ["/bin/systemctl", "is-active", "prospeccao-auto.timer"],
            capture_output=True, text=True, timeout=5,
        ).stdout.strip()
    except Exception as e:
        active = f"erro: {e}"
    return {"enabled": enabled, "active": active}


def _automacao_timer_ligar() -> tuple[bool, str]:
    try:
        subprocess.run(
            ["sudo", "-n", "/bin/systemctl", "enable", "prospeccao-auto.timer"],
            capture_output=True, text=True, timeout=10, check=True,
        )
        subprocess.run(
            ["sudo", "-n", "/bin/systemctl", "start", "prospeccao-auto.timer"],
            capture_output=True, text=True, timeout=10, check=True,
        )
        return True, ""
    except subprocess.CalledProcessError as e:
        return False, (e.stderr or e.stdout or str(e))[:400]
    except Exception as e:
        return False, str(e)[:400]


def _automacao_timer_desligar() -> tuple[bool, str]:
    err = ""
    try:
        subprocess.run(
            ["sudo", "-n", "/bin/systemctl", "stop", "prospeccao-auto.timer"],
            capture_output=True, text=True, timeout=10,
        )
        subprocess.run(
            ["sudo", "-n", "/bin/systemctl", "disable", "prospeccao-auto.timer"],
            capture_output=True, text=True, timeout=10,
        )
        return True, ""
    except Exception as e:
        return False, str(e)[:400]


@app.route("/automacao", methods=["GET"])
@login_required
def automacao_get():
    """Retorna a configuracao atual e status do timer."""
    config = _automacao_carregar_config()
    status = _automacao_timer_status()
    # Sanitiza token no retorno
    safe_config = dict(config)
    crm = config.get("crm") or {}
    safe_config["crm"] = {
        "location_id": crm.get("location_id", ""),
        "pipeline_id": crm.get("pipeline_id", ""),
        "stage_id": crm.get("stage_id", ""),
        "configurado": bool(crm.get("token")),
    }
    return jsonify({
        "ok": True,
        "config": safe_config,
        "timer": status,
        "dias_semana_validos": DIAS_SEMANA_VALIDOS,
    })


@app.route("/automacao/status", methods=["GET"])
@login_required
def automacao_status():
    """Retorna o progresso da rodada em andamento (arquivo escrito pelo prospeccao_auto.py)."""
    try:
        if not AUTOMACAO_PROGRESS_PATH.exists():
            return jsonify({"ok": True, "progresso": None})
        data = json.loads(AUTOMACAO_PROGRESS_PATH.read_text(encoding="utf-8"))
        return jsonify({"ok": True, "progresso": data})
    except Exception as e:
        return jsonify({"ok": False, "erro": str(e)}), 500


@app.route("/automacao", methods=["POST"])
@login_required
def automacao_save():
    """Salva nova config. NAO altera estado do timer aqui."""
    data = request.get_json() or {}

    # Aceita "palavras_chave" (novo) ou "nichos" (legado)
    palavras = data.get("palavras_chave") or data.get("nichos") or []
    cidades = data.get("cidades") or []
    limite = int(data.get("limite_por_busca") or 20)
    horario = (data.get("horario") or "08:00").strip()
    dias = data.get("dias_semana") or []
    ativo = bool(data.get("ativo", False))
    pais = (data.get("pais") or "BR").strip().upper()
    enriquecer = bool(data.get("enriquecer", False))
    filtro_sem_site = bool(data.get("filtro_sem_site", False))

    # Sanitiza listas
    palavras = [str(x).strip() for x in palavras if str(x).strip()]
    cidades = [str(x).strip() for x in cidades if str(x).strip()]
    dias = [d for d in dias if d in DIAS_SEMANA_VALIDOS]

    if pais not in ("BR", "US"):
        pais = "BR"

    # Valida horario HH:MM
    if not re.match(r"^\d{2}:\d{2}$", horario):
        return jsonify({"ok": False, "erro": "Horario invalido. Use HH:MM"}), 400

    if limite < 1 or limite > 500:
        return jsonify({"ok": False, "erro": "Limite deve estar entre 1 e 500"}), 400

    # ── Bloco CRM: se veio no request, snapshot dele; senão herda da conta ──
    crm_in = data.get("crm") or {}
    crm_block: dict = {}
    if crm_in.get("token") and crm_in.get("location_id"):
        crm_block = {
            "token": str(crm_in.get("token")).strip(),
            "location_id": str(crm_in.get("location_id")).strip(),
            "pipeline_id": str(crm_in.get("pipeline_id") or "").strip(),
            "stage_id": str(crm_in.get("stage_id") or "").strip(),
            "conta_id": current_user.conta_id,
        }
    else:
        # Fallback: puxa da conta do usuario logado se existir
        conta_id = current_user.conta_id
        if conta_id:
            try:
                cfg_conta = get_ghl_config(conta_id) or {}
                if cfg_conta.get("token") and cfg_conta.get("location_id"):
                    crm_block = {
                        "token": cfg_conta["token"],
                        "location_id": cfg_conta["location_id"],
                        "pipeline_id": cfg_conta.get("pipeline_id") or "",
                        "stage_id": cfg_conta.get("stage_id") or "",
                        "conta_id": conta_id,
                    }
            except Exception:
                pass

    config = {
        "ativo": ativo,
        "limite_por_busca": limite,
        "palavras_chave": palavras,
        "nichos": palavras,  # espelho pra compat com scripts antigos
        "cidades": cidades,
        "horario": horario,
        "dias_semana": dias,
        "pais": pais,
        "enriquecer": enriquecer,
        "filtro_sem_site": filtro_sem_site,
        "crm": crm_block,
    }

    try:
        _automacao_salvar_config(config)
    except Exception as e:
        return jsonify({"ok": False, "erro": f"Falha ao salvar: {e}"}), 500

    # Nao devolve token no response (evita eco de credencial)
    safe_config = dict(config)
    safe_config["crm"] = {
        "location_id": crm_block.get("location_id", ""),
        "pipeline_id": crm_block.get("pipeline_id", ""),
        "stage_id": crm_block.get("stage_id", ""),
        "configurado": bool(crm_block.get("token")),
    }
    return jsonify({"ok": True, "config": safe_config})


@app.route("/automacao/toggle", methods=["POST"])
@login_required
def automacao_toggle():
    """Liga/desliga a automacao. Atualiza campo `ativo` no JSON e o timer systemd."""
    data = request.get_json() or {}
    ativar = bool(data.get("ativar", False))

    config = _automacao_carregar_config()
    config["ativo"] = ativar

    try:
        _automacao_salvar_config(config)
    except Exception as e:
        return jsonify({"ok": False, "erro": f"Falha ao salvar config: {e}"}), 500

    if ativar:
        ok, err = _automacao_timer_ligar()
    else:
        ok, err = _automacao_timer_desligar()

    status = _automacao_timer_status()
    return jsonify({
        "ok": ok,
        "erro": err if not ok else None,
        "config": config,
        "timer": status,
    })


@app.route("/automacao/rodar-agora", methods=["POST"])
@login_required
def automacao_rodar_agora():
    """Executa a prospeccao imediatamente em background, sem alterar o timer."""
    if not AUTOMACAO_SCRIPT_PATH.exists():
        return jsonify({"ok": False, "erro": "Script prospeccao_auto.py nao encontrado"}), 500

    # Antes de disparar, força o campo "ativo" a valer (mesmo que o timer esteja off)
    # e garante que a config atual será usada. Não altera o toggle salvo.
    config = _automacao_carregar_config()
    if not (config.get("palavras_chave") or config.get("nichos")) or not config.get("cidades"):
        return jsonify({"ok": False, "erro": "Configure palavras-chave e cidades antes de rodar."}), 400
    if not (config.get("crm") or {}).get("token"):
        return jsonify({"ok": False, "erro": "CRM não configurado para a automação."}), 400

    # Reseta arquivo de progresso pra UI limpar
    try:
        AUTOMACAO_PROGRESS_PATH.write_text(json.dumps({
            "status": "iniciando",
            "inicio": datetime.now().isoformat(),
            "fim": None,
            "combinacoes_total": len(config.get("cidades", [])) * len(
                config.get("palavras_chave") or config.get("nichos") or []
            ),
            "combinacoes_feitas": 0,
            "combinacao_atual": None,
            "encontrados": 0,
            "enviados": 0,
            "duplicados": 0,
            "leads": [],
            "erro": None,
            "pais": config.get("pais", "BR"),
            "log_tail": [],
        }, ensure_ascii=False, indent=2), encoding="utf-8")
    except Exception:
        pass

    try:
        # Roda em background e nao bloqueia a resposta.
        AUTOMACAO_LOG_PATH.parent.mkdir(parents=True, exist_ok=True)
        log_fp = open(AUTOMACAO_LOG_PATH, "a", encoding="utf-8")
        # Precisa forçar o "ativo=true" na hora só pra permitir rodar-agora?
        # Não: o script já ignora `ativo` quando recebe overrides via CLI. Como não
        # passamos override aqui, deixamos o script respeitar o JSON. Se o toggle
        # estiver off, precisamos permitir rodar. Passamos --palavra e --cidade
        # explicitamente pra bypassar a checagem de `ativo`.
        cli_args = ["/usr/bin/python3", str(AUTOMACAO_SCRIPT_PATH)]
        for p in (config.get("palavras_chave") or config.get("nichos") or []):
            cli_args += ["--palavra", p]
        for c in config.get("cidades", []):
            cli_args += ["--cidade", c]
        cli_args += ["--limite", str(config.get("limite_por_busca") or 20)]
        if bool(config.get("filtro_sem_site", False)):
            cli_args += ["--filtro-sem-site"]

        subprocess.Popen(
            cli_args,
            cwd=str(AUTOMACAO_SCRIPT_PATH.parent),
            stdout=log_fp, stderr=log_fp,
            start_new_session=True,
        )
    except Exception as e:
        return jsonify({"ok": False, "erro": str(e)}), 500

    return jsonify({"ok": True, "mensagem": "Prospeccao disparada em background"})


@app.route("/geocode")
@login_required
def geocode():
    """Converte nome de cidade/endereço em lat/lng via Google Geocoding API.

    Usa a mesma chave GOOGLE_PLACES_KEY já configurada (a chave do Google Cloud
    normalmente cobre Places + Geocoding no mesmo projeto).
    """
    q = (request.args.get("q") or "").strip()
    if not q:
        return jsonify({"ok": False, "error": "q obrigatorio"}), 400

    try:
        from config import GOOGLE_PLACES_KEY as api_key
    except ImportError:
        api_key = os.environ.get("GOOGLE_PLACES_KEY", "")
    if not api_key:
        return jsonify({"ok": False, "error": "GOOGLE_PLACES_KEY nao configurada"}), 500

    try:
        r = requests.get(
            "https://maps.googleapis.com/maps/api/geocode/json",
            params={"address": q, "key": api_key, "language": "pt-BR"},
            timeout=10,
        )
        data = r.json()
    except Exception as e:
        return jsonify({"ok": False, "error": f"falha ao chamar geocoding: {str(e)[:120]}"}), 502

    if data.get("status") != "OK" or not data.get("results"):
        return jsonify({"ok": False, "error": "nao encontrado", "status": data.get("status")}), 404

    loc = data["results"][0]["geometry"]["location"]
    nome = data["results"][0].get("formatted_address", q)
    return jsonify({"ok": True, "lat": loc["lat"], "lng": loc["lng"], "nome": nome})


@app.route("/prospectar", methods=["POST"])
@login_required
def prospectar():
    if not current_user.ativo:
        return jsonify({"erro": "Conta inativa. Entre em contato com o administrador."}), 403
    data = request.get_json()
    nicho = (data.get("nicho") or "").strip()
    cidade = (data.get("cidade") or "").strip()
    limite = int(data.get("limite") or 20)
    crm_type = (data.get("crm_type") or "ghl").strip().lower()
    token = (data.get("token") or "").strip()

    if not nicho or not cidade:
        return jsonify({"erro": "Nicho e cidade sao obrigatorios"}), 400

    # Restrição por raio no mapa (opcional). Quando presente, a busca da Places API
    # usa locationRestriction (circle) em vez de depender só do texto da cidade.
    try:
        lat_raw = data.get("lat")
        lng_raw = data.get("lng")
        raio_raw = data.get("raio_km")
        lat_val = float(lat_raw) if lat_raw not in (None, "", "null") else None
        lng_val = float(lng_raw) if lng_raw not in (None, "", "null") else None
        raio_km_val = int(raio_raw) if raio_raw not in (None, "", "null") else 0
    except (TypeError, ValueError):
        lat_val = None
        lng_val = None
        raio_km_val = 0
    # Clamp de segurança. Google Places (New) aceita circle radius até 50000m em
    # locationRestriction. Pra raios maiores usamos locationBias (mais permissivo).
    # Aqui só evitamos valores absurdos/negativos.
    if raio_km_val and raio_km_val < 0:
        raio_km_val = 0
    if raio_km_val and raio_km_val > 500:
        raio_km_val = 500

    # Campos CRM
    location_id = (data.get("location_id") or "").strip()
    pipeline_id = (data.get("pipeline_id") or "").strip()
    stage_id = (data.get("stage_id") or "").strip()
    connection_id = (data.get("connection_id") or "").strip()

    pais = (data.get("pais") or "BR").strip().upper()
    # enriquecimento_auto=True significa modo "Enviar ao CRM automaticamente":
    # enriquece e envia cada lead. False = modo "Análise": leads ficam na tabela.
    enriquecimento_auto = bool(data.get("enriquecimento_auto", False))
    filtro_sem_site = bool(data.get("filtro_sem_site", False))
    # Filtros GMB pré-prospecção foram removidos (filtro_nota, filtro_site_gmb, score_minimo).
    # Coleta tudo agora — usuário filtra depois na tabela de resultados.
    filtro_nota = ""
    filtro_site_gmb = ""
    score_minimo = 0

    # Só exige credenciais de CRM quando o modo é envio automático
    if enriquecimento_auto:
        if not token:
            return jsonify({"erro": "Credenciais do CRM incompletas."}), 400
        if crm_type == "ghl" and not (location_id and pipeline_id and stage_id):
            return jsonify({"erro": "Credenciais do CRM incompletas. Configure o CRM antes de prospectar."}), 400

    limite = max(5, min(200, limite))

    # ── Verificação de quota mensal da conta ──
    conta_id = current_user.conta_id
    if conta_id:
        try:
            resetar_quota_se_novo_mes(conta_id)
            quota = get_quota(conta_id)
            if quota["usados"] >= quota["limite"]:
                return jsonify({
                    "ok": False,
                    "erro": f"Quota mensal atingida. Limite: {quota['limite']} leads/mês."
                }), 200
        except Exception:
            pass

    job_id = _novo_job(nicho, cidade, limite)
    jobs[job_id]["conta_id"] = conta_id
    try:
        jobs[job_id]["user_id"] = int(current_user.id)
    except Exception:
        jobs[job_id]["user_id"] = None
    try:
        salvar_job_historico(jobs[job_id])
    except Exception as e:
        print(f"[jobs_historico] erro ao vincular conta/user {job_id}: {e}")

    t = threading.Thread(
        target=run_prospector,
        args=(job_id, nicho, cidade, limite, token, location_id, pipeline_id, stage_id,
              crm_type, connection_id, pais, enriquecimento_auto, conta_id, filtro_sem_site,
              filtro_nota, filtro_site_gmb, score_minimo),
        kwargs={"lat": lat_val, "lng": lng_val, "raio_km": raio_km_val},
        daemon=True
    )
    t.start()

    return jsonify({"job_id": job_id, "mensagem": f"Job {job_id} iniciado"})


@app.route("/status/<job_id>")
@login_required
def status(job_id):
    job = jobs.get(job_id)
    # Isolamento por conta: quem não é super_admin só vê job da própria conta.
    if job and not current_user.is_super_admin:
        dono = job.get("conta_id")
        if dono is not None and dono != current_user.conta_id:
            return jsonify({"status": "erro", "erro": "Job nao encontrado", "pct": 0, "leads": [], "logs": []}), 404
    if not job:
        # Fallback: busca no histórico persistido (job saiu da memória)
        historico = listar_jobs_historico(limit=50)
        job_hist = next((j for j in historico if j["id"] == job_id), None)
        if job_hist:
            return jsonify({
                "id": job_hist["id"],
                "status": job_hist["status"] if job_hist["status"] != "aguardando" else "erro",
                "progresso": job_hist["progresso"],
                "total": job_hist["total"],
                "pct": 100 if job_hist["status"] == "concluido" else 0,
                "leads_enviados": job_hist["leads_enviados"],
                "logs": ["Job recuperado do histórico — dados em memória foram perdidos (limite de jobs atingido). Refaça o enriquecimento."],
                "erro": "Dados do job foram perdidos da memória. Refaça o enriquecimento.",
                "concluido_em": job_hist["concluido_em"],
                "leads": [],
            })
        return jsonify({"status": "erro", "erro": "Job nao encontrado", "pct": 0, "leads": [], "logs": []}), 404

    pct = 0
    if job["total"] > 0:
        pct = round((job["progresso"] / job["total"]) * 100)
    if job["status"] == "concluido":
        pct = 100

    return jsonify({
        "id": job["id"],
        "status": job["status"],
        "progresso": job["progresso"],
        "total": job["total"],
        "pct": pct,
        "leads_enviados": job["leads_enviados"],
        "logs": job["logs"][-10:],
        "erro": job["erro"],
        "concluido_em": job["concluido_em"],
        "leads": job["leads"],
        "crm_enviado": job.get("crm_enviado", False),
    })


@app.route("/jobs")
@login_required
def listar_jobs():
    # Jobs em memória (sessão atual) — quem não é super_admin só vê os da própria conta.
    if current_user.is_super_admin:
        em_memoria = {j["id"]: j for j in jobs.values()}
    else:
        em_memoria = {
            j["id"]: j for j in jobs.values()
            if j.get("conta_id") is None or j.get("conta_id") == current_user.conta_id
        }

    # Jobs históricos do banco (filtra pela conta do usuário ou todos se super_admin)
    conta_id = current_user.conta_id if not current_user.is_super_admin else None
    try:
        historico_db = listar_jobs_historico(conta_id=conta_id, limit=20)
    except Exception as e:
        print(f"[jobs_historico] erro ao listar: {e}")
        historico_db = []

    # Mescla: memória tem prioridade
    merged = {}
    for row in historico_db:
        merged[row["id"]] = row
    for jid, job in em_memoria.items():
        merged[jid] = {
            "id": job["id"],
            "tipo": job.get("tipo", "prospeccao"),
            "nicho": job.get("nicho", "—"),
            "cidade": job.get("cidade", "—"),
            "limite": job.get("limite", job.get("total", 0)),
            "status": job["status"],
            "progresso": job["progresso"],
            "total": job["total"],
            "leads_enviados": job["leads_enviados"],
            "criado_em": job["criado_em"],
            "concluido_em": job["concluido_em"],
        }

    lista = sorted(merged.values(), key=lambda j: j.get("criado_em") or "", reverse=True)[:10]
    return jsonify(lista)


@app.route("/cancelar/<job_id>", methods=["POST"])
@login_required
def cancelar(job_id):
    job = jobs.get(job_id)
    if not job:
        return jsonify({"erro": "Job nao encontrado"}), 404
    if job["status"] not in ("aguardando", "em_andamento"):
        return jsonify({"erro": f"Job ja esta em status '{job['status']}'"}), 400
    job["status"] = "cancelado"
    return jsonify({"mensagem": f"Job {job_id} marcado para cancelamento"})


@app.route("/pausar/<job_id>", methods=["POST"])
@login_required
def pausar(job_id):
    job = jobs.get(job_id)
    if not job:
        return jsonify({"erro": "Job nao encontrado"}), 404
    if job["status"] not in ("aguardando", "em_andamento"):
        return jsonify({"erro": f"Job ja esta em status '{job['status']}'"}), 400
    job["status"] = "pausado"
    return jsonify({"mensagem": f"Job {job_id} pausado - enviando leads coletados ao CRM"})


@app.route("/enviar-todos-crm", methods=["POST"])
@login_required
def enviar_todos_crm():
    """Envia todos os leads não enviados de um job ao CRM configurado.

    Aceita duas formas de fornecer os leads:
    1. job_id: busca o job em memória (comportamento padrão)
    2. leads: lista de leads enviados diretamente no body (fallback quando o
       job em memória foi perdido - ex: servidor reiniciado)
    """
    data = request.get_json() or {}
    job_id = data.get("job_id")
    job = jobs.get(job_id)

    # Prioriza os leads enviados no body (permite envio selecionado).
    # Só usa job["leads"] se o body não trouxer lista (não deve acontecer com o frontend atual).
    leads_raw = data.get("leads") or []
    if leads_raw:
        leads_para_enviar = [l for l in leads_raw if not l.get("enviado_crm")]
    elif job:
        leads_para_enviar = [l for l in job.get("leads", []) if not l.get("enviado_crm")]
    else:
        return jsonify({"erro": "Job nao encontrado e nenhum lead enviado"}), 404

    crm_type = (data.get("crm_type") or "ghl").strip().lower()
    token = (data.get("token") or "").strip()
    location_id = (data.get("location_id") or "").strip()
    pipeline_id = (data.get("pipeline_id") or "").strip()
    stage_id = (data.get("stage_id") or "").strip()
    connection_id = (data.get("connection_id") or "").strip()

    if not token:
        return jsonify({"erro": "Token do CRM ausente"}), 400
    if crm_type == "ghl" and not (location_id and pipeline_id and stage_id):
        return jsonify({"erro": "Credenciais GHL incompletas"}), 400

    enviados = 0
    falhas = 0
    erros_detalhes = []
    for lead in leads_para_enviar:
        if lead.get("enviado_crm"):
            continue
        try:
            if crm_type == "unnichat":
                contact_id, err = unnichat_criar_contato(
                    lead, token, ["prospeccao-ativa", lead.get("nicho", ""), lead.get("cidade", "")],
                    connection_id=connection_id
                )
                if contact_id:
                    unnichat_adicionar_pipeline(
                        contact_id, token, connection_id,
                        pipeline_id=pipeline_id or "8TgKxh0VK37bGNkbN8yO",
                        column_id=stage_id or "HThNTr24Yy2fN0vmeLm2",
                        business_name=lead["nome"]
                    )
                    lead["enviado_crm"] = True
                    lead["contact_id"] = contact_id
                    enviados += 1
                else:
                    falhas += 1
                    erros_detalhes.append({"nome": lead["nome"], "erro": err or "sem contact_id"})
            else:
                contact_id, err = criar_contato_ghl(
                    lead, token, location_id,
                    custom_fields_map=_carregar_campos_crm(location_id) or None,
                )
                if contact_id:
                    ok, opp_id, err_opp = criar_oportunidade_ghl(
                        lead, contact_id, token, location_id, pipeline_id, stage_id
                    )
                    is_dup = err_opp and "duplicate" in err_opp.lower()
                    if ok or is_dup:
                        lead["enviado_crm"] = True
                        lead["contact_id"] = contact_id
                        enviados += 1
                    else:
                        falhas += 1
                        erros_detalhes.append({"nome": lead["nome"], "erro": f"oportunidade: {err_opp}"})
                else:
                    falhas += 1
                    erros_detalhes.append({"nome": lead["nome"], "erro": f"contato: {err}"})
            time.sleep(0.3)
        except Exception as exc:
            falhas += 1
            erros_detalhes.append({"nome": lead.get("nome", "?"), "erro": str(exc)})

    if job:
        # Atualiza os objetos do job em memória para refletir o envio selecionado
        if leads_raw:
            enviados_nomes = {l["nome"] for l in leads_para_enviar if l.get("enviado_crm")}
            for jl in job.get("leads", []):
                if jl["nome"] in enviados_nomes:
                    jl["enviado_crm"] = True
            job["leads_enviados"] = sum(1 for l in job.get("leads", []) if l.get("enviado_crm"))
        else:
            job["leads_enviados"] = job.get("leads_enviados", 0) + enviados
    return jsonify({"ok": True, "enviados": enviados, "falhas": falhas, "erros": erros_detalhes})


@app.route("/enriquecer", methods=["POST"])
@login_required
def enriquecer():
    """Recebe CSV com nomes de empresas ou lista JSON, enriquece e envia ao CRM."""
    token = ""
    location_id = ""
    pipeline_id = ""
    stage_id = ""
    empresas = []

    if request.is_json:
        data = request.get_json()
        token = (data.get("token") or "").strip()
        location_id = (data.get("location_id") or "").strip()
        pipeline_id = (data.get("pipeline_id") or "").strip()
        stage_id = (data.get("stage_id") or "").strip()
        empresas = data.get("empresas", [])
    else:
        token = (request.form.get("token") or "").strip()
        location_id = (request.form.get("location_id") or "").strip()
        pipeline_id = (request.form.get("pipeline_id") or "").strip()
        stage_id = (request.form.get("stage_id") or "").strip()

        # Aceita CSV upload
        if "arquivo" in request.files:
            f = request.files["arquivo"]
            content = f.read().decode("utf-8", errors="ignore")
            reader = csv.reader(io.StringIO(content))
            for row in reader:
                if row and row[0].strip() and row[0].strip().lower() != "empresa":
                    empresas.append(row[0].strip())

        # Aceita texto colado (um nome por linha)
        lista_texto = (request.form.get("lista_texto") or "").strip()
        if lista_texto and not empresas:
            for linha in lista_texto.splitlines():
                nome = linha.strip()
                if nome:
                    empresas.append(nome)

    if not empresas:
        return jsonify({"erro": "Nenhuma empresa na lista"}), 400
    if not token or not location_id:
        return jsonify({"erro": "Token e Location ID obrigatórios"}), 400

    empresas = empresas[:100]  # limite de 100 por vez
    job_id = _novo_job_enriq(empresas)

    # Salva creds no job para uso posterior (envio manual ao CRM)
    jobs[job_id]["crm_token"] = token
    jobs[job_id]["crm_location_id"] = location_id
    jobs[job_id]["crm_pipeline_id"] = pipeline_id
    jobs[job_id]["crm_stage_id"] = stage_id
    jobs[job_id]["crm_enviado"] = False
    jobs[job_id]["conta_id"] = current_user.conta_id
    try:
        jobs[job_id]["user_id"] = int(current_user.id)
    except Exception:
        jobs[job_id]["user_id"] = None
    try:
        salvar_job_historico(jobs[job_id])
    except Exception as e:
        print(f"[jobs_historico] erro ao vincular conta/user {job_id}: {e}")

    t = threading.Thread(
        target=run_enriquecedor,
        args=(job_id, empresas, token, location_id, pipeline_id, stage_id),
        kwargs={"enviar_crm": False},
        daemon=True
    )
    t.start()

    return jsonify({"job_id": job_id, "mensagem": f"Enriquecendo {len(empresas)} empresas"})


@app.route("/enriquecer/exportar-csv")
@login_required
def enriquecer_exportar_csv():
    """Exporta os leads enriquecidos de um job como arquivo CSV.

    Uso: GET /enriquecer/exportar-csv?job_id=XXX
    Retorna um CSV com as colunas principais dos leads enriquecidos.
    """
    job_id = request.args.get("job_id")
    job = jobs.get(job_id)
    leads = (job or {}).get("leads") or _carregar_leads_enriq(job_id or "")
    if not leads:
        return jsonify({"erro": "Nenhum lead disponível"}), 404

    campos = [
        "nome", "telefone", "email", "site", "endereco", "cnpj",
        "decisor_nome", "decisor_cargo", "linkedin_url",
        "instagram", "facebook", "fontes",
    ]
    output = io.StringIO()
    writer = csv.DictWriter(output, fieldnames=campos, extrasaction="ignore")
    writer.writeheader()
    for lead in leads:
        row = {k: lead.get(k, "") for k in campos}
        if isinstance(row.get("fontes"), list):
            row["fontes"] = ", ".join(str(f) for f in row["fontes"])
        writer.writerow(row)

    from flask import Response
    return Response(
        output.getvalue(),
        mimetype="text/csv",
        headers={
            "Content-Disposition": f"attachment;filename=leads_enriquecidos_{job_id}.csv"
        },
    )


@app.route("/enriquecer/excluir-arquivo", methods=["POST"])
@login_required
def enriquecer_excluir_arquivo():
    """Remove o arquivo de enriquecimento carregado.

    O upload no fluxo atual é processado em memória (não é salvo em disco),
    então esta rota é idempotente: ela apenas confirma a limpeza e remove
    qualquer arquivo temporário associado ao usuário, se existir. A UI usa
    esta rota para resetar o estado do formulário no cliente.
    """
    import os as _os

    nome = ""
    try:
        if request.is_json:
            data = request.get_json(silent=True) or {}
            nome = (data.get("nome") or "").strip()
        else:
            nome = (request.form.get("nome") or "").strip()
    except Exception:
        nome = ""

    # Limpa qualquer arquivo temporário do diretório de uploads deste usuário, se existir
    removidos = []
    try:
        upload_dir = _os.path.join(_os.path.dirname(_os.path.abspath(__file__)), "uploads")
        if _os.path.isdir(upload_dir):
            user_id = str(session.get("user_id") or "").strip()
            for fname in _os.listdir(upload_dir):
                if user_id and not fname.startswith(f"{user_id}_"):
                    continue
                caminho = _os.path.join(upload_dir, fname)
                try:
                    _os.remove(caminho)
                    removidos.append(fname)
                except OSError:
                    pass
    except Exception:
        pass

    return jsonify({
        "ok": True,
        "mensagem": "Arquivo removido",
        "nome": nome,
        "removidos": removidos,
    })


# ── PLANILHA XLSX (upload multi-aba, enriquece antes, envia ao CRM depois) ────

@app.route("/planilha/abas", methods=["POST"])
@login_required
def planilha_abas():
    """Recebe .xlsx e retorna lista de abas com contagem de linhas não-vazias."""
    import openpyxl
    f = request.files.get("arquivo")
    if not f:
        return jsonify({"erro": "Arquivo não enviado"}), 400
    try:
        wb = openpyxl.load_workbook(f, read_only=True, data_only=True)
    except Exception as e:
        return jsonify({"erro": f"Falha ao ler planilha: {e}"}), 400

    abas = []
    for nome_aba in wb.sheetnames:
        ws = wb[nome_aba]
        total = 0
        for row in ws.iter_rows(min_row=1, values_only=True):
            val = row[0] if row else None
            if val is not None and str(val).strip():
                total += 1
        abas.append({"nome": nome_aba, "total": total})
    wb.close()
    return jsonify({"ok": True, "abas": abas})


@app.route("/planilha/carregar", methods=["POST"])
@login_required
def planilha_carregar():
    """Recebe .xlsx + abas selecionadas, retorna lista consolidada de nomes."""
    import openpyxl
    import json as _json
    f = request.files.get("arquivo")
    if not f:
        return jsonify({"erro": "Arquivo não enviado"}), 400
    try:
        abas_sel = _json.loads(request.form.get("abas_selecionadas", "[]"))
    except Exception:
        abas_sel = []
    if not abas_sel:
        return jsonify({"erro": "Nenhuma aba selecionada"}), 400

    try:
        wb = openpyxl.load_workbook(f, read_only=True, data_only=True)
    except Exception as e:
        return jsonify({"erro": f"Falha ao ler planilha: {e}"}), 400

    HEADERS = {"empresa", "nome", "lead", "company", "cliente", "razao social", "razão social"}
    itens = []
    for nome_aba in wb.sheetnames:
        if nome_aba not in abas_sel:
            continue
        ws = wb[nome_aba]
        for row in ws.iter_rows(min_row=1, values_only=True):
            val = row[0] if row else None
            if val is None:
                continue
            n = str(val).strip()
            if not n:
                continue
            if n.lower() in HEADERS:
                continue
            itens.append({"nome": n, "aba": nome_aba})
    wb.close()

    nomes = [i["nome"] for i in itens]
    return jsonify({"ok": True, "itens": itens, "nomes": nomes, "total": len(nomes)})


@app.route("/planilha/enriquecer", methods=["POST"])
@login_required
def planilha_enriquecer():
    """Enriquece lista SEM enviar ao CRM. Envio manual via /planilha/enviar-crm."""
    data = request.get_json(force=True) or {}
    empresas = data.get("empresas", [])
    if not empresas:
        return jsonify({"erro": "Lista vazia"}), 400
    empresas = [str(e).strip() for e in empresas if str(e).strip()]
    empresas = empresas[:200]
    job_id = _novo_job_enriq(empresas)
    jobs[job_id]["origem"] = "planilha"
    jobs[job_id]["conta_id"] = current_user.conta_id
    try:
        jobs[job_id]["user_id"] = int(current_user.id)
    except Exception:
        jobs[job_id]["user_id"] = None
    try:
        salvar_job_historico(jobs[job_id])
    except Exception as e:
        print(f"[jobs_historico] erro ao vincular conta/user {job_id}: {e}")

    t = threading.Thread(
        target=run_enriquecedor,
        args=(job_id, empresas, "", "", "", ""),
        kwargs={"enviar_crm": False},
        daemon=True,
    )
    t.start()
    return jsonify({"ok": True, "job_id": job_id, "total": len(empresas)})


@app.route("/planilha/enviar-crm", methods=["POST"])
@login_required
def planilha_enviar_crm():
    """Pega leads já enriquecidos de um job e envia ao CRM."""
    data = request.get_json(force=True) or {}
    job_id = data.get("job_id")
    token = (data.get("token") or "").strip()
    location_id = (data.get("location_id") or "").strip()
    pipeline_id = (data.get("pipeline_id") or "").strip()
    stage_id = (data.get("stage_id") or "").strip()

    if not job_id or job_id not in jobs:
        return jsonify({"erro": "Job não encontrado"}), 404

    job = jobs[job_id]

    if not token:
        token = job.get("crm_token", "")
    if not location_id:
        location_id = job.get("crm_location_id", "")
    if not pipeline_id:
        pipeline_id = job.get("crm_pipeline_id", "")
    if not stage_id:
        stage_id = job.get("crm_stage_id", "")

    if not token or not location_id:
        return jsonify({"erro": "Token e Location ID obrigatórios"}), 400

    leads_enriquecidos = list(job.get("leads", []))
    if not leads_enriquecidos:
        return jsonify({"erro": "Nenhum lead enriquecido neste job"}), 400

    if token.lower().startswith("bearer "):
        token = token[7:].strip()

    def _enviar():
        enviados = 0
        erros = 0
        job["logs"].append(f"[{_ts()}] Iniciando envio ao CRM: {len(leads_enriquecidos)} leads")
        for dados in leads_enriquecidos:
            nome_empresa = dados.get("nome", "")
            origem_tag = "lista-planilha" if job.get("origem") == "planilha" else "lista-enriquecida"
            tags = [origem_tag, "enriquecido"]
            if dados.get("linkedin_url"):
                tags.append("linkedin-encontrado")
            if dados.get("cnpj"):
                tags.append("cnpj-encontrado")
            lead = {
                "nome": dados.get("decisor_nome") or nome_empresa,
                "telefone": dados.get("telefone", ""),
                "email": dados.get("email", ""),
                "empresa": nome_empresa,
                "endereco": dados.get("endereco", ""),
                "site": dados.get("site", ""),
                "tags": tags,
                # Campos que alimentam custom fields no CRM (Sweet Angels etc.)
                "cnpj": dados.get("cnpj", ""),
                "decisor_nome": dados.get("decisor_nome", ""),
                "decisor_cargo": dados.get("decisor_cargo", ""),
                "linkedin_url": dados.get("linkedin_url", ""),
                "instagram": dados.get("instagram", ""),
                "facebook": dados.get("facebook", ""),
                "fontes": dados.get("fontes") or [],
            }
            try:
                contact_id, err = criar_contato_ghl(
                    lead, token, location_id, source="Planilha",
                    custom_fields_map=_carregar_campos_crm(location_id) or None,
                )
                if contact_id:
                    opportunity_id = None
                    if pipeline_id and stage_id:
                        ok, opportunity_id, _ = criar_oportunidade_ghl(
                            lead, contact_id, token, location_id, pipeline_id, stage_id, source="Planilha"
                        )
                    _adicionar_nota_ghl(contact_id, token, location_id, dados, opportunity_id=opportunity_id)
                    enviados += 1
                    job["logs"].append(f"[{_ts()}] ✓ {nome_empresa} → CRM (contato criado)")
                else:
                    erros += 1
                    job["logs"].append(f"[{_ts()}] ✗ {nome_empresa} → erro CRM: {err}")
            except Exception as e:
                erros += 1
                job["logs"].append(f"[{_ts()}] ✗ {nome_empresa} → exceção: {e}")
            time.sleep(0.5)

        job["leads_enviados"] = enviados
        job["crm_erros"] = erros
        job["crm_enviado"] = True
        job["logs"].append(
            f"[{_ts()}] Envio ao CRM concluído: {enviados}/{len(leads_enriquecidos)} contatos criados"
            + (f" ({erros} com erro)" if erros else "")
        )

    threading.Thread(target=_enviar, daemon=True).start()
    return jsonify({"ok": True, "total": len(leads_enriquecidos)})


# Arquivo local para guardar IDs dos campos personalizados criados no CRM
CAMPOS_PATH = Path("/opt/mia/workspace/prospeccao_ativa/ghl_campos.json")

CAMPOS_ENRIQUECIMENTO_BR = [
    {"name": "CNPJ",            "dataType": "TEXT", "placeholder": "00.000.000/0001-00"},
    {"name": "LinkedIn",        "dataType": "TEXT", "placeholder": "https://linkedin.com/in/..."},
    {"name": "Decisor",         "dataType": "TEXT", "placeholder": "Nome do sócio/decisor"},
    {"name": "Cargo do Decisor","dataType": "TEXT", "placeholder": "Ex: Sócio-Administrador"},
    {"name": "Site",            "dataType": "TEXT", "placeholder": "https://empresa.com.br"},
    {"name": "Instagram",       "dataType": "TEXT", "placeholder": "https://instagram.com/empresa"},
    {"name": "Facebook",        "dataType": "TEXT", "placeholder": "https://facebook.com/empresa"},
]

CAMPOS_ENRIQUECIMENTO_US = [
    {"name": "Email",           "dataType": "TEXT", "placeholder": "email@empresa.com"},
    {"name": "LinkedIn",        "dataType": "TEXT", "placeholder": "https://linkedin.com/in/..."},
    {"name": "Decisor",         "dataType": "TEXT", "placeholder": "Owner / CEO / Founder"},
    {"name": "Cargo do Decisor","dataType": "TEXT", "placeholder": "Ex: Owner, CEO, Manager"},
    {"name": "Site",            "dataType": "TEXT", "placeholder": "https://company.com"},
    {"name": "Instagram",       "dataType": "TEXT", "placeholder": "https://instagram.com/company"},
    {"name": "Facebook",        "dataType": "TEXT", "placeholder": "https://facebook.com/company"},
]

# Mantém compatibilidade — padrão BR
CAMPOS_ENRIQUECIMENTO = CAMPOS_ENRIQUECIMENTO_BR


def _carregar_campos_ghl() -> dict:
    """Retorna mapeamento nome_campo -> field_id salvo localmente."""
    try:
        if CAMPOS_PATH.exists():
            return json.loads(CAMPOS_PATH.read_text())
    except Exception:
        pass
    return {}


def _salvar_campos_ghl(campos: dict):
    try:
        CAMPOS_PATH.write_text(json.dumps(campos))
    except Exception:
        pass


@app.route("/api/configurar-campos-ghl", methods=["POST"])
@login_required
def configurar_campos_ghl():
    """Cria campos personalizados no CRM e salva os IDs localmente."""
    data = request.get_json(force=True) or {}
    token = (data.get("token") or "").strip()
    location_id = (data.get("location_id") or "").strip()

    if not token or not location_id:
        return jsonify({"erro": "Token e Location ID obrigatórios"}), 400

    # Remove "Bearer " se vier com prefixo
    if token.lower().startswith("bearer "):
        token = token[7:].strip()

    pais = (data.get("pais") or "BR").strip().upper()
    campos_a_criar = CAMPOS_ENRIQUECIMENTO_US if pais == "US" else CAMPOS_ENRIQUECIMENTO_BR

    campos_salvos = _carregar_campos_ghl()
    criados = []
    ja_existiam = []
    erros = []

    # Busca campos existentes para não duplicar
    try:
        r_list = requests.get(
            f"{GHL_BASE}/locations/{location_id}/customFields",
            headers=_ghl_headers(token),
            timeout=10
        )
        existentes = {}
        if r_list.status_code == 200:
            for f in r_list.json().get("customFields", []):
                existentes[f.get("name", "").strip()] = f.get("id")
    except Exception:
        existentes = {}

    for campo in campos_a_criar:
        nome = campo["name"]
        # Já existe no CRM - salva o ID
        if nome in existentes:
            campos_salvos[nome] = existentes[nome]
            ja_existiam.append(nome)
            continue
        # Tenta criar - experimenta os dois parâmetros que o GHL pode aceitar
        criado = False
        for body_extra in [{"model": "opportunity"}, {"objectKey": "opportunity"}, {}]:
            try:
                body = {"name": nome, "dataType": campo["dataType"], "position": 0}
                body.update(body_extra)
                r = requests.post(
                    f"{GHL_BASE}/locations/{location_id}/customFields",
                    headers=_ghl_headers(token),
                    json=body,
                    timeout=10
                )
                if r.status_code in (200, 201):
                    resp = r.json()
                    field_id = (resp.get("customField") or {}).get("id") or resp.get("id")
                    obj_key = (resp.get("customField") or {}).get("objectKey") or resp.get("objectKey", "contact")
                    if field_id:
                        campos_salvos[nome] = field_id
                        tipo = "oportunidade" if "opportun" in str(obj_key).lower() else "contato"
                        criados.append(f"{nome} (tipo: {tipo})")
                        criado = True
                    break
            except Exception as e:
                erros.append(f"{nome}: {str(e)[:60]}")
                break
        if not criado and nome not in campos_salvos:
            try:
                msg = r.json().get("message", r.text[:80])
            except Exception:
                msg = "sem resposta"
            erros.append(f"{nome}: {msg}")

    _salvar_campos_ghl(campos_salvos)

    return jsonify({
        "criados": criados,
        "ja_existiam": ja_existiam,
        "erros": erros,
        "campos": campos_salvos,
    })


@app.route("/api/campos-ghl")
@login_required
def listar_campos_ghl():
    """Retorna os campos salvos localmente."""
    return jsonify(_carregar_campos_ghl())


# ── WEBHOOK DINASTIA API ──────────────────────────────────────

import logging
_dinastia_log = logging.getLogger("dinastia_webhook")

# Credenciais GHL/Linkia da PX3 Lab (fonte: /opt/mia/workspace/clientes/px3lab/config_ghl.py)
sys.path.insert(0, "/opt/mia/workspace/clientes/px3lab")
try:
    from config_ghl import PX3_TOKEN as _PX3_GHL_TOKEN, PX3_LOCATION_ID as _PX3_GHL_LOCATION
except Exception as _e:
    _dinastia_log.error(f"Falha ao importar config_ghl PX3: {_e}")
    _PX3_GHL_TOKEN = ""
    _PX3_GHL_LOCATION = ""

_GHL_BASE = "https://services.leadconnectorhq.com"

_DINASTIA_BASE = "https://dinastiapi.agentesclimb.us"
_DINASTIA_INSTANCIAS_PATH = Path("/opt/mia/workspace/clientes/px3lab/dinastia_webhook/instancias.json")


def _dinastia_carregar_instancia(nome: str) -> dict | None:
    """Lê config da instância pelo nome. Retorna dict ou None se não encontrada."""
    try:
        with open(_DINASTIA_INSTANCIAS_PATH) as f:
            instancias = json.load(f)
        return instancias.get(nome)
    except Exception as e:
        _dinastia_log.error(f"Erro ao carregar instâncias: {e}")
        return None


def _ghl_headers(ghl_token: str) -> dict:
    return {
        "Authorization": f"Bearer {ghl_token}",
        "Version": "2021-07-28",
        "Content-Type": "application/json",
        "Accept": "application/json",
    }

# OpenAI key para transcrição Whisper
def _get_openai_key() -> str:
    key = os.environ.get("OPENAI_API_KEY", "")
    if not key:
        try:
            with open("/opt/mia-bot/.env") as _f:
                for _line in _f:
                    if _line.startswith("OPENAI_API_KEY="):
                        key = _line.strip().split("=", 1)[1]
                        break
        except Exception:
            pass
    return key


def _dinastia_extrair_evento(payload: dict) -> dict | None:
    """Normaliza o payload da Dinastia. O envelope real vem como data.data.event,
    mas aceita também data.event para retrocompatibilidade."""
    if not isinstance(payload, dict):
        return None
    # Caminho real observado: payload["data"]["data"]["event"]
    data = payload.get("data") or {}
    inner = data.get("data") if isinstance(data, dict) else None
    if isinstance(inner, dict) and "event" in inner:
        return inner.get("event")
    # Caminho alternativo: payload["data"]["event"]
    if isinstance(data, dict) and "event" in data:
        return data.get("event")
    # Caminho direto
    if "event" in payload and isinstance(payload["event"], dict):
        return payload["event"]
    return None


def _dinastia_extrair_texto(message: dict) -> str:
    """Extrai texto da mensagem em diferentes formatos."""
    if not isinstance(message, dict):
        return ""
    # conversation simples
    if isinstance(message.get("conversation"), str) and message["conversation"].strip():
        return message["conversation"].strip()
    # extendedTextMessage
    ext = message.get("extendedTextMessage") or {}
    if isinstance(ext, dict) and isinstance(ext.get("text"), str):
        return ext["text"].strip()
    return ""


def _dinastia_normalizar_telefone(chat: str) -> str:
    """Recebe '5511999999999@s.whatsapp.net' ou '5511999999999:70@s.whatsapp.net' e devolve '+5511999999999'."""
    if not isinstance(chat, str) or not chat:
        return ""
    # Remove domínio (@s.whatsapp.net etc)
    fone = chat.split("@", 1)[0]
    # Remove sufixo de sessão de dispositivo (:XX) antes de limpar
    fone = fone.split(":")[0]
    fone = re.sub(r"\D", "", fone)
    if not fone:
        return ""
    return "+" + fone


def _dinastia_baixar_midia(endpoint: str, msg_obj: dict, dinastia_token: str = "") -> bytes | None:
    """Baixa mídia (audio/image) da Dinastia API. Devolve bytes ou None."""
    try:
        body = {
            "Url": msg_obj.get("url") or msg_obj.get("Url") or "",
            "DirectPath": msg_obj.get("directPath") or msg_obj.get("DirectPath") or "",
            "MediaKey": msg_obj.get("mediaKey") or msg_obj.get("MediaKey") or "",
            "Mimetype": msg_obj.get("mimetype") or msg_obj.get("Mimetype") or "",
            "FileEncSHA256": msg_obj.get("fileEncSha256") or msg_obj.get("FileEncSHA256") or "",
            "FileSHA256": msg_obj.get("fileSha256") or msg_obj.get("FileSHA256") or "",
            "FileLength": msg_obj.get("fileLength") or msg_obj.get("FileLength") or 0,
        }
        r = requests.post(
            f"{_DINASTIA_BASE}/{endpoint.lstrip('/')}",
            headers={"token": dinastia_token, "Content-Type": "application/json"},
            json=body,
            timeout=30,
        )
        if r.status_code != 200:
            _dinastia_log.warning(f"Dinastia {endpoint} {r.status_code}: {r.text[:200]}")
            return None
        # Resposta pode ser JSON com base64 ou binary direto
        ct = r.headers.get("Content-Type", "")
        if "application/json" in ct:
            data = r.json()
            # Tenta campos comuns de base64
            b64 = data.get("data") or data.get("file") or data.get("audio") or data.get("image")
            if b64 and isinstance(b64, str):
                import base64
                return base64.b64decode(b64)
            return None
        return r.content if r.content else None
    except Exception as e:
        _dinastia_log.error(f"Erro baixar mídia {endpoint}: {e}")
        return None


def _dinastia_transcrever_audio(audio_bytes: bytes, mimetype: str) -> str:
    """Transcreve áudio com OpenAI Whisper. Devolve texto ou string vazia."""
    import tempfile
    ext = ".ogg"
    if mimetype:
        if "mp4" in mimetype or "mpeg" in mimetype:
            ext = ".mp4"
        elif "wav" in mimetype:
            ext = ".wav"
        elif "webm" in mimetype:
            ext = ".webm"
        elif "ogg" in mimetype:
            ext = ".ogg"
    try:
        from openai import OpenAI
        client = OpenAI(api_key=_get_openai_key())
        with tempfile.NamedTemporaryFile(suffix=ext, delete=False) as tmp:
            tmp.write(audio_bytes)
            tmp_path = tmp.name
        with open(tmp_path, "rb") as audio_file:
            result = client.audio.transcriptions.create(model="whisper-1", file=audio_file)
        try:
            os.unlink(tmp_path)
        except Exception:
            pass
        return result.text or ""
    except Exception as e:
        _dinastia_log.error(f"Erro Whisper: {e}")
        return ""


def _ghl_buscar_contato(telefone: str, ghl_token: str, ghl_location: str) -> str | None:
    """Procura contato no Linkia pelo telefone. Retorna contactId ou None."""
    try:
        url = f"{_GHL_BASE}/contacts/"
        params = {"locationId": ghl_location, "query": telefone}
        r = requests.get(url, headers=_ghl_headers(ghl_token), params=params, timeout=10)
        if r.status_code != 200:
            _dinastia_log.warning(f"GHL busca {r.status_code}: {r.text[:200]}")
            return None
        body = r.json() if r.content else {}
        contatos = body.get("contacts") or []
        return contatos[0].get("id") if contatos else None
    except Exception as e:
        _dinastia_log.error(f"Erro buscar contato GHL: {e}")
        return None


def _ghl_criar_contato(telefone: str, ghl_token: str, ghl_location: str) -> str | None:
    """Cria contato no Linkia com o telefone. Retorna contactId ou None."""
    try:
        url = f"{_GHL_BASE}/contacts/"
        body = {"phone": telefone, "locationId": ghl_location, "source": "WhatsApp Dinastia"}
        r = requests.post(url, headers=_ghl_headers(ghl_token), json=body, timeout=10)
        if r.status_code in (200, 201):
            data = r.json() if r.content else {}
            contato = data.get("contact") or data
            cid = contato.get("id") if isinstance(contato, dict) else None
            if cid:
                return cid
        # Em duplicidade o GHL devolve 400 com meta.contactId
        if r.status_code == 400:
            try:
                meta = r.json().get("meta") or {}
                cid = meta.get("contactId")
                if cid:
                    return cid
            except Exception:
                pass
        _dinastia_log.warning(f"GHL criar {r.status_code}: {r.text[:200]}")
        return None
    except Exception as e:
        _dinastia_log.error(f"Erro criar contato GHL: {e}")
        return None


def _ghl_adicionar_nota(contact_id: str, texto: str, ghl_token: str) -> bool:
    """Adiciona nota no contato."""
    try:
        url = f"{_GHL_BASE}/contacts/{contact_id}/notes"
        body = {"body": f"📱 WhatsApp: {texto}", "userId": ""}
        r = requests.post(url, headers=_ghl_headers(ghl_token), json=body, timeout=10)
        if r.status_code in (200, 201):
            return True
        _dinastia_log.warning(f"GHL nota {r.status_code}: {r.text[:200]}")
        return False
    except Exception as e:
        _dinastia_log.error(f"Erro nota GHL: {e}")
        return False


def _ghl_buscar_ou_criar_conversa(contact_id: str, ghl_token: str, ghl_location: str) -> str | None:
    """Busca ou cria conversa no GHL para o contato. Retorna conversationId."""
    try:
        # Buscar conversa existente
        url = f"{_GHL_BASE}/conversations/search"
        params = {"locationId": ghl_location, "contactId": contact_id}
        r = requests.get(url, headers=_ghl_headers(ghl_token), params=params, timeout=10)
        if r.status_code == 200:
            body = r.json() if r.content else {}
            convs = body.get("conversations") or []
            if convs:
                return convs[0].get("id")
        # Criar nova conversa
        url = f"{_GHL_BASE}/conversations/"
        body = {"contactId": contact_id, "locationId": ghl_location}
        r = requests.post(url, headers=_ghl_headers(ghl_token), json=body, timeout=10)
        if r.status_code in (200, 201):
            data = r.json() if r.content else {}
            conv = data.get("conversation") or data
            return conv.get("id") if isinstance(conv, dict) else None
        _dinastia_log.warning(f"GHL conversa {r.status_code}: {r.text[:200]}")
        return None
    except Exception as e:
        _dinastia_log.error(f"Erro conversa GHL: {e}")
        return None


def _ghl_postar_mensagem_conversa(conversation_id: str, contact_id: str, texto: str, instancia: str, ghl_token: str) -> bool:
    """Posta mensagem numa conversa do GHL como inbound (lado esquerdo no Linkia)."""
    try:
        # /conversations/messages/inbound garante direction=inbound (lado esquerdo no Linkia)
        url = f"{_GHL_BASE}/conversations/messages/inbound"
        body = {
            "type": "WhatsApp",
            "conversationId": conversation_id,
            "contactId": contact_id,
            "message": texto,
        }
        r = requests.post(url, headers=_ghl_headers(ghl_token), json=body, timeout=10)
        if r.status_code in (200, 201):
            return True
        _dinastia_log.warning(f"GHL msg inbound status={r.status_code}: {r.text[:200]}")
        # Fallback: endpoint padrão com direction explícito
        url_fallback = f"{_GHL_BASE}/conversations/messages"
        body_fallback = {"type": "WhatsApp", "conversationId": conversation_id,
                         "contactId": contact_id, "message": texto, "direction": "inbound"}
        r2 = requests.post(url_fallback, headers=_ghl_headers(ghl_token), json=body_fallback, timeout=10)
        if r2.status_code in (200, 201):
            return True
        _dinastia_log.warning(f"GHL msg fallback status={r2.status_code}: {r2.text[:200]}")
        return False
    except Exception as e:
        _dinastia_log.error(f"Erro postar mensagem GHL: {e}")
        return False


def _ghl_adicionar_tag_instancia(contact_id: str, instancia: str, ghl_token: str) -> None:
    """Adiciona tag 'via-<instancia>' ao contato."""
    try:
        url = f"{_GHL_BASE}/contacts/{contact_id}"
        tag = f"via-{instancia}"
        r = requests.put(url, headers=_ghl_headers(ghl_token), json={"tags": [tag]}, timeout=10)
        if r.status_code not in (200, 201):
            _dinastia_log.warning(f"GHL tag {r.status_code}: {r.text[:100]}")
    except Exception as e:
        _dinastia_log.error(f"Erro tag GHL: {e}")


def _webhook_dinastia_processar(nome_instancia: str):
    """Lógica principal do webhook — usada por todas as rotas de instância."""
    status = {"processed": False, "reason": "", "contact_id": None, "instancia": nome_instancia}
    try:
        # Carregar config da instância
        cfg = _dinastia_carregar_instancia(nome_instancia)
        if not cfg:
            _dinastia_log.warning(f"Instância desconhecida: {nome_instancia}")
            status["reason"] = f"instancia_nao_configurada:{nome_instancia}"
            return jsonify({"ok": True, **status}), 200

        ghl_token = cfg.get("ghl_token") or _PX3_GHL_TOKEN
        ghl_location = cfg.get("ghl_location") or _PX3_GHL_LOCATION
        dinastia_token = cfg.get("dinastia_token") or ""

        payload = request.get_json(force=True, silent=True) or {}
        event_type = payload.get("event") or payload.get("type") or "unknown"
        _dinastia_log.info(f"[{nome_instancia}] Evento: {event_type} | {str(payload)[:200]}")

        # Log bruto
        log_path = Path(f"/opt/mia/workspace/clientes/px3lab/dinastia_webhook/eventos_{nome_instancia}.jsonl")
        log_path.parent.mkdir(parents=True, exist_ok=True)
        with open(log_path, "a") as f:
            f.write(json.dumps({"ts": datetime.now().isoformat(), "event": event_type, "data": payload}) + "\n")

        evento = _dinastia_extrair_evento(payload)
        if not isinstance(evento, dict):
            status["reason"] = "sem_event"
            return jsonify({"ok": True, "received": event_type, **status}), 200

        info = evento.get("Info") or {}
        message = evento.get("Message") or {}

        if info.get("IsFromMe") is True:
            status["reason"] = "ignorado_is_from_me"
            return jsonify({"ok": True, "received": event_type, **status}), 200
        if info.get("IsGroup") is True:
            status["reason"] = "ignorado_grupo"
            return jsonify({"ok": True, "received": event_type, **status}), 200
        msg_type = (info.get("Type") or "").lower()
        if msg_type not in ("text", "audio", "image"):
            status["reason"] = f"ignorado_tipo_{msg_type}"
            return jsonify({"ok": True, "received": event_type, **status}), 200

        # Extrair telefone — LID fix
        chat_raw = info.get("Chat") or ""
        if chat_raw.endswith("@lid"):
            chat_raw = info.get("SenderAlt") or chat_raw
        telefone = _dinastia_normalizar_telefone(chat_raw)

        if not telefone:
            status["reason"] = "sem_telefone"
            return jsonify({"ok": True, "received": event_type, **status}), 200

        # Filtro de números ignorados (ex: familiares, contatos pessoais)
        numeros_ignorados = cfg.get("numeros_ignorados") or []
        fone_limpo = re.sub(r"\D", "", telefone)
        if any(re.sub(r"\D", "", n) == fone_limpo for n in numeros_ignorados):
            status["reason"] = f"ignorado_numero_pessoal:{telefone}"
            _dinastia_log.info(f"[{nome_instancia}] Número ignorado: {telefone}")
            return jsonify({"ok": True, "received": event_type, **status}), 200

        # Montar nota conforme o tipo
        nota_texto = ""
        if msg_type == "text":
            nota_texto = _dinastia_extrair_texto(message)
            if not nota_texto:
                status["reason"] = "sem_texto"
                return jsonify({"ok": True, "received": event_type, **status}), 200

        elif msg_type == "audio":
            audio_obj = (message.get("audioMessage") or message.get("pttMessage")
                         or message.get("AudioMessage") or {})
            mimetype = audio_obj.get("mimetype") or "audio/ogg"
            audio_bytes = _dinastia_baixar_midia("chat/downloadaudio", audio_obj, dinastia_token)
            if audio_bytes:
                transcricao = _dinastia_transcrever_audio(audio_bytes, mimetype)
                nota_texto = f"🎤 Áudio (transcrição): {transcricao}" if transcricao else "🎤 Áudio recebido (transcrição indisponível)"
            else:
                nota_texto = "🎤 Áudio recebido (download falhou)"

        elif msg_type == "image":
            img_obj = message.get("imageMessage") or message.get("ImageMessage") or {}
            caption = img_obj.get("caption") or img_obj.get("Caption") or ""
            nota_texto = "🖼️ Imagem recebida" + (f": {caption}" if caption else "")

        # Buscar ou criar contato
        contact_id = _ghl_buscar_contato(telefone, ghl_token, ghl_location)
        if not contact_id:
            contact_id = _ghl_criar_contato(telefone, ghl_token, ghl_location)
        if not contact_id:
            status["reason"] = "falha_contato_ghl"
            return jsonify({"ok": True, "received": event_type, **status}), 200

        # Verificar se é chat interno entre membros da equipe
        numeros_equipe = []
        for inst_cfg in _dinastia_carregar_instancias().values():
            numeros_equipe += [re.sub(r"\D", "", n) for n in (inst_cfg.get("numeros_equipe") or [])]
        fone_limpo_sender = re.sub(r"\D", "", telefone)
        is_chat_interno = fone_limpo_sender in numeros_equipe

        # Tag por instância (só para leads reais)
        if not is_chat_interno:
            _ghl_adicionar_tag_instancia(contact_id, nome_instancia, ghl_token)

        # Buscar ou criar conversa
        conversation_id = _ghl_buscar_ou_criar_conversa(contact_id, ghl_token, ghl_location)

        ok_msg = False
        ok_nota = False
        if is_chat_interno:
            # Chat interno entre vendedores → nota interna, não conversa de lead
            ok_nota = _ghl_adicionar_nota(contact_id, f"💬 [CHAT INTERNO via {nome_instancia.upper()}]: {nota_texto}", ghl_token)
            status["reason"] = "chat_interno_nota"
        elif conversation_id:
            msg_para_conversa = f"📩 [{nome_instancia.upper()} ← LEAD]: {nota_texto}"
            ok_msg = _ghl_postar_mensagem_conversa(conversation_id, contact_id, msg_para_conversa, nome_instancia, ghl_token)
            if not ok_msg:
                ok_nota = _ghl_adicionar_nota(contact_id, nota_texto, ghl_token)
        else:
            ok_nota = _ghl_adicionar_nota(contact_id, nota_texto, ghl_token)

        status.update({
            "processed": True,
            "contact_id": contact_id,
            "conversation_id": conversation_id,
            "msg_ok": ok_msg,
            "nota_ok": ok_nota,
            "telefone": telefone,
            "tipo": msg_type,
        })
        _dinastia_log.info(f"[{nome_instancia}] Sincronizado | tipo={msg_type} | tel={telefone} | contact={contact_id} | conv={conversation_id} | msg_ok={ok_msg}")

        # SDR webhook direto: se a instância tem sdr_webhook_url configurada,
        # dispara sem depender do workflow do Linkia (contorna race conditions
        # e limitações dos triggers "Cliente Respondeu" / "Contato Criado").
        sdr_url = (cfg.get("sdr_webhook_url") or "").strip()
        if sdr_url and not is_chat_interno:
            try:
                sdr_payload = {
                    "contact_id": contact_id,
                    "conversation_id": conversation_id,
                    "phone": telefone,
                    "message": nota_texto,
                    "body": nota_texto,
                    "direction": "inbound",
                    "type": "WhatsApp",
                    "location_id": ghl_location,
                    "source": f"dinastia:{nome_instancia}",
                }
                r = requests.post(sdr_url, json=sdr_payload, timeout=10)
                _dinastia_log.info(
                    f"[{nome_instancia}] SDR webhook direto {r.status_code}: {r.text[:120]}"
                )
                status["sdr_webhook_status"] = r.status_code
            except Exception as e:
                _dinastia_log.exception(f"[{nome_instancia}] Erro SDR webhook direto: {e}")
                status["sdr_webhook_error"] = str(e)

        return jsonify({"ok": True, "received": event_type, **status}), 200

    except Exception as e:
        _dinastia_log.exception(f"[{nome_instancia}] Erro webhook: {e}")
        return jsonify({"ok": False, "erro": str(e)}), 200


@app.route("/webhook/dinastia/<nome>", methods=["POST"])
def webhook_dinastia_instancia(nome):
    """Webhook multi-instância — /webhook/dinastia/<nome_instancia>"""
    if not _webhook_autorizado():
        return jsonify({"ok": False, "erro": "nao autorizado"}), 401
    return _webhook_dinastia_processar(nome)


@app.route("/webhook/dinastia", methods=["POST"])
def webhook_dinastia():
    """Alias legado — redireciona para instância 'px3'."""
    if not _webhook_autorizado():
        return jsonify({"ok": False, "erro": "nao autorizado"}), 401
    return _webhook_dinastia_processar("px3")


# ── ENVIO BIDIRECIONAL: GHL (Linkia) → Dinastia (WhatsApp) ────────────

def _dinastia_carregar_instancias() -> dict:
    """Retorna o dict completo de instâncias do JSON."""
    try:
        with open(_DINASTIA_INSTANCIAS_PATH) as f:
            return json.load(f) or {}
    except Exception as e:
        _dinastia_log.error(f"Erro ao carregar instâncias (all): {e}")
        return {}


def _dinastia_enviar_mensagem(telefone: str, mensagem: str, dinastia_token: str) -> bool:
    """Envia mensagem de texto via Dinastia API.

    Header de auth é `token: <dinastia_token>` (não Bearer).
    Endpoint observado: POST {_DINASTIA_BASE}/chat/send/text
    Body esperado: {"Phone": "...", "Body": "...", "NumberCheck": true}
    """
    try:
        url = f"{_DINASTIA_BASE}/chat/send/text"
        headers = {"token": dinastia_token, "Content-Type": "application/json"}
        body = {"Phone": telefone, "Body": mensagem, "NumberCheck": True}
        r = requests.post(url, headers=headers, json=body, timeout=20)
        if r.status_code in (200, 201):
            try:
                data = r.json()
                if isinstance(data, dict) and data.get("success") is False:
                    _dinastia_log.warning(f"Dinastia envio success=false: {str(data)[:200]}")
                    return False
            except Exception:
                pass
            return True
        _dinastia_log.warning(f"Dinastia envio status={r.status_code}: {r.text[:200]}")
        return False
    except Exception as e:
        _dinastia_log.error(f"Erro envio Dinastia: {e}")
        return False


@app.route("/webhook/ghl-custom-provider", methods=["GET", "POST"])
def webhook_ghl_custom_provider():
    """Endpoint de verificação/healthcheck do Custom Provider GHL (papiWebhook)."""
    return jsonify({"status": "ok"}), 200


@app.route("/webhook/ghl-outbound", methods=["POST"])
def webhook_ghl_outbound():
    """Recebe resposta do atendente no Linkia (GHL Custom Provider outgoingWebhook)
    e envia via Dinastia API pela instância correta (identificada pela tag via-XXX)."""
    if not _webhook_autorizado():
        return jsonify({"ok": False, "erro": "nao autorizado"}), 401
    try:
        payload = request.get_json(force=True, silent=True) or {}
        _dinastia_log.info(f"[GHL-OUTBOUND] payload: {str(payload)[:500]}")

        contact_id = payload.get("contactId") or ""
        message = payload.get("message") or ""
        conversation_id = payload.get("conversationId") or ""

        if not contact_id or not message:
            return jsonify({"ok": False, "erro": "contactId ou message ausente"}), 400

        ghl_token = _PX3_GHL_TOKEN

        # Buscar contato no GHL para pegar telefone e tags
        url = f"{_GHL_BASE}/contacts/{contact_id}"
        r = requests.get(url, headers=_ghl_headers(ghl_token), timeout=10)
        if r.status_code != 200:
            _dinastia_log.warning(f"[GHL-OUTBOUND] contato não encontrado {r.status_code}: {r.text[:200]}")
            return jsonify({"ok": False, "erro": f"contato não encontrado: {r.status_code}"}), 200

        data = r.json() or {}
        contato = data.get("contact") or data
        telefone = contato.get("phone") or ""
        tags = contato.get("tags") or []

        if not telefone:
            return jsonify({"ok": False, "erro": "contato sem telefone"}), 200

        # Identificar instância pela tag via-XXX
        instancia = "px3"  # fallback padrão
        for tag in tags:
            if isinstance(tag, str) and tag.startswith("via-"):
                instancia = tag.replace("via-", "", 1).strip().lower()
                break

        # Buscar token da instância
        instancias = _dinastia_carregar_instancias()
        cfg = instancias.get(instancia) or instancias.get("px3") or {}
        dinastia_token = cfg.get("dinastia_token") or ""

        if not dinastia_token:
            _dinastia_log.warning(f"[GHL-OUTBOUND] sem token para instancia={instancia}")
            return jsonify({"ok": False, "erro": f"sem token para instância {instancia}"}), 200

        # Enviar via Dinastia API
        fone_limpo = re.sub(r"\D", "", telefone)
        ok = _dinastia_enviar_mensagem(fone_limpo, message, dinastia_token)

        _dinastia_log.info(f"[GHL-OUTBOUND] instancia={instancia} tel={telefone} conv={conversation_id} ok={ok}")
        return jsonify({"ok": ok, "instancia": instancia, "telefone": telefone}), 200

    except Exception as e:
        _dinastia_log.exception(f"[GHL-OUTBOUND] erro: {e}")
        return jsonify({"ok": False, "erro": str(e)}), 200


# ── POLLING LINKIA → DINASTIA (respostas de vendedores) ──────────────────
_LINKIA_PIT_TOKEN = "pit-51341ae4-f642-4dc8-a8a8-d89613dde205"
# Cache em memória dos IDs já processados (populado no boot via SQL).
# Persistência agora é em users.db.linkia_processed (Fase 1 ClimbLeads).
_linkia_processed_ids: set = set()


def _linkia_marcar_processado(msg_id: str, conta_id=None):
    """Persiste um message_id como processado. Idempotente."""
    if not msg_id:
        return
    try:
        conn = _users_db_conn()
        cur = conn.cursor()
        cur.execute(
            "INSERT OR IGNORE INTO linkia_processed (conta_id, message_id, processado_em) VALUES (?,?,?)",
            (int(conta_id) if conta_id else None, msg_id, datetime.now().isoformat())
        )
        conn.commit()
        conn.close()
    except Exception as e:
        print(f"[linkia_processed] erro ao gravar msg_id={msg_id}: {e}")


def _linkia_carregar_processados() -> set:
    """Carrega TODOS os IDs processados (loop é global — SDR Maria roda pro Linkia PX3)."""
    try:
        conn = _users_db_conn()
        cur = conn.cursor()
        rows = cur.execute("SELECT message_id FROM linkia_processed").fetchall()
        conn.close()
        return {r[0] for r in rows}
    except Exception as e:
        print(f"[linkia_processed] erro ao carregar: {e}")
        return set()


def _linkia_rotacionar(limite: int = 5000):
    """Mantém apenas os últimos `limite` IDs (equivalente ao TTL do JSON antigo)."""
    try:
        conn = _users_db_conn()
        cur = conn.cursor()
        cur.execute(
            "DELETE FROM linkia_processed WHERE id NOT IN "
            "(SELECT id FROM linkia_processed ORDER BY id DESC LIMIT ?)",
            (int(limite),)
        )
        conn.commit()
        conn.close()
    except Exception as e:
        print(f"[linkia_processed] erro ao rotacionar: {e}")


def _linkia_poll_once():
    """Verifica conversas recentes no Linkia e encaminha respostas de vendedores via Dinastia."""
    global _linkia_processed_ids
    ghl_location = _PX3_GHL_LOCATION
    token = _LINKIA_PIT_TOKEN
    since_ms = int((datetime.now(timezone.utc) - timedelta(minutes=2)).timestamp() * 1000)

    try:
        r = requests.get(
            f"{_GHL_BASE}/conversations/search",
            headers=_ghl_headers(token),
            params={"locationId": ghl_location, "lastMessageType": "TYPE_WHATSAPP",
                    "lastMessageDirection": "outbound", "limit": 50},
            timeout=15,
        )
        if r.status_code != 200:
            return
        convs = r.json().get("conversations", [])
    except Exception as e:
        _dinastia_log.error(f"[LINKIA-POLL] erro ao buscar conversas: {e}")
        return

    for conv in convs:
        if (conv.get("lastMessageDate") or 0) < since_ms:
            continue
        conv_id = conv.get("id") or ""
        contact_id = conv.get("contactId") or ""
        if not conv_id or not contact_id:
            continue

        try:
            mr = requests.get(
                f"{_GHL_BASE}/conversations/{conv_id}/messages",
                headers=_ghl_headers(token),
                params={"limit": 10},
                timeout=10,
            )
            if mr.status_code != 200:
                continue
            raw = mr.json().get("messages", {})
            msgs = raw.get("messages", []) if isinstance(raw, dict) else (raw or [])
        except Exception:
            continue

        for msg in msgs:
            msg_id = msg.get("id") or ""
            if not msg_id or msg_id in _linkia_processed_ids:
                continue

            direction = msg.get("direction") or ""
            msg_type = msg.get("messageType") or ""
            body = (msg.get("body") or "").strip()
            date_added = msg.get("dateAdded") or ""

            _linkia_processed_ids.add(msg_id)
            _linkia_marcar_processado(msg_id)

            # Só mensagens outbound WhatsApp que NÃO são injeções nossas
            if direction != "outbound" or msg_type != "TYPE_WHATSAPP":
                continue
            if body.startswith("📩 ["):
                continue  # mensagem de lead que nós postamos
            if not body:
                continue

            # Verificar se é recente (últimos 2 minutos)
            try:
                msg_dt = datetime.fromisoformat(date_added.replace("Z", "+00:00"))
                if (datetime.now(timezone.utc) - msg_dt).total_seconds() > 120:
                    continue
            except Exception:
                continue

            # Buscar telefone e tags do contato
            try:
                cr = requests.get(
                    f"{_GHL_BASE}/contacts/{contact_id}",
                    headers=_ghl_headers(_PX3_GHL_TOKEN),
                    timeout=10,
                )
                if cr.status_code != 200:
                    continue
                contato = cr.json().get("contact") or cr.json()
                telefone = contato.get("phone") or ""
                tags = contato.get("tags") or []
            except Exception:
                continue

            if not telefone:
                continue

            # Identificar instância pela tag via-XXX
            instancia = "px3"
            for tag in tags:
                if isinstance(tag, str) and tag.startswith("via-"):
                    instancia = tag.replace("via-", "", 1).lower()
                    break

            instancias = _dinastia_carregar_instancias()
            cfg = instancias.get(instancia) or instancias.get("px3") or {}
            dinastia_token = cfg.get("dinastia_token") or ""
            if not dinastia_token:
                _dinastia_log.warning(f"[LINKIA-POLL] sem token para instancia={instancia}")
                continue

            fone_limpo = re.sub(r"\D", "", telefone)
            ok = _dinastia_enviar_mensagem(fone_limpo, body, dinastia_token)
            _dinastia_log.info(
                f"[LINKIA-POLL] repassado via={instancia} tel={fone_limpo} ok={ok} msg={body[:60]}"
            )

    # Rotaciona cache em memória (mantém últimos 5000) e purga o DB.
    if len(_linkia_processed_ids) > 10000:
        _linkia_processed_ids = set(list(_linkia_processed_ids)[-5000:])
        _linkia_rotacionar(5000)


def _linkia_poll_loop():
    """Thread daemon: roda _linkia_poll_once() a cada 30 segundos."""
    _dinastia_log.info("[LINKIA-POLL] Thread iniciada (intervalo 30s).")
    # Boot: hidrata cache em memória com histórico persistido no DB.
    try:
        _linkia_processed_ids.update(_linkia_carregar_processados())
        _dinastia_log.info(f"[LINKIA-POLL] Hidratado do DB: {len(_linkia_processed_ids)} IDs.")
    except Exception as e:
        _dinastia_log.error(f"[LINKIA-POLL] erro hidratando cache: {e}")
    while True:
        try:
            _linkia_poll_once()
        except Exception as e:
            _dinastia_log.error(f"[LINKIA-POLL] erro no ciclo: {e}")
        time.sleep(30)


# ═════════════════════════════════════════════════════════════════════════
# WEBHOOK SDR MARIA — GHL/Linkia → Claude → GHL
# ═════════════════════════════════════════════════════════════════════════
# Fluxo:
# 1) GHL Workflow "Cliente Respondeu" dispara POST neste endpoint
# 2) Extrai contactId/conversationId/mensagem
# 3) Busca histórico da conversa via API do GHL
# 4) Chama Claude API (modelo claude-sonnet-4-6) com o system prompt da Maria
# 5) Posta resposta como outbound na conversa GHL (aparece do lado direito)
# 6) Move stage/aplica tags se detectar palavras-chave

_SDR_MARIA_PROMPT_PATH = Path("/opt/mia/workspace/clientes/px3lab/sdr_alma_cloud_photorf.md")
_SDR_MARIA_PIPELINE_ID = "iV13JEVbMgfcad9Xw6gz"
_SDR_MARIA_STAGES = {
    "novo_contato":      "d06d8954-e19d-4962-9b86-c25da4278de7",
    "em_conversa":       "1efeab5b-8f4f-46ef-be3d-29ec42366833",
    "qualificando":      "d189c3e2-b321-4288-8ba7-011acabf62e9",
    "qualificado":       "bfca2016-a09c-42e5-86ef-a78836b67b48",
    "em_negociacao":     "70733ae5-a5d0-4ec5-8235-97859b8ef10f",
    "aguardando":        "8f95ea16-18d5-45b3-8575-c965e0e583cf",
    "encaminhado":       "faf155d0-a694-435c-9015-ccc974f72167",
    "demo_agendada":     "1dddadaa-6c50-423d-a2fa-f2487e0a6baa",
    "ganhamos":          "72587990-1a07-4196-9582-117deebf2ae3",
    "perdemos":          "07ab9b2a-e5d4-4313-9c69-6a593fb7b193",
}

_sdr_maria_system_prompt_cache = {"content": "", "loaded_at": 0}

# Buffer de concatenação de mensagens: quando o lead manda várias mensagens em
# sequência rápida, aguardamos BUFFER_DELAY segundos após a última antes de
# processar tudo junto — evita responder mensagem por mensagem.
# Estrutura: contact_id -> {"msgs": [str], "conv_id": str, "contact_id": str,
#                           "ghl_token": str, "timer": threading.Timer | None}
_sdr_maria_buffer: dict = {}
_sdr_maria_buffer_lock = threading.Lock()
BUFFER_DELAY = 4.0  # segundos


def _sdr_maria_enviar_typing(conversation_id: str, ghl_token: str) -> None:
    """Envia typing indicator ao GHL para o lead ver 'digitando...' no WhatsApp.
    Tenta primeiro o endpoint PUT em /conversations/{id}; se não suportado,
    tenta POST em /conversations/messages com type=ACTIVITY. Nunca bloqueia
    o fluxo — se falhar, apenas loga e segue."""
    if not conversation_id or not ghl_token:
        return
    # Tentativa 1: PUT /conversations/{id}
    try:
        url = f"{_GHL_BASE}/conversations/{conversation_id}"
        body = {"typing": True, "typing_timeout": 5}
        r = requests.put(url, headers=_ghl_headers(ghl_token), json=body, timeout=8)
        if r.status_code in (200, 201, 204):
            return
        if r.status_code not in (404, 405):
            _dinastia_log.info(f"[SDR-MARIA] typing PUT status={r.status_code}: {r.text[:200]}")
    except Exception as e:
        _dinastia_log.warning(f"[SDR-MARIA] typing PUT erro: {e}")

    # Tentativa 2: POST /conversations/messages com type=ACTIVITY
    try:
        url = f"{_GHL_BASE}/conversations/messages"
        body = {
            "type": "ACTIVITY",
            "conversationId": conversation_id,
            "body": "typing",
        }
        r = requests.post(url, headers=_ghl_headers(ghl_token), json=body, timeout=8)
        if r.status_code in (200, 201, 204):
            return
        _dinastia_log.info(f"[SDR-MARIA] typing não suportado: status={r.status_code}")
    except Exception as e:
        _dinastia_log.warning(f"[SDR-MARIA] typing ACTIVITY erro: {e}")


def _sdr_maria_processar_buffer(contact_id: str) -> None:
    """Processa mensagens acumuladas no buffer para um contact_id.
    Roda em thread separada (chamado por threading.Timer), fora do contexto Flask.
    Faz: concatena mensagens -> busca histórico -> typing -> Claude -> post -> ações."""
    try:
        # 1) Retira dados do buffer
        with _sdr_maria_buffer_lock:
            dados = _sdr_maria_buffer.pop(contact_id, None)
        if not dados:
            _dinastia_log.info(f"[SDR-MARIA-BUFFER] nada no buffer para contact={contact_id}")
            return

        msgs = dados.get("msgs") or []
        conversation_id = dados.get("conv_id") or ""
        ghl_token = dados.get("ghl_token") or ""

        if not msgs:
            _dinastia_log.info(f"[SDR-MARIA-BUFFER] msgs vazio contact={contact_id}")
            return

        # Concatena com " / " se forem várias, senão só a mensagem única
        if len(msgs) == 1:
            mensagem = msgs[0]
        else:
            mensagem = " / ".join(m.strip() for m in msgs if m and m.strip())

        _dinastia_log.info(f"[SDR-MARIA-BUFFER] processando contact={contact_id} "
                           f"msgs={len(msgs)} concatenado={mensagem[:200]}")

        if not ghl_token:
            _dinastia_log.error(f"[SDR-MARIA-BUFFER] sem ghl_token contact={contact_id}")
            return

        # 2) Se não veio conversation_id, tenta buscar/criar
        if not conversation_id:
            conversation_id = _ghl_buscar_ou_criar_conversa(contact_id, ghl_token, _PX3_GHL_LOCATION) or ""

        # 3) Busca histórico
        historico = _sdr_maria_buscar_historico(conversation_id, ghl_token, limite=20) if conversation_id else []

        # 4) Envia typing indicator (não bloqueia se falhar)
        if conversation_id:
            _sdr_maria_enviar_typing(conversation_id, ghl_token)

        # 5) Chama Claude
        resposta = _sdr_maria_gerar_resposta(historico, mensagem)
        if not resposta:
            _dinastia_log.warning(f"[SDR-MARIA-BUFFER] Claude retornou vazio contact={contact_id}")
            return

        # 6) Posta resposta no GHL (outbound)
        posted = False
        if conversation_id:
            posted = _sdr_maria_postar_outbound(conversation_id, contact_id, resposta, ghl_token)

        # 7) Detecta ações automáticas na resposta
        acoes = _sdr_maria_detectar_acoes(resposta)
        stage_moved = False
        if acoes["stage"]:
            opp_id = _sdr_maria_buscar_opportunity(contact_id, ghl_token)
            if opp_id:
                stage_moved = _sdr_maria_mover_stage(opp_id, acoes["stage"], ghl_token)
        if acoes["tags"]:
            _sdr_maria_adicionar_tags(contact_id, acoes["tags"], ghl_token)

        _dinastia_log.info(f"[SDR-MARIA-BUFFER] ok contact={contact_id} conv={conversation_id} "
                           f"posted={posted} stage_moved={stage_moved} tags={acoes['tags']}")

    except Exception as e:
        _dinastia_log.exception(f"[SDR-MARIA-BUFFER] erro contact={contact_id}: {e}")


def _sdr_maria_system_prompt() -> str:
    """Carrega o prompt da Maria em cache (reload a cada 5 min)."""
    now = time.time()
    if _sdr_maria_system_prompt_cache["content"] and (now - _sdr_maria_system_prompt_cache["loaded_at"] < 300):
        return _sdr_maria_system_prompt_cache["content"]
    try:
        with open(_SDR_MARIA_PROMPT_PATH, "r", encoding="utf-8") as f:
            content = f.read()
        _sdr_maria_system_prompt_cache["content"] = content
        _sdr_maria_system_prompt_cache["loaded_at"] = now
        return content
    except Exception as e:
        _dinastia_log.error(f"[SDR-MARIA] erro lendo prompt: {e}")
        return ""


def _sdr_maria_buscar_historico(conversation_id: str, ghl_token: str, limite: int = 20) -> list[dict]:
    """Busca as últimas N mensagens da conversa no GHL. Retorna lista [{role, content}]."""
    if not conversation_id:
        return []
    try:
        url = f"{_GHL_BASE}/conversations/{conversation_id}/messages"
        params = {"limit": limite}
        r = requests.get(url, headers=_ghl_headers(ghl_token), params=params, timeout=10)
        if r.status_code != 200:
            _dinastia_log.warning(f"[SDR-MARIA] histórico status={r.status_code}: {r.text[:200]}")
            return []
        data = r.json() or {}
        # O GHL retorna geralmente em {"messages": {"messages": [...]}} ou {"messages": [...]}
        msgs_raw = data.get("messages")
        if isinstance(msgs_raw, dict):
            msgs = msgs_raw.get("messages") or []
        elif isinstance(msgs_raw, list):
            msgs = msgs_raw
        else:
            msgs = []

        # GHL retorna do mais recente pro mais antigo — invertemos
        msgs = list(reversed(msgs))

        historico = []
        for m in msgs:
            body = (m.get("body") or m.get("message") or "").strip()
            if not body:
                continue
            direction = (m.get("direction") or "").lower()
            # inbound = lead falando (user), outbound = Maria/atendente (assistant)
            role = "user" if direction == "inbound" else "assistant"
            historico.append({"role": role, "content": body})

        # Consolida mensagens consecutivas do mesmo papel (Claude exige alternância)
        consolidado = []
        for item in historico:
            if consolidado and consolidado[-1]["role"] == item["role"]:
                consolidado[-1]["content"] += "\n" + item["content"]
            else:
                consolidado.append(item)

        # Claude exige que a primeira mensagem seja do user
        while consolidado and consolidado[0]["role"] != "user":
            consolidado.pop(0)

        return consolidado
    except Exception as e:
        _dinastia_log.error(f"[SDR-MARIA] erro histórico: {e}")
        return []


def _sdr_maria_gerar_resposta(historico: list[dict], mensagem_atual: str) -> str:
    """Chama Claude API e retorna resposta da Maria."""
    try:
        import anthropic
    except ImportError:
        _dinastia_log.error("[SDR-MARIA] anthropic SDK não instalado")
        return ""

    api_key = os.environ.get("ANTHROPIC_API_KEY", "")
    if not api_key:
        # Tenta ler dos .env conhecidos (mia-bot tem a chave)
        for env_path in ("/opt/mia-bot/.env", "/opt/mia/.env"):
            try:
                with open(env_path, "r") as f:
                    for _line in f:
                        _line = _line.strip()
                        if _line.startswith("ANTHROPIC_API_KEY="):
                            api_key = _line.split("=", 1)[1].strip().strip('"').strip("'")
                            break
                if api_key:
                    break
            except Exception:
                continue
    if not api_key:
        _dinastia_log.error("[SDR-MARIA] ANTHROPIC_API_KEY ausente")
        return ""

    system_prompt = _sdr_maria_system_prompt()
    if not system_prompt:
        return ""

    # Se a mensagem atual não está no histórico como último user, adiciona
    messages = list(historico)
    if not messages or messages[-1].get("role") != "user" or mensagem_atual not in (messages[-1].get("content") or ""):
        if messages and messages[-1].get("role") == "user":
            messages[-1]["content"] += "\n" + mensagem_atual
        else:
            messages.append({"role": "user", "content": mensagem_atual})

    # Garante que começa com user
    while messages and messages[0]["role"] != "user":
        messages.pop(0)

    if not messages:
        messages = [{"role": "user", "content": mensagem_atual}]

    try:
        client = anthropic.Anthropic(api_key=api_key)
        response = client.messages.create(
            model="claude-opus-4-7",
            max_tokens=1024,
            system=system_prompt,
            messages=messages,
        )
        texto = ""
        for block in response.content:
            if getattr(block, "type", None) == "text":
                texto += block.text
        return texto.strip()
    except Exception as e:
        _dinastia_log.error(f"[SDR-MARIA] erro Claude: {e}")
        return ""


def _sdr_maria_postar_outbound(conversation_id: str, contact_id: str, texto: str, ghl_token: str) -> bool:
    """Posta mensagem como outbound (lado direito no Linkia — resposta da Maria)."""
    try:
        url = f"{_GHL_BASE}/conversations/messages"
        body = {
            "type": "WhatsApp",
            "conversationId": conversation_id,
            "contactId": contact_id,
            "message": texto,
        }
        r = requests.post(url, headers=_ghl_headers(ghl_token), json=body, timeout=15)
        if r.status_code in (200, 201):
            return True
        _dinastia_log.warning(f"[SDR-MARIA] post outbound status={r.status_code}: {r.text[:300]}")
        return False
    except Exception as e:
        _dinastia_log.error(f"[SDR-MARIA] erro post outbound: {e}")
        return False


def _sdr_maria_buscar_opportunity(contact_id: str, ghl_token: str) -> str | None:
    """Busca opportunity ativa do contato no pipeline SDR Maria."""
    try:
        url = f"{_GHL_BASE}/opportunities/search"
        params = {
            "location_id": _PX3_GHL_LOCATION,
            "contact_id": contact_id,
            "pipeline_id": _SDR_MARIA_PIPELINE_ID,
        }
        r = requests.get(url, headers=_ghl_headers(ghl_token), params=params, timeout=10)
        if r.status_code != 200:
            return None
        data = r.json() or {}
        opps = data.get("opportunities") or []
        if opps:
            return opps[0].get("id")
        return None
    except Exception as e:
        _dinastia_log.error(f"[SDR-MARIA] erro busca opp: {e}")
        return None


def _sdr_maria_mover_stage(opportunity_id: str, stage_id: str, ghl_token: str) -> bool:
    """Move a opportunity para outro stage."""
    if not opportunity_id or not stage_id:
        return False
    try:
        url = f"{_GHL_BASE}/opportunities/{opportunity_id}"
        body = {"pipelineStageId": stage_id}
        r = requests.put(url, headers=_ghl_headers(ghl_token), json=body, timeout=10)
        if r.status_code in (200, 201):
            return True
        _dinastia_log.warning(f"[SDR-MARIA] mover stage status={r.status_code}: {r.text[:200]}")
        return False
    except Exception as e:
        _dinastia_log.error(f"[SDR-MARIA] erro mover stage: {e}")
        return False


def _sdr_maria_adicionar_tags(contact_id: str, tags: list[str], ghl_token: str) -> None:
    """Adiciona tags ao contato (mantém as existentes via endpoint específico)."""
    if not contact_id or not tags:
        return
    try:
        url = f"{_GHL_BASE}/contacts/{contact_id}/tags"
        r = requests.post(url, headers=_ghl_headers(ghl_token), json={"tags": tags}, timeout=10)
        if r.status_code not in (200, 201):
            _dinastia_log.warning(f"[SDR-MARIA] add tags status={r.status_code}: {r.text[:200]}")
    except Exception as e:
        _dinastia_log.error(f"[SDR-MARIA] erro add tags: {e}")


_SDR_MARIA_OPTOUT_KEYWORDS = (
    "bloquear", "bloqueia", "bloqueado", "bloqueada",
    "parar", "para de", "pare de", "não quero mais", "nao quero mais",
    "me tire", "me remova", "remover", "tirar da lista",
    "não me mande", "nao me mande", "não me envie", "nao me envie",
    "não me contacte", "nao me contacte", "não me contate", "nao me contate",
    "chega", "obrigado não", "obrigado nao", "obrigado mas não", "obrigado mas nao",
    "descadastrar", "cancelar", "não tenho interesse", "nao tenho interesse",
    "vai tomar", "me deixa em paz", "some", "não perturbe", "nao perturbe",
)


def _sdr_maria_detectar_optout(mensagem: str) -> bool:
    """Retorna True se a mensagem do lead indicar intenção de bloqueio/opt-out."""
    if not mensagem:
        return False
    txt = mensagem.lower()
    return any(k in txt for k in _SDR_MARIA_OPTOUT_KEYWORDS)


def _sdr_maria_detectar_acoes(resposta: str) -> dict:
    """Detecta indicações na resposta da Maria para mover stage/tags automaticamente."""
    txt = (resposta or "").lower()
    acoes = {"stage": None, "tags": []}

    if any(k in txt for k in ["calendly", "agendei", "agendamento confirmado", "demo agendada",
                              "agendada para", "reservei o horário"]):
        acoes["stage"] = _SDR_MARIA_STAGES["demo_agendada"]
        acoes["tags"].append("demo-agendada")
    elif any(k in txt for k in ["encaminhei para", "vou passar para o consultor",
                                "consultor especializado", "monique vai te chamar"]):
        acoes["stage"] = _SDR_MARIA_STAGES["encaminhado"]
        acoes["tags"].append("encaminhado-consultor")
    elif any(k in txt for k in ["link de pagamento", "pagamento", "vou te enviar a proposta",
                                "enviar a proposta"]):
        acoes["stage"] = _SDR_MARIA_STAGES["em_negociacao"]
    elif any(k in txt for k in ["qualificado", "faz sentido pra vocês"]):
        acoes["stage"] = _SDR_MARIA_STAGES["qualificado"]

    return acoes


@app.route("/webhook/ghl-sdr-maria", methods=["POST"])
def webhook_ghl_sdr_maria():
    """Recebe webhook do GHL quando lead responde no WhatsApp.
    Bufferiza as mensagens por BUFFER_DELAY segundos e delega o processamento
    (Claude + post + ações) para _sdr_maria_processar_buffer em thread separada.
    Opt-out é tratado imediatamente, sem passar pelo buffer."""
    if not _webhook_autorizado():
        return jsonify({"ok": False, "erro": "nao autorizado"}), 401
    payload = request.get_json(force=True, silent=True) or {}
    _dinastia_log.info(f"[SDR-MARIA] payload: {str(payload)[:600]}")

    try:
        # Extração defensiva — GHL varia bastante o payload
        contact = payload.get("contact") or {}
        conversation = payload.get("conversation") or {}
        message = payload.get("message") or {}

        contact_id = (contact.get("id") or payload.get("contactId") or
                      payload.get("contact_id") or "")
        conversation_id = (conversation.get("id") or payload.get("conversationId") or
                           payload.get("conversation_id") or "")
        mensagem = (message.get("body") or payload.get("message_body") or
                    payload.get("body") or payload.get("text") or "")

        if not contact_id or not mensagem:
            _dinastia_log.warning(f"[SDR-MARIA] payload incompleto: contact_id={contact_id} msg={bool(mensagem)}")
            return jsonify({"ok": True, "skip": "payload_incompleto"}), 200

        ghl_token = _PX3_GHL_TOKEN
        if not ghl_token:
            _dinastia_log.error("[SDR-MARIA] _PX3_GHL_TOKEN ausente")
            return jsonify({"ok": True, "skip": "sem_token"}), 200

        # Detecta opt-out ANTES de bufferizar — se lead pediu pra parar, processa
        # imediatamente: move oportunidade pro stage Perdemos, aplica tags e encerra.
        if _sdr_maria_detectar_optout(mensagem):
            _dinastia_log.info(f"[SDR-MARIA] OPT-OUT detectado | contact={contact_id} | msg={mensagem[:50]}")
            # Cancela buffer pendente (se houver) — lead pediu pra parar
            with _sdr_maria_buffer_lock:
                pend = _sdr_maria_buffer.pop(contact_id, None)
                if pend and pend.get("timer"):
                    try:
                        pend["timer"].cancel()
                    except Exception:
                        pass
            _sdr_maria_adicionar_tags(contact_id, ["perdido", "opt-out"], ghl_token)
            opp_id = _sdr_maria_buscar_opportunity(contact_id, ghl_token)
            stage_moved = False
            if opp_id:
                stage_moved = _sdr_maria_mover_stage(opp_id, _SDR_MARIA_STAGES["perdemos"], ghl_token)
            return jsonify({
                "ok": True,
                "optout": True,
                "posted": False,
                "stage_moved": stage_moved,
                "tags_added": ["perdido", "opt-out"],
            }), 200

        # Bufferiza: adiciona mensagem, cancela timer anterior e agenda novo.
        with _sdr_maria_buffer_lock:
            slot = _sdr_maria_buffer.get(contact_id)
            if slot is None:
                slot = {
                    "msgs": [],
                    "conv_id": conversation_id,
                    "contact_id": contact_id,
                    "ghl_token": ghl_token,
                    "timer": None,
                }
                _sdr_maria_buffer[contact_id] = slot
            else:
                # Atualiza conv_id se ainda não tinha e agora veio
                if not slot.get("conv_id") and conversation_id:
                    slot["conv_id"] = conversation_id
                # Refresh token (caso tenha rotacionado)
                slot["ghl_token"] = ghl_token

            slot["msgs"].append(mensagem)

            # Cancela timer anterior (se estava agendado)
            old_timer = slot.get("timer")
            if old_timer is not None:
                try:
                    old_timer.cancel()
                except Exception:
                    pass

            # Agenda novo timer
            new_timer = threading.Timer(BUFFER_DELAY, _sdr_maria_processar_buffer, args=[contact_id])
            new_timer.daemon = True
            slot["timer"] = new_timer
            total_msgs = len(slot["msgs"])

        new_timer.start()

        _dinastia_log.info(f"[SDR-MARIA-BUFFER] bufferizado contact={contact_id} "
                           f"total_msgs={total_msgs} delay={BUFFER_DELAY}s")

        return jsonify({"ok": True, "buffered": True, "queued_msgs": total_msgs}), 200

    except Exception as e:
        _dinastia_log.exception(f"[SDR-MARIA] erro: {e}")
        # Sempre retorna 200 pro GHL não retentar
        return jsonify({"ok": True, "erro": str(e)}), 200


# Inicia thread de polling ao carregar o módulo (funciona com gunicorn e python direto)
threading.Thread(target=_linkia_poll_loop, daemon=True, name="linkia-poll").start()


# ══════════════════════════════════════════════════════════════
# ANALISADOR DE LEADS (super_admin only)
# ══════════════════════════════════════════════════════════════

from flask import send_file
import analisador as _analisador

_analisador.ensure_table()


# Contas com acesso liberado ao Analisador de Leads (além de super_admin).
# 3 = Climb/Digital+
_CONTAS_COM_ANALISE = {3}


def pode_analisar_leads(user) -> bool:
    """True se o usuário pode acessar a feature Analisar Lead.

    Libera para: super_admin OU usuários pertencentes às contas em _CONTAS_COM_ANALISE.
    Cobre também o caso de super_admin impersonando: current_user vira o user impersonado,
    então o gate segue a conta_id do impersonado (correto — se a conta impersonada não
    tem acesso, o botão some).
    """
    if not user or not getattr(user, "is_authenticated", False):
        return False
    if getattr(user, "plano", "") == "super_admin":
        return True
    try:
        cid = int(getattr(user, "conta_id", 0) or 0)
    except (TypeError, ValueError):
        cid = 0
    return cid in _CONTAS_COM_ANALISE


def _bloqueia_sem_acesso_analise():
    """Retorna response 403 se o usuário logado não puder acessar o Analisador de Leads."""
    if not current_user.is_authenticated:
        return jsonify({"ok": False, "erro": "não autenticado"}), 401
    if pode_analisar_leads(current_user):
        return None
    return jsonify({"ok": False, "erro": "acesso restrito"}), 403


@app.route("/analisar-lead", methods=["POST"])
@login_required
def analisar_lead():
    blk = _bloqueia_sem_acesso_analise()
    if blk is not None:
        return blk
    data = request.get_json(silent=True) or {}
    nome = (data.get("nome") or "").strip()
    if not nome:
        return jsonify({"ok": False, "erro": "nome obrigatório"}), 400

    lead_data = {
        "nome": nome,
        "telefone": (data.get("telefone") or "").strip(),
        "cidade": (data.get("cidade") or "").strip(),
        "categoria": (data.get("categoria") or "").strip(),
        "instagram": (data.get("instagram") or "").strip(),
        "site_atual": (data.get("site_atual") or data.get("site") or "").strip(),
    }
    aid = _analisador.criar_analise(lead_data, int(current_user.id))

    threading.Thread(
        target=_analisador.run_analise_lead,
        args=(aid, lead_data),
        daemon=True,
        name=f"analise-{aid}"
    ).start()

    return jsonify({"ok": True, "analise_id": aid})


@app.route("/analisar-lead/<int:aid>/status")
@login_required
def analisar_lead_status(aid):
    blk = _bloqueia_sem_acesso_analise()
    if blk is not None:
        return blk
    row = _analisador.get_analise(aid)
    if not row:
        return jsonify({"ok": False, "erro": "análise não encontrada"}), 404
    # só o próprio dono ou super_admin vê
    if row.get("user_id") != int(current_user.id) and current_user.plano != "super_admin":
        return jsonify({"ok": False, "erro": "acesso negado"}), 403
    return jsonify({
        "ok": True,
        "id": row["id"],
        "status": row["status"],
        "site_url": row.get("site_url"),
        "pdf_disponivel": bool(row.get("pdf_path")),
        "erro": row.get("erro"),
        "criado_em": row.get("criado_em"),
        "concluido_em": row.get("concluido_em"),
        "lead_nome": row.get("lead_nome"),
    })


@app.route("/analisar-lead/<int:aid>/pdf")
@login_required
def analisar_lead_pdf(aid):
    blk = _bloqueia_sem_acesso_analise()
    if blk is not None:
        return blk
    row = _analisador.get_analise(aid)
    if not row or not row.get("pdf_path"):
        return jsonify({"ok": False, "erro": "PDF indisponível"}), 404
    pdf = Path(row["pdf_path"])
    if not pdf.exists():
        return jsonify({"ok": False, "erro": "arquivo não existe"}), 404
    nome_arq = f"dossie_{(row.get('lead_nome') or 'lead').replace(' ','_')}.pdf"
    return send_file(str(pdf), as_attachment=True, download_name=nome_arq, mimetype="application/pdf")


@app.route("/analises")
@login_required
def listar_analises_rota():
    blk = _bloqueia_sem_acesso_analise()
    if blk is not None:
        return blk
    rows = _analisador.listar_analises(int(current_user.id))
    return jsonify({
        "ok": True,
        "analises": [
            {
                "id": r["id"],
                "lead_nome": r["lead_nome"],
                "lead_cidade": r["lead_cidade"],
                "lead_categoria": r["lead_categoria"],
                "status": r["status"],
                "site_url": r.get("site_url"),
                "pdf_disponivel": bool(r.get("pdf_path")),
                "criado_em": r.get("criado_em"),
                "concluido_em": r.get("concluido_em"),
            }
            for r in rows
        ]
    })


# ══════════════════════════════════════════════════════════════
# API PROGRAMÁTICA (Bearer token) - uso por agentes externos
# ══════════════════════════════════════════════════════════════

# Garante que a migração idempotente das colunas api_token / api_allowed_ips
# rode em produção (gunicorn não executa o bloco __main__).
try:
    init_db()
except Exception as _e:
    print(f"[api] Falha em init_db no import: {_e}")


# Rate limiting em memória (por token) — 30 req/min, sem dependência externa.
from collections import defaultdict as _defaultdict_api
import time as _time_api
import threading as _threading_api

_rate_limits: dict = _defaultdict_api(list)
_rate_limits_lock = _threading_api.Lock()
_RATE_LIMIT_PER_MINUTE = 30


def _check_rate_limit(token: str) -> bool:
    """Retorna True se dentro do limite, False se excedeu."""
    now = _time_api.time()
    with _rate_limits_lock:
        timestamps = [t for t in _rate_limits[token] if now - t < 60]
        if len(timestamps) >= _RATE_LIMIT_PER_MINUTE:
            _rate_limits[token] = timestamps
            return False
        timestamps.append(now)
        _rate_limits[token] = timestamps
    return True


def _auth_token():
    """Retorna (User, token_str) autenticado pelo header Authorization: Bearer <token>,
    ou (None, None). Também valida IP whitelist se configurada."""
    auth = request.headers.get("Authorization", "")
    if not auth.startswith("Bearer "):
        return None, None
    token = auth[7:].strip()
    user = buscar_por_token(token)
    if not user:
        return None, None
    # Checar IP whitelist (JSON array em user.api_allowed_ips)
    allowed_raw = getattr(user, "api_allowed_ips", None)
    if allowed_raw:
        try:
            allowed = _json_api.loads(allowed_raw)
            client_ip = (request.headers.get("X-Forwarded-For") or request.remote_addr or "").split(",")[0].strip()
            if allowed and client_ip not in allowed:
                return None, None
        except Exception:
            pass
    return user, token


def require_token(f):
    """Decorator para endpoints da API programática. Injeta 'user' como primeiro arg."""
    @wraps(f)
    def decorated(*args, **kwargs):
        user, token = _auth_token()
        if not user:
            return jsonify({"ok": False, "error": "unauthorized"}), 401
        if not _check_rate_limit(token):
            return jsonify({"ok": False, "error": "rate_limit_exceeded"}), 429
        return f(user, *args, **kwargs)
    return decorated


def _job_visivel_para(user, job) -> bool:
    """Isolamento por conta: user só vê job da própria conta (super_admin vê tudo)."""
    if not job:
        return False
    if getattr(user, "is_super_admin", False):
        return True
    dono = job.get("conta_id")
    if dono is None:
        return True
    return dono == user.conta_id


# ── POST /api/v1/prospectar ──
@app.route("/api/v1/prospectar", methods=["POST"])
@require_token
def api_v1_prospectar(user):
    data = request.get_json(silent=True) or {}
    nicho = (data.get("nicho") or "").strip()
    cidade = (data.get("cidade") or "").strip()
    limite = int(data.get("limite") or 10)
    pais = (data.get("pais") or "BR").strip().upper()

    if not nicho or not cidade:
        return jsonify({"ok": False, "error": "nicho e cidade sao obrigatorios"}), 400
    if not user.ativo:
        return jsonify({"ok": False, "error": "usuario inativo"}), 403

    limite = max(1, min(200, limite))

    # Quota mensal da conta (mesma logica do endpoint web)
    conta_id = user.conta_id
    if conta_id:
        try:
            resetar_quota_se_novo_mes(conta_id)
            quota = get_quota(conta_id)
            if quota["usados"] >= quota["limite"]:
                return jsonify({
                    "ok": False,
                    "error": f"quota mensal atingida ({quota['limite']} leads/mes)"
                }), 403
        except Exception:
            pass

    # Cria job usando helper existente, mas fora do contexto de current_user.
    # Precisamos setar conta_id/user_id manualmente porque _novo_job usa _dono_atual().
    job_id = str(uuid.uuid4())[:8]
    jobs[job_id] = {
        "id": job_id,
        "tipo": "prospeccao",
        "nicho": nicho,
        "cidade": cidade,
        "limite": limite,
        "conta_id": conta_id,
        "user_id": int(user.id),
        "status": "aguardando",
        "progresso": 0,
        "total": limite,
        "leads_enviados": 0,
        "leads": [],
        "logs": [],
        "erro": None,
        "criado_em": datetime.now().isoformat(),
        "concluido_em": None,
    }
    if len(jobs) > MAX_JOBS_HISTORY:
        ids_ordenados = sorted(jobs.keys(), key=lambda k: jobs[k]["criado_em"])
        for old_id in ids_ordenados[:-MAX_JOBS_HISTORY]:
            del jobs[old_id]
    try:
        salvar_job_historico(jobs[job_id])
    except Exception as e:
        print(f"[api_v1_prospectar] erro histórico {job_id}: {e}")

    # Modo API sempre roda em modo "coleta" (sem envio CRM automático).
    t = threading.Thread(
        target=run_prospector,
        args=(job_id, nicho, cidade, limite,
              "",           # token CRM
              "", "", "",   # location/pipeline/stage
              "ghl",        # crm_type default
              "",           # connection_id
              pais,
              False,        # enriquecimento_auto
              conta_id,
              False,        # filtro_sem_site
              "", "", 0,    # filtro_nota, filtro_site_gmb, score_minimo
        ),
        daemon=True,
        name=f"api-prospec-{job_id}",
    )
    t.start()

    return jsonify({"ok": True, "job_id": job_id})


# ── GET /api/v1/jobs/<job_id>/status ──
@app.route("/api/v1/jobs/<job_id>/status", methods=["GET"])
@require_token
def api_v1_job_status(user, job_id):
    job = jobs.get(job_id)
    if not job:
        # Fallback histórico
        try:
            historico = listar_jobs_historico(limit=50)
            job_hist = next((j for j in historico if j["id"] == job_id), None)
        except Exception:
            job_hist = None
        if not job_hist:
            return jsonify({"ok": False, "error": "job_not_found"}), 404
        # Checar acesso pelo conta_id do histórico
        if not user.is_super_admin and job_hist.get("conta_id") not in (None, user.conta_id):
            return jsonify({"ok": False, "error": "job_not_found"}), 404
        return jsonify({
            "ok": True,
            "job_id": job_hist["id"],
            "status": job_hist["status"],
            "progresso": job_hist["progresso"],
            "total": job_hist["total"],
            "leads_enviados": job_hist["leads_enviados"],
            "criado_em": job_hist.get("criado_em"),
            "concluido_em": job_hist.get("concluido_em"),
        })
    if not _job_visivel_para(user, job):
        return jsonify({"ok": False, "error": "job_not_found"}), 404
    return jsonify({
        "ok": True,
        "job_id": job["id"],
        "status": job["status"],
        "progresso": job["progresso"],
        "total": job["total"],
        "leads_enviados": job["leads_enviados"],
        "erro": job.get("erro"),
        "criado_em": job.get("criado_em"),
        "concluido_em": job.get("concluido_em"),
    })


# ── GET /api/v1/jobs/<job_id>/leads ──
@app.route("/api/v1/jobs/<job_id>/leads", methods=["GET"])
@require_token
def api_v1_job_leads(user, job_id):
    job = jobs.get(job_id)
    if not job:
        return jsonify({"ok": False, "error": "job_not_found"}), 404
    if not _job_visivel_para(user, job):
        return jsonify({"ok": False, "error": "job_not_found"}), 404
    if job["status"] not in ("concluido", "pausado", "cancelado", "erro"):
        return jsonify({
            "ok": False,
            "error": "job_not_finished",
            "status": job["status"],
            "progresso": job["progresso"],
            "total": job["total"],
        }), 409
    return jsonify({
        "ok": True,
        "job_id": job["id"],
        "status": job["status"],
        "total": len(job.get("leads") or []),
        "leads": job.get("leads") or [],
    })


# ── POST /api/v1/analisar ──
@app.route("/api/v1/analisar", methods=["POST"])
@require_token
def api_v1_analisar(user):
    if not pode_analisar_leads(user):
        return jsonify({"ok": False, "error": "forbidden"}), 403
    data = request.get_json(silent=True) or {}
    nome = (data.get("nome") or "").strip()
    if not nome:
        return jsonify({"ok": False, "error": "nome obrigatório"}), 400

    lead_data = {
        "nome": nome,
        "telefone": (data.get("telefone") or "").strip(),
        "cidade": (data.get("cidade") or "").strip(),
        "categoria": (data.get("categoria") or "").strip(),
        "instagram": (data.get("instagram") or "").strip(),
        "site_atual": (data.get("site_atual") or data.get("site") or "").strip(),
    }

    aid = _analisador.criar_analise(lead_data, int(user.id))

    threading.Thread(
        target=_analisador.run_analise_lead,
        args=(aid, lead_data),
        daemon=True,
        name=f"api-analise-{aid}",
    ).start()

    return jsonify({"ok": True, "analise_id": aid})


# ── GET /api/v1/analises/<aid>/status ──
@app.route("/api/v1/analises/<int:aid>/status", methods=["GET"])
@require_token
def api_v1_analise_status(user, aid):
    row = _analisador.get_analise(aid)
    if not row:
        return jsonify({"ok": False, "error": "analise_not_found"}), 404
    # Isolamento: dono da análise ou super_admin
    if row.get("user_id") != int(user.id) and not user.is_super_admin:
        return jsonify({"ok": False, "error": "analise_not_found"}), 404
    return jsonify({
        "ok": True,
        "analise_id": row["id"],
        "status": row["status"],
        "etapa": row.get("status"),
        "lead_nome": row.get("lead_nome"),
        "site_url": row.get("site_url"),
        "pdf_disponivel": bool(row.get("pdf_path")),
        "erro": row.get("erro"),
        "criado_em": row.get("criado_em"),
        "concluido_em": row.get("concluido_em"),
    })


# ── GET /api/v1/analises/<aid>/resultado ──
@app.route("/api/v1/analises/<int:aid>/resultado", methods=["GET"])
@require_token
def api_v1_analise_resultado(user, aid):
    row = _analisador.get_analise(aid)
    if not row:
        return jsonify({"ok": False, "error": "analise_not_found"}), 404
    if row.get("user_id") != int(user.id) and not user.is_super_admin:
        return jsonify({"ok": False, "error": "analise_not_found"}), 404
    if row.get("status") != "concluido":
        return jsonify({
            "ok": False,
            "error": "analise_not_finished",
            "status": row.get("status"),
        }), 409
    # URL do PDF: reaproveita a rota web /analisar-lead/<aid>/pdf (mas exige login web).
    # Como agente externo, expomos um endpoint alternativo autenticado por token.
    pdf_url = None
    if row.get("pdf_path"):
        pdf_url = f"/api/v1/analises/{aid}/pdf"
    return jsonify({
        "ok": True,
        "analise_id": row["id"],
        "lead_nome": row.get("lead_nome"),
        "lp_url": row.get("site_url"),
        "pdf_url": pdf_url,
    })


# ── GET /api/v1/analises/<aid>/pdf ── (download do PDF via token)
@app.route("/api/v1/analises/<int:aid>/pdf", methods=["GET"])
@require_token
def api_v1_analise_pdf(user, aid):
    row = _analisador.get_analise(aid)
    if not row:
        return jsonify({"ok": False, "error": "analise_not_found"}), 404
    if row.get("user_id") != int(user.id) and not user.is_super_admin:
        return jsonify({"ok": False, "error": "analise_not_found"}), 404
    if not row.get("pdf_path"):
        return jsonify({"ok": False, "error": "pdf_unavailable"}), 404
    pdf = Path(row["pdf_path"])
    if not pdf.exists():
        return jsonify({"ok": False, "error": "pdf_file_missing"}), 404
    nome_arq = f"dossie_{(row.get('lead_nome') or 'lead').replace(' ', '_')}.pdf"
    return send_file(str(pdf), as_attachment=True, download_name=nome_arq, mimetype="application/pdf")


if __name__ == "__main__":
    init_db()
    # Migração automática: cria perfil de CRM a partir do ghl_config para contas
    # que ainda não têm nenhum perfil salvo (idempotente).
    try:
        criados = migrar_ghl_config_para_perfis()
        if criados:
            print(f"[migracao] {criados} perfil(is) de CRM criado(s) a partir de ghl_config.")
    except Exception as e:
        print(f"[migracao] Falha ao migrar ghl_config -> crm_perfis: {e}")
    print("Iniciando Prospeccao Web na porta 5050...")
    app.run(host="0.0.0.0", port=5050, debug=False)
