from __future__ import annotations

import asyncio
import hashlib
import json
import logging
import time
from pathlib import Path
from typing import Any

import httpx
from fastapi import Depends, FastAPI, HTTPException, Request
from fastapi.responses import HTMLResponse, JSONResponse, StreamingResponse
from fastapi.templating import Jinja2Templates

from .logging_setup import configure_logging
from .tracing import setup_tracing, teardown_tracing
from .langfuse_client import setup_langfuse, teardown_langfuse

logger = logging.getLogger(__name__)

from .agent_router import AgentKind, all_agents_summary
from .chatbot_bridge import ChatbotBridge, ChatbotBridgeError
from .config import BASE_DIR, get_settings
from .models import (
    AgentChatRequest,
    AnalyticsTrackRequest,
    AutomationRuleRunRequest,
    AutomationRuleToggleRequest,
    ChatbotConsentCreateRequest,
    CommandRequest,
    EventInboundRequest,
    MemoryImportRequest,
    MemoryPromotion,
    ProposalDecision,
    WhatsAppCampaignCreateRequest,
    WhatsAppCampaignSendRequest,
)
from .metrics_collector import generate_metrics, inc
from .orchestrator import MultiAgentOrchestrator
from .providers import build_provider
from .rate_limiter import get_stats as rl_stats
from .rate_limiter import rate_limit_medium, rate_limit_strict
from .scheduler import AutomationScheduler
from .ssh_exec import SSHCommandRunner
from .storage import SupervisorStore
from .telegram_bot import TelegramBotService
from .whatsapp_client import WhatsAppService
from . import security_scanner
from .rag import rag_store

settings = get_settings()
store = SupervisorStore(settings.db_path)
runner = SSHCommandRunner(settings)
provider = build_provider(settings)
chatbot_bridge = ChatbotBridge(settings)
orchestrator = MultiAgentOrchestrator(store=store, provider=provider, runner=runner, chatbot_bridge=chatbot_bridge)
telegram_bot = TelegramBotService(settings=settings, store=store, orchestrator=orchestrator)
whatsapp = WhatsAppService(
    enabled=settings.whatsapp_enabled,
    phone_number_id=settings.whatsapp_phone_number_id,
    access_token=settings.whatsapp_access_token,
    webhook_verify_token=settings.whatsapp_webhook_verify_token,
    app_secret=settings.whatsapp_app_secret,
    allowed_phones=settings.whatsapp_allowed_phones,
    api_version=settings.whatsapp_api_version,
)
scheduler = AutomationScheduler(store=store, orchestrator=orchestrator)
templates = Jinja2Templates(directory=str(BASE_DIR / "app" / "templates"))

app = FastAPI(title=settings.app_name, version="0.2.0")


@app.middleware("http")
async def request_logging_middleware(request: Request, call_next: Any) -> Any:
    start = time.perf_counter()
    response = await call_next(request)
    elapsed_ms = round((time.perf_counter() - start) * 1000)
    path = request.url.path
    # No loguear el SSE stream ni assets para evitar spam
    if not path.startswith("/api/stream"):
        logger.info(
            "http.request",
            extra={
                "method": request.method,
                "path": path,
                "status": response.status_code,
                "ms": elapsed_ms,
            },
        )
    return response


@app.exception_handler(Exception)
async def unhandled_exception_handler(request: Request, exc: Exception) -> JSONResponse:
    logger.exception(
        "http.unhandled_error",
        extra={"path": request.url.path, "error": str(exc)},
    )
    return JSONResponse(
        status_code=500,
        content={"detail": "Error interno del servidor"},
    )


@app.on_event("startup")
async def startup_event() -> None:
    configure_logging()
    setup_tracing(
        service_name=settings.otel_service_name,
        otlp_endpoint=settings.otel_otlp_endpoint,
        enabled=settings.otel_enabled,
    )
    setup_langfuse(
        enabled=settings.langfuse_enabled,
        public_key=settings.langfuse_public_key,
        secret_key=settings.langfuse_secret_key,
        host=settings.langfuse_host,
    )
    logger.info("supervisor.starting", extra={"app": settings.app_name, "provider": provider.name})
    await telegram_bot.start()
    await whatsapp.start()

    # ── Registrar handler de mensajes entrantes de WhatsApp ── #
    async def _wa_message_handler(phone: str, text: str, msg_id: str, raw: dict[str, Any]) -> None:  # noqa: ARG001
        actor = f"wa:{phone}"

        # Botones interactivos: id con formato "approve:xxx" o "reject:xxx"
        if text.startswith("approve:"):
            proposal_id = text.split(":", 1)[1].strip()
            try:
                proposal = await orchestrator.approve_proposal(proposal_id, actor, "Aprobado desde WhatsApp")
                await whatsapp.send_text(phone, f"✅ *Propuesta aprobada*\n{proposal.get('title', '')}\nID: {proposal_id}\nEstado: {proposal.get('status', 'approved')}")
            except KeyError:
                await whatsapp.send_text(phone, f"❌ Propuesta {proposal_id} no encontrada.")
            except Exception as exc:
                await whatsapp.send_text(phone, f"❌ Error al aprobar: {exc}")
            return

        if text.startswith("reject:"):
            proposal_id = text.split(":", 1)[1].strip()
            try:
                proposal = orchestrator.reject_proposal(proposal_id, actor, "Rechazado desde WhatsApp")
                await whatsapp.send_text(phone, f"❌ Propuesta rechazada: {proposal.get('title', '')}\nID: {proposal_id}")
            except KeyError:
                await whatsapp.send_text(phone, f"❌ Propuesta {proposal_id} no encontrada.")
            except Exception as exc:
                await whatsapp.send_text(phone, f"❌ Error al rechazar: {exc}")
            return

        # Comandos de texto compactos
        if text.startswith("/aprobar "):
            proposal_id = text.split(None, 1)[1].strip()
            try:
                proposal = await orchestrator.approve_proposal(proposal_id, actor, "Aprobado desde WhatsApp")
                exec_log = (proposal.get("execution_log") or "")[:800]
                body = f"✅ Aprobado: {proposal.get('title', '')}\nID: {proposal_id}\nEstado: {proposal.get('status', 'approved')}"
                if exec_log:
                    body += f"\n\nSalida:\n{exec_log}"
                await whatsapp.send_text(phone, body)
            except KeyError:
                await whatsapp.send_text(phone, f"❌ Propuesta no encontrada: {proposal_id}")
            except Exception as exc:
                await whatsapp.send_text(phone, f"❌ Error: {exc}")
            return

        if text.startswith("/rechazar "):
            proposal_id = text.split(None, 1)[1].strip()
            try:
                proposal = orchestrator.reject_proposal(proposal_id, actor, "Rechazado desde WhatsApp")
                await whatsapp.send_text(phone, f"❌ Rechazado: {proposal.get('title', '')}\nID: {proposal_id}")
            except KeyError:
                await whatsapp.send_text(phone, f"❌ Propuesta no encontrada: {proposal_id}")
            except Exception as exc:
                await whatsapp.send_text(phone, f"❌ Error: {exc}")
            return

        if text.startswith("/propuestas"):
            proposals = store.list_proposals(status="pending", limit=5)
            if not proposals:
                await whatsapp.send_text(phone, "✅ No hay propuestas pendientes.")
            else:
                lines = [f"📌 {len(proposals)} propuesta(s) pendiente(s):\n"]
                for p in proposals:
                    lines.append(f"• {p['id'][:8]} — {p.get('title', '')}")
                lines.append("\nResponde: /aprobar <id> o /rechazar <id>")
                await whatsapp.send_text(phone, "\n".join(lines))
            return

        # Mensaje libre → crear plan vía orchestrator
        try:
            result = await orchestrator.create_command_plan(text, actor, "whatsapp")
            study = result["study"]
            proposals = result["proposals"]
            summary = str(study.get("summary") or study.get("title") or "Estudio generado")[:300]
            await whatsapp.send_text(phone, f"📊 *Estudio*\n{summary}")
            for proposal in proposals[:3]:  # Máx 3 para no saturar WA
                pid = proposal.get("id", "")
                title = str(proposal.get("title") or "Propuesta")[:60]
                risk = proposal.get("risk", "medium")
                risk_emoji = {"⬇️ low": "🟢", "low": "🟢", "medium": "🟡", "high": "🔴"}.get(risk, "🟡")
                body = f"{risk_emoji} *{title}*\nID: {pid[:8]}\nRiesgo: {risk}\n\n¿Aprobar esta propuesta?"
                await whatsapp.send_interactive_buttons(
                    phone,
                    body_text=body,
                    buttons=[
                        {"id": f"approve:{pid}", "title": "Aprobar"},
                        {"id": f"reject:{pid}", "title": "Rechazar"},
                    ],
                    footer_text="XZonas Ops Supervisor",
                )
        except Exception as exc:
            logger.exception("wa.handler_error", extra={"phone": phone, "error": str(exc)})
            await whatsapp.send_text(phone, f"❌ Error procesando el mensaje: {exc}")

    whatsapp.set_message_handler(_wa_message_handler)
    scheduler.attach_telegram(telegram_bot)
    await scheduler.start()
    # Registrar función de envío de campañas WA en el bot de Telegram
    telegram_bot.set_wa_campaign_sender(_do_wa_campaign_send)
    # ── RAG: inicializar y auto-indexar docs si está habilitado ── #
    if settings.rag_enabled:
        await asyncio.to_thread(
            rag_store.setup,
            settings.qdrant_url,
            settings.rag_collection,
            settings.rag_embed_model,
        )
        if rag_store.status()["ready"]:
            asyncio.create_task(
                rag_store.ingest_directory(),
                name="rag-ingest-docs",
            )
    logger.info("supervisor.ready", extra={"bind": settings.public_base_url})


@app.on_event("shutdown")
async def shutdown_event() -> None:
    logger.info("supervisor.stopping")
    await scheduler.stop()
    await telegram_bot.stop()
    await whatsapp.stop()
    teardown_tracing()
    teardown_langfuse()
    logger.info("supervisor.stopped")


def require_api_token(request: Request) -> str:
    token = request.headers.get("X-Supervisor-Token") or request.query_params.get("token")
    if token is None or token != settings.dashboard_token:
        raise HTTPException(status_code=401, detail="Token de supervisor inválido")
    return token


@app.get("/health")
def health() -> dict[str, Any]:
    counts = store.count_summary()
    sched = scheduler.get_status()
    return {
        "ok": True,
        "app": settings.app_name,
        "provider": provider.name,
        "scheduler": sched,
        "db": {
            "studies": counts.get("studies", 0),
            "proposals": counts.get("proposals", 0),
            "pending": counts.get("pending", 0),
            "memories": counts.get("memories", 0),
        },
        "telegram": {
            "enabled": telegram_bot.enabled,
            "mode": settings.telegram_mode if telegram_bot.enabled else "disabled",
            "allowed_chats": len(settings.telegram_allowed_chat_ids),
        },
        "whatsapp": {
            "enabled": whatsapp.enabled,
            "phone_number_id": bool(settings.whatsapp_phone_number_id),
            "allowed_phones": len(settings.whatsapp_allowed_phones),
            "api_version": settings.whatsapp_api_version,
        },
        "bridge": {
            "configured": bool(settings.chatbot_bridge_url and settings.chatbot_admin_token),
        },
        "rate_limiter": rl_stats(),
    }


@app.get("/", response_class=HTMLResponse)
def dashboard(request: Request) -> HTMLResponse:
    return templates.TemplateResponse(
        request,
        "dashboard.html",
        {
            "app_name": settings.app_name,
            "provider_name": provider.name,
            "public_base_url": settings.public_base_url,
            "telegram_enabled": telegram_bot.enabled,
            "telegram_mode": settings.telegram_mode if telegram_bot.enabled else "disabled",
            "telegram_allowed_chats": len(settings.telegram_allowed_chat_ids),
        },
    )


@app.get("/api/status", dependencies=[Depends(require_api_token)])
def api_status() -> dict[str, Any]:
    from . import langfuse_client as _lf

    langfuse_probe: dict[str, Any] = {
        "host": settings.langfuse_host,
        **_lf.get_summary(),
        "http_ok": False,
    }
    if settings.langfuse_host:
        health_url = settings.langfuse_host.rstrip("/") + "/api/public/health"
        langfuse_probe["health_url"] = health_url
        try:
            resp = httpx.get(health_url, timeout=2.0)
            langfuse_probe["http_ok"] = resp.status_code < 400
            langfuse_probe["http_status"] = resp.status_code
        except Exception as exc:
            langfuse_probe["http_error"] = str(exc)

    return {
        "settings": settings.public_summary(),
        "counts": store.count_summary(),
        "allowed_commands": runner.list_allowed_commands(),
        "scheduler": scheduler.get_status(),
        "langfuse": langfuse_probe,
        "agents": all_agents_summary(),
    }


@app.get("/api/metrics-summary", dependencies=[Depends(require_api_token)])
def api_metrics_summary() -> dict[str, Any]:
    """Resumen numérico de métricas del supervisor para el panel del dashboard."""
    counts = store.count_summary()
    sched = scheduler.get_status()
    from . import langfuse_client as _lf
    return {
        "studies": counts.get("studies", 0),
        "proposals": counts.get("proposals", 0),
        "pending": counts.get("pending", 0),
        "executed": counts.get("executed", 0),
        "rejected": counts.get("rejected", 0),
        "memories": counts.get("memories", 0),
        "automation_rules": counts.get("automation_rules", 0),
        "automation_enabled": counts.get("automation_enabled", 0),
        "scheduler_running": sched.get("running", False),
        "scheduler_executing": sched.get("executing_now", []),
        "rate_limiter": rl_stats(),
        "analytics_total": counts.get("analytics_total", 0),
        "langfuse": _lf.get_summary(),
    }


@app.get("/api/langfuse/status", dependencies=[Depends(require_api_token)])
def api_langfuse_status() -> dict[str, Any]:
    """Estado de Langfuse para el dashboard y diagnóstico."""
    from . import langfuse_client as _lf
    return {
        "host": settings.langfuse_host,
        **_lf.get_summary(),
    }


# ------------------------------------------------------------------ #
# Seguridad — pip-audit + trivy                                        #
# ------------------------------------------------------------------ #

@app.post("/api/security/scan", dependencies=[Depends(require_api_token), Depends(rate_limit_strict)])
async def api_security_scan(request: Request) -> dict[str, Any]:
    """Lanza un escaneo de seguridad completo (pip-audit + trivy si disponible)."""
    body: dict = {}
    try:
        body = await request.json()
    except Exception:
        pass
    actor = body.get("actor", "api")
    result = await security_scanner.run_full_scan(actor=actor)
    store.save_security_scan(result)
    inc("security_scan_total")
    if result.get("critical", 0) > 0 or result.get("high", 0) > 0:
        inc("security_scan_findings_high")
    return result


@app.get("/api/security/scans", dependencies=[Depends(require_api_token)])
def api_security_scans(limit: int = 20) -> list[dict[str, Any]]:
    """Devuelve el historial de escaneos de seguridad."""
    return store.list_security_scans(limit=limit)


@app.get("/api/security/latest", dependencies=[Depends(require_api_token)])
def api_security_latest() -> dict[str, Any]:
    """Devuelve el último escaneo completo con resultados detallados."""
    return store.get_latest_security_scan() or {}


# ------------------------------------------------------------------ #
# SSE — eventos en tiempo real para dashboard                          #
# ------------------------------------------------------------------ #

@app.get("/api/stream")
async def api_stream(request: Request, token: str = "", after: int = 0) -> StreamingResponse:
    if not token or token != settings.dashboard_token:
        raise HTTPException(status_code=401, detail="Token inválido")

    async def event_generator():
        last_id = after
        yield f"data: {json.dumps({'type': 'connected', 'last_id': last_id})}\n\n"
        while not await request.is_disconnected():
            events = store.poll_sse_events(after_id=last_id, limit=10)
            for ev in events:
                last_id = ev["id"]
                yield f"id: {last_id}\ndata: {json.dumps(ev)}\n\n"
            await asyncio.sleep(2)

    return StreamingResponse(
        event_generator(),
        media_type="text/event-stream",
        headers={
            "Cache-Control": "no-cache",
            "X-Accel-Buffering": "no",
        },
    )


# ------------------------------------------------------------------ #
# Eventos inbound (de PHP portal o servicios externos)                 #
# ------------------------------------------------------------------ #

@app.post("/api/events/inbound", dependencies=[Depends(rate_limit_medium)])
async def api_events_inbound(payload: EventInboundRequest, request: Request) -> dict[str, Any]:
    token = request.headers.get("X-Supervisor-Token") or request.query_params.get("token")
    if not token or token != settings.dashboard_token:
        raise HTTPException(status_code=401, detail="Token de supervisor inválido")
    inc("supervisor_inbound_events_total")

    event_id = store.add_inbound_event(
        event_type=payload.event_type,
        source=payload.source,
        payload=payload.payload,
    )
    store.add_audit("inbound_event.received", {
        "event_id": event_id,
        "event_type": payload.event_type,
        "source": payload.source,
    })
    store.push_sse_event("inbound_event", {
        "event_id": event_id,
        "event_type": payload.event_type,
        "source": payload.source,
    })

    study_id: str | None = None
    if payload.auto_study:
        text = f"Evento del portal: {payload.event_type} desde {payload.source}. Datos: {json.dumps(payload.payload)}"
        result = await orchestrator.create_command_plan(text, actor=payload.source, channel="event_inbound")
        study_id = result.get("study", {}).get("id")
        store.mark_inbound_event_processed(event_id, study_id=study_id)

    return {"ok": True, "event_id": event_id, "study_id": study_id}


@app.get("/api/events/inbound", dependencies=[Depends(require_api_token)])
def api_events_list(limit: int = 30, unprocessed: bool = False) -> list[dict[str, Any]]:
    return store.list_inbound_events(limit=limit, unprocessed_only=unprocessed)


# ------------------------------------------------------------------ #
# Analytics — visitas, interacciones y tracking de eventos             #
# ------------------------------------------------------------------ #

@app.post("/api/analytics/track", dependencies=[Depends(rate_limit_medium)])
async def api_analytics_track(payload: AnalyticsTrackRequest, request: Request) -> dict[str, Any]:
    """Recibe eventos de tracking desde el portal PHP o el dashboard JS."""
    raw_ip = request.client.host if request.client else ""
    ip_hash = hashlib.sha256(raw_ip.encode()).hexdigest()[:16] if raw_ip else ""
    event_id = store.add_analytics_event(
        event_type=payload.event_type,
        page=payload.page,
        source=payload.source,
        actor=payload.actor,
        ip_hash=ip_hash,
        meta=payload.meta,
    )
    inc("supervisor_analytics_events_total")
    store.push_sse_event("analytics_event", {
        "event_id": event_id,
        "event_type": payload.event_type,
        "page": payload.page,
        "source": payload.source,
    })
    return {"ok": True, "event_id": event_id}


@app.get("/api/analytics/summary", dependencies=[Depends(require_api_token)])
def api_analytics_summary() -> dict[str, Any]:
    """Resumen de analytics: totales, breakdown por tipo, páginas top y evolución 7 días."""
    return store.count_analytics_summary()


@app.get("/api/analytics/events", dependencies=[Depends(require_api_token)])
def api_analytics_events(limit: int = 50) -> list[dict[str, Any]]:
    """Lista de eventos de analytics más recientes."""
    return store.list_analytics_events(limit=min(limit, 200))


# ------------------------------------------------------------------ #
# Agentes — chat directo con agente especializado                      #
# ------------------------------------------------------------------ #

@app.post("/api/agents/chat", dependencies=[Depends(require_api_token), Depends(rate_limit_strict)])
async def api_agent_chat(payload: AgentChatRequest) -> dict[str, Any]:
    result = await orchestrator.agent_chat(payload.text, payload.agent_kind, payload.actor)
    inc("supervisor_agent_chat_total", {"kind": payload.agent_kind or "general"})
    store.push_sse_event("agent.chat", {
        "agent": result.get("agent"),
        "kind": result.get("kind"),
        "actor": payload.actor,
    })
    return result


@app.get("/api/agents", dependencies=[Depends(require_api_token)])
def api_agents_list() -> list[dict[str, str]]:
    return all_agents_summary()


# ------------------------------------------------------------------ #
# Chatbot bridge                                                        #
# ------------------------------------------------------------------ #

@app.get("/api/chatbot/overview", dependencies=[Depends(require_api_token)])
async def api_chatbot_overview(limit: int = 8) -> dict[str, Any]:
    try:
        return await orchestrator.get_chatbot_overview(limit=limit)
    except (ChatbotBridgeError, RuntimeError) as exc:
        raise HTTPException(status_code=503, detail=str(exc)) from exc


@app.get("/api/chatbot/campaigns", dependencies=[Depends(require_api_token)])
async def api_chatbot_campaigns(limit: int = 12) -> dict[str, Any]:
    try:
        return await orchestrator.list_chatbot_campaigns(limit=limit)
    except (ChatbotBridgeError, RuntimeError) as exc:
        raise HTTPException(status_code=503, detail=str(exc)) from exc


@app.get("/api/chatbot/campaigns/{campaign_id}", dependencies=[Depends(require_api_token)])
async def api_chatbot_campaign_detail(campaign_id: int, source: str | None = None) -> dict[str, Any]:
    try:
        return await orchestrator.get_chatbot_campaign(campaign_id, source)
    except (ChatbotBridgeError, RuntimeError) as exc:
        raise HTTPException(status_code=503, detail=str(exc)) from exc


@app.get("/api/chatbot/consents", dependencies=[Depends(require_api_token)])
async def api_chatbot_consents(
    limit: int = 50,
    status: str | None = None,
    source: str | None = None,
    phone: str | None = None,
    advertiser_id: int | None = None,
) -> dict[str, Any]:
    try:
        return await orchestrator.list_chatbot_consents(limit=limit, status=status, source=source, phone=phone, advertiser_id=advertiser_id)
    except (ChatbotBridgeError, RuntimeError) as exc:
        raise HTTPException(status_code=503, detail=str(exc)) from exc


@app.post("/api/chatbot/consents", dependencies=[Depends(require_api_token)])
async def api_chatbot_create_consent(payload: ChatbotConsentCreateRequest) -> dict[str, Any]:
    try:
        return await orchestrator.create_chatbot_consent(
            actor=payload.actor,
            phone=payload.phone,
            source=payload.source,
            name=payload.name,
            proof=payload.proof,
            advertiser_id=payload.advertiser_id,
        )
    except (ChatbotBridgeError, RuntimeError, ValueError) as exc:
        raise HTTPException(status_code=503, detail=str(exc)) from exc


@app.post("/api/chatbot/campaigns/create", dependencies=[Depends(require_api_token)])
async def api_chatbot_create_campaign(payload: WhatsAppCampaignCreateRequest) -> dict[str, Any]:
    try:
        return await orchestrator.create_chatbot_campaign(
            actor=payload.actor,
            advertiser_id=payload.advertiser_id,
            name=payload.name,
            message_template=payload.message_template,
            description=payload.description,
            scheduled_at=payload.scheduled_at,
            use_ai_personalization=payload.use_ai_personalization,
            targets=[target.model_dump() for target in payload.targets],
        )
    except (ChatbotBridgeError, RuntimeError, ValueError) as exc:
        raise HTTPException(status_code=503, detail=str(exc)) from exc


@app.post("/api/chatbot/campaigns/{campaign_id}/send", dependencies=[Depends(require_api_token)])
async def api_chatbot_send_campaign(campaign_id: int, payload: WhatsAppCampaignSendRequest) -> dict[str, Any]:
    try:
        return await orchestrator.send_chatbot_campaign(
            actor=payload.actor,
            campaign_id=campaign_id,
            immediate=payload.immediate,
            source=payload.source,
        )
    except (ChatbotBridgeError, RuntimeError, ValueError) as exc:
        raise HTTPException(status_code=503, detail=str(exc)) from exc


@app.post("/api/chatbot/marketing/prepare", dependencies=[Depends(require_api_token)])
async def api_chatbot_prepare_marketing(payload: CommandRequest) -> dict[str, Any]:
    try:
        return await orchestrator.create_marketing_plan(payload.text, payload.actor, payload.channel or "dashboard-marketing")
    except (ChatbotBridgeError, RuntimeError) as exc:
        raise HTTPException(status_code=503, detail=str(exc)) from exc


# ------------------------------------------------------------------ #
# Comandos, propuestas, memorias                                        #
# ------------------------------------------------------------------ #

@app.post("/api/commands", dependencies=[Depends(require_api_token), Depends(rate_limit_strict)])
async def api_commands(payload: CommandRequest) -> dict[str, Any]:
    inc("supervisor_commands_total")
    return await orchestrator.create_command_plan(payload.text, payload.actor, payload.channel)


@app.get("/api/studies", dependencies=[Depends(require_api_token)])
def api_studies(limit: int = 50) -> list[dict[str, Any]]:
    return store.list_studies(limit=limit)


@app.get("/api/proposals", dependencies=[Depends(require_api_token)])
def api_proposals(status: str | None = None, limit: int = 100) -> list[dict[str, Any]]:
    return store.list_proposals(status=status, limit=limit)


@app.post("/api/proposals/{proposal_id}/approve", dependencies=[Depends(require_api_token)])
async def api_approve_proposal(proposal_id: str, payload: ProposalDecision) -> dict[str, Any]:
    try:
        result = await orchestrator.approve_proposal(proposal_id, payload.actor, payload.note)
        store.push_sse_event("proposal.approved", {"proposal_id": proposal_id, "actor": payload.actor})
        return result
    except KeyError as exc:
        raise HTTPException(status_code=404, detail=str(exc)) from exc
    except RuntimeError as exc:
        raise HTTPException(status_code=500, detail=str(exc)) from exc


@app.post("/api/proposals/{proposal_id}/reject", dependencies=[Depends(require_api_token)])
def api_reject_proposal(proposal_id: str, payload: ProposalDecision) -> dict[str, Any]:
    try:
        result = orchestrator.reject_proposal(proposal_id, payload.actor, payload.note)
        store.push_sse_event("proposal.rejected", {"proposal_id": proposal_id, "actor": payload.actor})
        return result
    except KeyError as exc:
        raise HTTPException(status_code=404, detail=str(exc)) from exc
    except RuntimeError as exc:
        raise HTTPException(status_code=500, detail=str(exc)) from exc


@app.post("/api/studies/{study_id}/memorize", dependencies=[Depends(require_api_token)])
def api_memorize_study(study_id: str, payload: MemoryPromotion) -> dict[str, Any]:
    try:
        return orchestrator.memorize_study(study_id, payload.actor, payload.note)
    except KeyError as exc:
        raise HTTPException(status_code=404, detail=str(exc)) from exc


@app.post("/api/integrations/telegram/webhook/{secret}")
async def api_telegram_webhook(secret: str, payload: dict[str, Any]) -> dict[str, Any]:
    if not telegram_bot.enabled:
        raise HTTPException(status_code=404, detail="Telegram no está habilitado")
    if settings.telegram_mode != "webhook":
        raise HTTPException(status_code=409, detail="El supervisor no está en modo webhook")
    if secret != settings.telegram_webhook_secret:
        raise HTTPException(status_code=403, detail="Secreto de webhook inválido")
    await telegram_bot.handle_update(payload)
    return {"ok": True}


@app.get("/api/memories", dependencies=[Depends(require_api_token)])
def api_memories(limit: int = 50) -> list[dict[str, Any]]:
    return store.list_memories(limit=limit)


@app.get("/api/memories/export", dependencies=[Depends(require_api_token)])
def api_memories_export() -> dict[str, Any]:
    """Exporta memorias y feedback para transferir aprendizaje a otra instancia."""
    return store.export_memories()


@app.post("/api/memories/import", dependencies=[Depends(require_api_token), Depends(rate_limit_strict)])
def api_memories_import(payload: MemoryImportRequest) -> dict[str, Any]:
    """Importa un dump generado por /api/memories/export (de otra instancia)."""
    result = store.import_memories(
        {"memories": payload.memories, "training_feedback": payload.training_feedback},
        actor=payload.actor,
    )
    store.push_sse_event("memories.imported", result)
    logger.info("memories.import", extra=result)
    return {"ok": True, **result}


@app.get("/api/audit", dependencies=[Depends(require_api_token)])
def api_audit(limit: int = 60) -> list[dict[str, Any]]:
    """Últimas entradas del audit_log (trazabilidad interna)."""
    return store.list_audit_log(limit=limit)


# ------------------------------------------------------------------ #
# Automatizaciones                                                      #
# ------------------------------------------------------------------ #

@app.get("/api/automation/state", dependencies=[Depends(require_api_token)])
def api_automation_state() -> dict[str, Any]:
    state = orchestrator.get_autonomy_state()
    state["scheduler"] = scheduler.get_status()
    return state


@app.get("/api/automation/rules", dependencies=[Depends(require_api_token)])
def api_automation_rules() -> list[dict[str, Any]]:
    return orchestrator.list_automation_rules()


@app.get("/api/automation/runs", dependencies=[Depends(require_api_token)])
def api_automation_runs(limit: int = 30) -> list[dict[str, Any]]:
    """Historial de ejecuciones del scheduler (auditoría)."""
    return store.list_automation_runs(limit=limit)


# ------------------------------------------------------------------ #
# WhatsApp Business API — Webhook                                       #
# ------------------------------------------------------------------ #

@app.get("/whatsapp/webhook")
async def whatsapp_webhook_verify(request: Request) -> Any:
    """Verificación del webhook por Meta (GET con hub.challenge)."""
    mode = request.query_params.get("hub.mode", "")
    token = request.query_params.get("hub.verify_token", "")
    challenge = request.query_params.get("hub.challenge", "")
    result = whatsapp.verify_webhook(mode, token, challenge)
    if result is None:
        raise HTTPException(status_code=403, detail="Token de verificación incorrecto")
    from fastapi.responses import PlainTextResponse
    return PlainTextResponse(result)


@app.post("/whatsapp/webhook")
async def whatsapp_webhook_receive(request: Request) -> dict[str, str]:
    """Recibe eventos de mensajes entrantes desde Meta."""
    body = await request.body()
    signature = request.headers.get("x-hub-signature-256", "")
    if not whatsapp.validate_signature(body, signature):
        raise HTTPException(status_code=403, detail="Firma HMAC inválida")
    try:
        payload = await request.json()
    except Exception:
        raise HTTPException(status_code=400, detail="JSON inválido")
    await whatsapp.handle_webhook(payload)
    return {"status": "ok"}


@app.get("/api/whatsapp/status", dependencies=[Depends(require_api_token)])
def api_whatsapp_status() -> dict[str, Any]:
    """Estado del servicio WhatsApp."""
    return whatsapp.status_summary()


@app.post("/api/whatsapp/send", dependencies=[Depends(require_api_token)])
async def api_whatsapp_send(request: Request) -> dict[str, Any]:
    """Envía un mensaje de texto o plantilla a un número de teléfono.

    Body JSON:
        {"to": "+34600000000", "text": "Hola", "template": null}
        {"to": "+34600000000", "template": "renovacion_anuncio", "lang": "es_ES", "components": [...]}
    """
    if not whatsapp.enabled:
        raise HTTPException(status_code=503, detail="WhatsApp no está habilitado")
    data = await request.json()
    to = data.get("to", "").strip()
    if not to:
        raise HTTPException(status_code=400, detail="Campo 'to' requerido")
    template = data.get("template")
    if template:
        result = await whatsapp.send_template(
            to,
            template_name=template,
            lang=data.get("lang", "es_ES"),
            components=data.get("components"),
        )
    else:
        text = data.get("text", "").strip()
        if not text:
            raise HTTPException(status_code=400, detail="Campo 'text' o 'template' requerido")
        result = await whatsapp.send_text(to, text)
    return result


async def _do_wa_campaign_send(proposal_id: str) -> dict[str, Any]:
    """Lógica de envío de campaña WA. Compartida por el endpoint HTTP y el bot de Telegram.

    Raises:
        HTTPException / ValueError / KeyError según el estado de la propuesta.
    """
    if not whatsapp.enabled:
        raise ValueError("WhatsApp no está habilitado en este supervisor")
    proposal = store.get_proposal(proposal_id)
    if proposal is None:
        raise KeyError(f"Propuesta no encontrada: {proposal_id}")
    if proposal.get("action_type") != "whatsapp_campaign":
        raise ValueError(f"La propuesta '{proposal_id}' no es de tipo whatsapp_campaign")
    current_status = proposal.get("status")
    # Guarda de idempotencia: evita re-envío si ya se está procesando o se ejecutó
    if current_status in {"sending", "executed", "failed"}:
        raise ValueError(f"La campaña ya fue procesada (estado: {current_status})")
    if current_status != "approved":
        raise ValueError(f"La propuesta debe estar aprobada para enviarse (estado actual: {current_status})")
    payload = proposal.get("payload") or {}
    contacts = payload.get("contacts") or []
    template_name = payload.get("template_name", "")
    lang = payload.get("lang", "es_ES")
    if not contacts:
        raise ValueError("La propuesta no incluye contactos (payload.contacts vacío)")
    if not template_name:
        raise ValueError("La propuesta no incluye template_name")
    store.set_proposal_execution(proposal_id, "sending", f"Enviando a {len(contacts)} contactos…")
    store.add_audit("whatsapp.campaign.started", {"proposal_id": proposal_id, "contacts": len(contacts), "template": template_name})
    results = await whatsapp.send_campaign(contacts, template_name, lang)
    sent = sum(1 for r in results if r.get("status") == "sent")
    failed = len(results) - sent
    summary = f"Campaña enviada: {sent}/{len(contacts)} OK, {failed} errores"
    store.set_proposal_execution(proposal_id, "executed", summary)
    store.add_audit("whatsapp.campaign.completed", {"proposal_id": proposal_id, "sent": sent, "failed": failed})
    return {"summary": summary, "sent": sent, "failed": failed, "results": results}


@app.post("/api/whatsapp/campaign/{proposal_id}/send", dependencies=[Depends(require_api_token)])
async def api_whatsapp_campaign_send(proposal_id: str) -> dict[str, Any]:
    """Ejecuta el envío de una campaña WhatsApp aprobada.

    El proposal debe tener action_type="whatsapp_campaign" y status="approved".
    En payload debe incluir: contacts (list), template_name, lang (opcional).
    Esta es una acción irreversible — se registra en el audit log.
    """
    try:
        return await _do_wa_campaign_send(proposal_id)
    except KeyError as exc:
        raise HTTPException(status_code=404, detail=str(exc)) from exc
    except ValueError as exc:
        # Distinguir 503 (WA deshabilitado) de 409 (estado incorrecto)
        detail = str(exc)
        code = 503 if "no está habilitado" in detail else 409
        raise HTTPException(status_code=code, detail=detail) from exc


# ------------------------------------------------------------------ #
# Métricas Prometheus                                                   #
# ------------------------------------------------------------------ #

@app.get("/metrics")
def metrics_endpoint(request: Request) -> Any:
    """Endpoint de métricas compatible con Prometheus scraper.
    
    Requiere token igual que el dashboard para no exponer contadores al público.
    Sin token se devuelve 403 en lugar de texto de métricas.
    """
    # Acepta: X-Supervisor-Token, ?token=, o Authorization: Bearer <token> (estándar Prometheus)
    auth_header = request.headers.get("Authorization", "")
    bearer = auth_header.removeprefix("Bearer ").strip() if auth_header.startswith("Bearer ") else ""
    token = request.headers.get("X-Supervisor-Token") or request.query_params.get("token") or bearer
    if not token or token != settings.dashboard_token:
        raise HTTPException(status_code=403, detail="Token requerido para /metrics")
    from fastapi.responses import PlainTextResponse
    body = generate_metrics(store, scheduler)
    return PlainTextResponse(body, media_type="text/plain; version=0.0.4; charset=utf-8")


@app.post("/api/automation/rules/{rule_key}/enable", dependencies=[Depends(require_api_token)])
def api_automation_rule_enable(rule_key: str, payload: AutomationRuleToggleRequest) -> dict[str, Any]:
    try:
        result = orchestrator.set_automation_rule_enabled(rule_key, payload.enabled, payload.actor, payload.note)
        store.push_sse_event("automation.rule.toggled", {"rule_key": rule_key, "enabled": payload.enabled})
        return result
    except KeyError as exc:
        raise HTTPException(status_code=404, detail=str(exc)) from exc


@app.post("/api/automation/rules/{rule_key}/run", dependencies=[Depends(require_api_token)])
async def api_automation_rule_run(rule_key: str, payload: AutomationRuleRunRequest) -> dict[str, Any]:
    try:
        result = await orchestrator.run_automation_rule(rule_key, payload.actor, payload.note)
        store.push_sse_event("automation.rule.executed", {"rule_key": rule_key, "actor": payload.actor})
        return result
    except KeyError as exc:
        raise HTTPException(status_code=404, detail=str(exc)) from exc


# ------------------------------------------------------------------ #
# Backups SQLite                                                        #
# ------------------------------------------------------------------ #

@app.get("/api/backups", dependencies=[Depends(require_api_token)])
def api_list_backups() -> list[dict[str, Any]]:
    """Lista los backups SQLite disponibles en data/backups/."""
    backup_dir = Path(settings.db_path).parent / "backups"
    if not backup_dir.exists():
        return []
    files = sorted(backup_dir.glob("supervisor_*.sqlite3"), reverse=True)
    result = []
    for f in files:
        stat = f.stat()
        result.append({
            "name": f.name,
            "size_kb": round(stat.st_size / 1024, 1),
            "created_at": time.strftime("%Y-%m-%d %H:%M:%S", time.gmtime(stat.st_mtime)),
        })
    return result


@app.post("/api/backups/trigger", dependencies=[Depends(require_api_token)])
async def api_trigger_backup() -> dict[str, Any]:
    """Lanza un backup manual inmediato fuera del cron."""
    from datetime import datetime, timezone
    now = datetime.now(timezone.utc)
    asyncio.create_task(scheduler._run_backup(now), name="manual-db-backup")
    return {"queued": True, "ts": now.isoformat()}


# ═══════════════════════════════════════════════════════════════════════════
# SSH — Biblioteca de comandos y propuesta directa
# ═══════════════════════════════════════════════════════════════════════════

@app.get("/api/ssh/commands", dependencies=[Depends(require_api_token)])
def api_ssh_commands() -> dict[str, Any]:
    """Lista todos los comandos SSH permitidos con descripción y nivel de riesgo."""
    try:
        commands = runner.list_allowed_commands()
        profile = settings.default_ssh_profile
        host = settings.ssh_host
        return {
            "profile": profile,
            "host": host,
            "port": settings.ssh_port,
            "user": settings.ssh_user,
            "use_agent": settings.ssh_use_agent,
            "total_commands": len(commands),
            "commands": [
                {
                    "key": key,
                    "description": meta.get("description", ""),
                    "risk": meta.get("risk", "medium"),
                }
                for key, meta in commands.items()
            ],
        }
    except Exception as exc:
        raise HTTPException(status_code=503, detail=str(exc)) from exc


@app.post("/api/ssh/propose", dependencies=[Depends(require_api_token)])
async def api_ssh_propose(body: dict[str, Any]) -> dict[str, Any]:
    """Crea directamente una propuesta SSH para aprobación humana.

    Body: { command_key: str, actor: str, reason?: str, profile?: str }
    La propuesta queda pendiente — ejecuta cuando el operador la apruebe.
    """
    command_key = str(body.get("command_key", "")).strip()
    actor = str(body.get("actor", "dashboard")).strip() or "dashboard"
    reason = str(body.get("reason", "")).strip()
    profile_key = str(body.get("profile", settings.default_ssh_profile)).strip() or settings.default_ssh_profile

    if not command_key:
        raise HTTPException(status_code=422, detail="command_key es obligatorio")

    try:
        spec = runner.get_command_spec(command_key, profile_key)
    except KeyError as exc:
        raise HTTPException(status_code=404, detail=str(exc)) from exc

    title = f"SSH: {spec['description']}"
    summary = f"Ejecutar `{command_key}` en {settings.ssh_host}:{settings.ssh_port}"
    if reason:
        summary += f"\n\nMotivo: {reason}"

    proposal = store.create_proposal(
        study_id="__ssh__",
        agent_name="ops",
        title=title,
        summary=summary,
        action_type="ssh_command",
        risk_level=spec["risk"],
        payload={"command_key": command_key, "profile": profile_key},
    )
    proposal_id = proposal["id"]
    store.add_audit(
        "ssh.proposed",
        {"command_key": command_key, "actor": actor, "reason": reason, "proposal_id": proposal_id},
    )
    return {
        "ok": True,
        "proposal_id": proposal_id,
        "command_key": command_key,
        "description": spec["description"],
        "risk": spec["risk"],
        "message": f"Propuesta creada (id={proposal_id[:8]}…). Apruébala en el Dashboard o Telegram para ejecutarla.",
    }


@app.post("/api/ssh/allowlist/suggest", dependencies=[Depends(require_api_token)])
async def api_ssh_allowlist_suggest(body: dict[str, Any]) -> dict[str, Any]:
    """Crea una propuesta de sugerencia para ampliar la allowlist SSH (sin ejecutar nada)."""
    actor = str(body.get("actor", "dashboard")).strip() or "dashboard"
    command_template = str(body.get("command_template", "")).strip()
    description = str(body.get("description", "")).strip()
    reason = str(body.get("reason", "")).strip()
    risk = str(body.get("risk", "medium")).strip().lower() or "medium"
    profile_key = str(body.get("profile", settings.default_ssh_profile)).strip() or settings.default_ssh_profile

    if not command_template:
        raise HTTPException(status_code=422, detail="command_template es obligatorio")
    if not description:
        raise HTTPException(status_code=422, detail="description es obligatoria")
    if risk not in {"low", "medium", "high"}:
        raise HTTPException(status_code=422, detail="risk debe ser low, medium o high")

    payload = {
        "profile": profile_key,
        "command_template": command_template,
        "description": description,
        "risk": risk,
        "reason": reason,
        "suggested_by": actor,
    }

    proposal = store.create_proposal(
        study_id="__ssh_allowlist__",
        agent_name="ops",
        title=f"SSH allowlist: sugerencia {description}",
        summary=(
            "Sugerencia de ampliación de allowlist SSH. Revisión humana obligatoria.\n\n"
            f"Perfil: {profile_key}\n"
            f"Comando: {command_template}\n"
            f"Riesgo: {risk}\n"
            f"Motivo: {reason or 'sin motivo'}"
        ),
        action_type="ssh_allowlist_suggestion",
        risk_level=risk,
        payload=payload,
    )
    proposal_id = proposal["id"]

    store.add_audit(
        "ssh.allowlist.suggested",
        {
            "proposal_id": proposal_id,
            "actor": actor,
            "profile": profile_key,
            "risk": risk,
        },
    )
    store.push_sse_event(
        "ssh.allowlist.suggested",
        {
            "proposal_id": proposal_id,
            "profile": profile_key,
            "risk": risk,
            "description": description,
        },
    )

    return {
        "ok": True,
        "proposal_id": proposal_id,
        "message": "Sugerencia creada. Apruébala para revisarla y aplicar cambios en allowlist.",
    }


@app.get("/api/ssh/test", dependencies=[Depends(require_api_token)])
async def api_ssh_test() -> dict[str, Any]:
    """Prueba la conectividad SSH con un echo inocuo (no registra propuesta)."""
    import asyncio as _asyncio
    import shlex

    host = settings.ssh_host
    port = settings.ssh_port
    user = settings.ssh_user

    if not host or not user:
        return {"ok": False, "error": "SSH_HOST o SSH_USER no configurados en .env"}

    ssh_args = [
        "ssh", "-p", str(port),
        "-o", "StrictHostKeyChecking=accept-new",
        "-o", "ConnectTimeout=8",
        "-o", "BatchMode=yes",
    ]
    if settings.ssh_use_agent:
        ssh_args += ["-o", "IdentitiesOnly=no"]
    elif settings.ssh_key_path:
        ssh_args += ["-i", settings.ssh_key_path]
    ssh_args += [f"{user}@{host}", "echo OK; hostname; uptime | cut -d, -f1"]

    try:
        process = await _asyncio.create_subprocess_exec(
            *ssh_args,
            stdout=_asyncio.subprocess.PIPE,
            stderr=_asyncio.subprocess.PIPE,
        )
        try:
            stdout, stderr = await _asyncio.wait_for(process.communicate(), timeout=15)
        except _asyncio.TimeoutError:
            process.kill()
            return {"ok": False, "error": "Timeout al conectar (15 s)"}
        output = (stdout or b"").decode("utf-8", errors="replace").strip()
        err = (stderr or b"").decode("utf-8", errors="replace").strip()
        if process.returncode == 0:
            return {"ok": True, "host": host, "port": port, "user": user, "output": output}
        return {"ok": False, "host": host, "port": port, "returncode": process.returncode, "stderr": err[:400]}
    except FileNotFoundError:
        return {"ok": False, "error": "Comando ssh no encontrado en el sistema"}
    except Exception as exc:
        return {"ok": False, "error": str(exc)}


# ═══════════════════════════════════════════════════════════════════════════
# RAG — Retrieval-Augmented Generation (LlamaIndex + Qdrant)
# ═══════════════════════════════════════════════════════════════════════════

@app.get("/api/rag/status", dependencies=[Depends(require_api_token)])
def api_rag_status() -> dict[str, Any]:
    """Estado del módulo RAG: colección, modelo de embeddings, número de vectores."""
    return rag_store.status()


@app.post("/api/rag/ingest", dependencies=[Depends(require_api_token)])
async def api_rag_ingest(payload: dict[str, Any]) -> dict[str, Any]:
    """Indexa documentos en Qdrant.

    Body (opcional):
    - `path`: ruta del directorio a indexar (default: docs/)
    - `texts`: lista de {text, id, source, type} para indexar directamente
    """
    if not rag_store.status().get("ready"):
        raise HTTPException(status_code=503, detail="RAG no está disponible (Qdrant no conectado o RAG_ENABLED=false)")

    path = payload.get("path")
    texts = payload.get("texts")

    if texts:
        if not isinstance(texts, list):
            raise HTTPException(status_code=422, detail="'texts' debe ser un array de objetos {text, id, source}")
        return await rag_store.ingest_texts(texts)

    return await rag_store.ingest_directory(path)


@app.post("/api/rag/query", dependencies=[Depends(require_api_token)])
async def api_rag_query(payload: dict[str, Any]) -> dict[str, Any]:
    """Recupera los fragmentos más relevantes para una consulta.

    Body:
    - `query` (str, requerido): texto a buscar semánticamente
    - `top_k` (int, default=5): número de fragmentos a devolver
    """
    query_text = str(payload.get("query", "")).strip()
    if not query_text:
        raise HTTPException(status_code=422, detail="El campo 'query' es obligatorio")
    if not rag_store.status().get("ready"):
        raise HTTPException(status_code=503, detail="RAG no está disponible")

    top_k = int(payload.get("top_k", 5))
    context = await rag_store.query(query_text, top_k=top_k)
    return {"query": query_text, "context": context, "top_k": top_k}

