"""
Servidor local de Leonex.

Sirve los ficheros estaticos del proyecto (dashboard, assets, etc.) y ademas
expone un endpoint:

    POST /api/refresh
        Ejecuta la cadena agente_datos -> agente_regimen -> agente_senales ->
        agente_validacion y devuelve un JSON con el resultado de cada paso.
        BLOQUEA hasta que termine el pipeline entero (puede ser ~40 minutos).

    POST /api/refresh-async
        Misma cadena, pero "fire-and-forget": lanza el pipeline en un thread
        en segundo plano y devuelve {"ok": true, "status": "started"} en <1s.
        Pensado para schedulers externos (n8n cron) que no quieren — ni deben —
        mantener una conexion HTTP abierta durante 40 minutos. Evita timeouts
        y desconexiones por redeploy. Solo permite un pipeline a la vez.

Asi el boton "Refrescar datos" del dashboard puede regenerar los JSONs en
caliente, sin necesidad de ejecutar manualmente los agentes desde la terminal.

Uso:
    python leonex_server.py
    python leonex_server.py --port 8765
    python leonex_server.py --strategy regime_adaptive_v1_lo

Compatible con Python 3.10+ (sin dependencias externas).
"""

from __future__ import annotations

import argparse
import json
import os
import socket
import subprocess
import sys
import threading
import time
from datetime import datetime, timezone
from http.server import SimpleHTTPRequestHandler, ThreadingHTTPServer
from pathlib import Path
from urllib.parse import urlparse


ROOT = Path(__file__).resolve().parent
DEFAULT_PORT = 8765
DEFAULT_STRATEGY = "regime_adaptive_v1_lo"
PIPELINE_TIMEOUT_S = 3600  # 60 min/agente: la descarga del S&P 500 completo
                           # (diario + intradia + scalping) no cabe en 10 min;
                           # con el tope antiguo de 600s se cortaba a medias.
RUN_AGENT_TIMEOUT_S = 3600  # 60 min: labs y descargas sobre el S&P 500

# Scheduler interno del bridge intradia. Ventana US market en UTC (L-V).
# Cada tick re-evalua las estrategias promovidas y, si disparan, manda orden a
# Alpaca paper. Corre dentro del contenedor: no depende de n8n y no gasta
# tokens. Alineado a multiplos del intervalo (con 15 min: :00/:15/:30/:45).
# Configurable via env LEONEX_BRIDGE_INTERVAL_MIN (default 15, minimo 5).
# OJO: bajar el intervalo NO genera mas senales — el plan diario se computa
# sobre datos diarios y las promovidas disparan cuando su barra dispara. Lo
# unico que mejora con 5 min es la resolucion de triggers en promovidas de
# timeframe 5m/15m. A cambio, cada ciclo ejecuta executor+close_monitor
# (~40-90s), asi que a 5 min el sistema pasa gran parte del tiempo ocupado.
try:
    _BRIDGE_INTERVAL_MIN = max(5, int(float(
        os.environ.get("LEONEX_BRIDGE_INTERVAL_MIN", "15").strip() or "15")))
except ValueError:
    _BRIDGE_INTERVAL_MIN = 15
INTRADAY_BRIDGE_INTERVAL_S = _BRIDGE_INTERVAL_MIN * 60
MARKET_OPEN_UTC = (13, 30)             # 13:30 UTC
MARKET_CLOSE_UTC = (20, 0)             # 20:00 UTC

# Estado compartido entre peticiones (un solo refresh a la vez)
_state_lock = threading.Lock()


class RefreshLock:
    """Wrapper sobre threading.Lock que registra holder + timestamp.

    Si una subprocess hace timeout o un job petó sin liberar, el lock se
    queda atascado y todas las acciones devuelven 409. Este wrapper:
    1) Mantiene metadata del holder para que las respuestas 409 digan
       quien lo tiene y desde hace cuanto tiempo (asi sabes si esta
       legitimamente ocupado o atascado).
    2) Permite forzar release via endpoint /api/release-lock cuando se
       detecta atasco (job ya terminado pero el lock no se libero).
    """

    def __init__(self) -> None:
        self._lock = threading.Lock()
        self._holder: str | None = None
        self._acquired_at: float | None = None

    def acquire(self, blocking: bool = False, holder: str = "unknown") -> bool:
        ok = self._lock.acquire(blocking=blocking)
        if ok:
            self._holder = holder
            self._acquired_at = time.time()
        return ok

    def release(self) -> None:
        # release defensivo: no peta si no se tenia
        try:
            self._lock.release()
        except RuntimeError:
            pass
        self._holder = None
        self._acquired_at = None

    def force_release(self) -> dict:
        """Devuelve metadata del holder previo (si lo tenia) y lo libera."""
        prev_holder = self._holder
        prev_age = (time.time() - self._acquired_at) if self._acquired_at else None
        was_held = self._lock.locked()
        try:
            self._lock.release()
        except RuntimeError:
            # No estaba tomado, no pasa nada
            pass
        self._holder = None
        self._acquired_at = None
        return {
            "was_held": was_held,
            "previous_holder": prev_holder,
            "age_seconds": round(prev_age, 1) if prev_age is not None else None,
        }

    def status(self) -> dict:
        held = self._lock.locked()
        age = (time.time() - self._acquired_at) if (held and self._acquired_at) else None
        return {
            "held": held,
            "holder": self._holder if held else None,
            "age_seconds": round(age, 1) if age is not None else None,
        }


_refresh_lock = RefreshLock()


def _python_executable() -> str:
    """Usa el venv del proyecto si existe, si no el python del PATH."""
    venv_win = ROOT / "venv" / "Scripts" / "python.exe"
    venv_unix = ROOT / "venv" / "bin" / "python"
    if venv_win.exists():
        return str(venv_win)
    if venv_unix.exists():
        return str(venv_unix)
    return sys.executable


def _dry_run_mode_forced() -> bool:
    """True si LEONEX_DRY_RUN_MODE esta activo: modo observacion global.
    El pipeline, las senales y el journal funcionan al 100%, pero NINGUNA
    orden automatica sale hacia Alpaca — scheduler interno, /api/execute,
    /api/close-monitor, /api/bridge-intraday y /api/daily-auto-trade quedan
    forzados a dry_run aunque el invocador pida mode=execute. La defensa
    real vive ademas DENTRO de agente_executor.py y agente_close_monitor.py
    (leen la misma env var), asi que ni un subproceso lanzado a mano con
    --execute puede enviar ordenes. Los flujos MANUALES de emergencia
    (boton Close por posicion, reset total) siguen operativos: un
    kill-switch humano nunca se bloquea."""
    return os.environ.get("LEONEX_DRY_RUN_MODE", "").strip().lower() in (
        "1", "true", "yes")


def _within_market_hours(now: datetime | None = None) -> bool:
    """True si `now` (UTC) cae en dia laborable (L-V) dentro de la ventana
    MARKET_OPEN_UTC..MARKET_CLOSE_UTC."""
    now = now or datetime.now(timezone.utc)
    if now.weekday() >= 5:                      # 5=sab, 6=dom
        return False
    open_m = MARKET_OPEN_UTC[0] * 60 + MARKET_OPEN_UTC[1]
    close_m = MARKET_CLOSE_UTC[0] * 60 + MARKET_CLOSE_UTC[1]
    cur_m = now.hour * 60 + now.minute
    return open_m <= cur_m < close_m


def _is_intraday_cache_stale() -> bool:
    """True si el cache 1h del SQLite tiene last_bar de antes de hoy UTC.
    Usado para que el bridge fuerce fetch=1 al detectar cache rancio. El
    daily corre pre-market (08:00 NY) y solo deja barras del dia anterior;
    el primer bridge del dia tiene que descargar incremental para poder
    evaluar contra precios reales del dia."""
    try:
        import sqlite3 as _sql
        db_file = ROOT / "data" / "Leonex.sqlite"
        if not db_file.exists():
            return False
        with _sql.connect(str(db_file)) as conn:
            row = conn.execute(
                "SELECT MAX(ts) FROM prices_intraday "
                "WHERE timeframe = '1h'"
            ).fetchone()
            last_ts = row[0] if row and row[0] else None
        if not last_ts:
            return True  # cache vacio = rancio
        today_utc = datetime.now(timezone.utc).strftime("%Y-%m-%d")
        return last_ts[:10] < today_utc
    except Exception:
        return False  # si falla la deteccion, no forzar


def _run_bridge_intraday_once(strategy: str, mode: str = "execute",
                              fetch: bool = False) -> dict:
    """Un ciclo del bridge intradia: executor (con bridge promovidas) +
    close_monitor. Reutilizado por el scheduler interno y por el endpoint
    HTTP /api/bridge-intraday. Toma el _refresh_lock de forma NO bloqueante:
    si hay otra accion en curso (refresh, daily-auto-trade, otro bridge), se
    salta este tick en vez de encolar. Devuelve {ok, skipped, steps}.

    Auto-promote fetch=False -> True si el cache 1h esta rancio y estamos en
    market hours. Esto evita que el scheduler interno evalue con datos del
    dia anterior cuando el daily corrio pre-market."""
    fetch_was_forced = False
    if not fetch and _within_market_hours() and _is_intraday_cache_stale():
        fetch = True
        fetch_was_forced = True

    if mode == "execute" and _dry_run_mode_forced():
        mode = "dry_run"

    if not _refresh_lock.acquire(blocking=False,
                                 holder=f"bridge_intraday_auto:{mode}"):
        return {"ok": False, "skipped": "lock_busy",
                "holder": _refresh_lock.status().get("holder"), "steps": []}

    py = _python_executable()
    # Mode del executor + close_monitor: ambos respetan el mode.
    # mode=execute -> envia ordenes Y cierra TP/SL automaticamente.
    # mode=dry_run -> solo simula ambos.
    if mode == "execute":
        executor_args = ["--strategy", strategy, "--execute", "--confirm",
                         "--auto-approve"]
        close_args = ["--execute", "--confirm"]
        close_label = "close_monitor (execute)"
    else:
        executor_args = ["--strategy", strategy, "--dry-run"]
        close_args = ["--dry-run"]
        close_label = "close_monitor (dry-run)"

    steps_def: list[tuple[str, list[str]]] = []
    if fetch:
        steps_def.append((
            "intraday data swing (1h/4h, incremental)",
            [py, "agents/agente_datos_intraday.py"],
        ))
    steps_def.extend([
        (f"executor ({mode}) + bridge",
         [py, "agents/agente_executor.py", *executor_args]),
        (close_label,
         [py, "agents/agente_close_monitor.py", *close_args]),
    ])
    # En modo execute anadimos risk_monitor al final para revalidar y
    # auto-limpiar el pause si los near_stops ya cerraron (evita el bug
    # donde el pause quedaba vivo con posiciones fantasma hasta el daily).
    if mode == "execute":
        steps_def.append((
            "risk_monitor (revalidate pause)",
            [py, "agents/agente_risk_monitor.py"],
        ))

    steps: list[dict] = []
    overall_ok = True
    try:
        for label, cmd in steps_def:
            t0 = time.time()
            step = {"step": label, "cmd": cmd}
            try:
                proc = subprocess.run(
                    cmd, cwd=str(ROOT), capture_output=True, text=True,
                    timeout=PIPELINE_TIMEOUT_S, encoding="utf-8",
                    errors="replace",
                )
                step.update({
                    "ok": proc.returncode == 0,
                    "returncode": proc.returncode,
                    "duration": round(time.time() - t0, 2),
                    "stdout_tail": (proc.stdout or "").splitlines()[-5:],
                    "stderr_tail": (proc.stderr or "").splitlines()[-3:],
                })
            except subprocess.TimeoutExpired:
                step.update({"ok": False,
                             "duration": round(time.time() - t0, 2),
                             "error": f"timeout_after_{PIPELINE_TIMEOUT_S}s"})
            except Exception as exc:
                step.update({"ok": False,
                             "duration": round(time.time() - t0, 2),
                             "error": repr(exc)})
            steps.append(step)
            if not step.get("ok"):
                overall_ok = False
                break
    finally:
        _refresh_lock.release()

    return {"ok": overall_ok, "skipped": None, "steps": steps,
            "fetch": fetch, "fetch_was_forced": fetch_was_forced}


def _intraday_bridge_scheduler(strategy: str) -> None:
    """Daemon: cada INTRADAY_BRIDGE_INTERVAL_S dispara el bridge intradia en
    modo execute, pero SOLO dentro de horario de mercado (L-V, ventana
    MARKET_OPEN_UTC..MARKET_CLOSE_UTC). Fuera de la ventana duerme. Se alinea
    a multiplos de 15 min del reloj para no derivar."""
    while True:
        now = datetime.now(timezone.utc)
        # Dormir hasta el proximo multiplo del intervalo configurado
        # (con 15 min: :00/:15/:30/:45; con 5 min: :00/:05/:10/...).
        secs_into = (now.minute % _BRIDGE_INTERVAL_MIN) * 60 + now.second
        sleep_s = INTRADAY_BRIDGE_INTERVAL_S - secs_into
        if sleep_s <= 0:
            sleep_s = INTRADAY_BRIDGE_INTERVAL_S
        time.sleep(sleep_s)

        if not _within_market_hours():
            continue
        try:
            res = _run_bridge_intraday_once(strategy, mode="execute",
                                            fetch=False)
            ts = datetime.now(timezone.utc).strftime("%Y-%m-%d %H:%M UTC")
            if res.get("skipped"):
                print(f"[intraday-bridge] {ts} saltado ({res['skipped']})")
            else:
                print(f"[intraday-bridge] {ts} ejecutado: ok={res.get('ok')}, "
                      f"pasos={len(res.get('steps') or [])}")
            sys.stdout.flush()
        except Exception as exc:  # nunca matar el hilo
            print(f"[intraday-bridge] error en ciclo: {exc!r}")
            sys.stdout.flush()


def _build_pipeline(strategy: str) -> list[tuple[str, list[str], bool]]:
    """Pipeline de pasos. Cada tupla es (label, cmd, fatal).

    fatal=True  → si falla, parar el pipeline (datos imprescindibles)
    fatal=False → si falla, registrar el error pero CONTINUAR. Pasos de
                  observabilidad o que dependen del estado del mercado
                  (p.ej. HRP sin plan, Executor sin Alpaca) no deben tumbar
                  todo el refresh.
    """
    py = _python_executable()
    return [
        ("Agente Universe (S&P 500 top-30 + legacy)",
         [py, "agents/agente_universe.py"],                               False),
        ("Agente de Datos", [py, "agents/agente_datos.py"],               True),
        ("Agente Datos Intraday (1h/4h equities+forex)",
         [py, "agents/agente_datos_intraday.py"],                         False),
        ("Agente Datos Intraday SCALPING (30m/15m/5m)",
         [py, "agents/agente_datos_intraday.py",
          "--timeframes", "30m,15m,5m"],                                  False),
        ("Agente Data Coverage (frescura/completitud de los datos)",
         [py, "agents/agente_data_coverage.py"],                          False),
        ("Agente Timeframe Selector (daily vs 1h vs 4h por activo)",
         [py, "agents/agente_timeframe_selector.py"],                     False),
        ("Agente de Regimen", [py, "agents/agente_regimen.py"],           True),
        ("Agente de Senales", [py, "agents/agente_senales.py"],           True),
        ("Agente Triple Barrier",
         [py, "agents/agente_triple_barrier.py", "--strategy", strategy], True),
        ("Agente Causal (PC algorithm — features causa vs spurious)",
         [py, "agents/agente_causal.py", "--explain"],                    False),
        # Extended (todas las features): el OOF AUC del modo causal-only (0.509)
        # quedaba por DEBAJO del extended (0.522) -> causal-only dejaba edge en
        # la mesa. Solo 2 features causales (days_held, rsi_14) no bastan.
        ("Agente Meta-Labeling (extended: todas las features)",
         [py, "agents/agente_meta.py", "--strategy", strategy],            True),
        ("Agente Validacion (Triple Barrier + Meta filter 0.55)",
         [py, "agents/agente_validacion.py", "--strategy", strategy,
          "--triple-barrier", "--meta-filter", "0.55"],                   True),
        ("Agente CPCV (Combinatorial Purged CV + PBO — anti-overfitting)",
         [py, "agents/agente_cpcv.py", "--strategy", strategy],           False),
        ("Agente Drift Monitor (revalidacion + flag pausa)",
         [py, "agents/agente_drift_monitor.py", "--window", "60"],        False),
        ("Agente Executor (Alpaca lectura)",
         [py, "agents/agente_executor.py", "--strategy", strategy,
          "--threshold", "0.55"],                                         False),
        ("Agente HRP (sizing por contribucion al riesgo)",
         [py, "agents/agente_hrp.py", "--plan"],                          False),
        ("Agente Close Monitor (lectura TP/SL/timeout)",
         [py, "agents/agente_close_monitor.py"],                          False),
        ("Agente Risk Monitor (drawdown, leverage, near-stop, vol spike)",
         [py, "agents/agente_risk_monitor.py"],                           False),
        ("Agente Strategy Lab (mejor estrategia/activo con DSR anti-overfitting)",
         [py, "agents/agente_strategy_lab.py"],                           False),
        ("Agente Strategy Lab SCALPING (30m/15m/5m con costes de transaccion)",
         [py, "agents/agente_strategy_lab_scalping.py"],                  False),
        ("Agente Portfolio Lab (Ley Fundamental — combina estrategias en cartera)",
         [py, "agents/agente_portfolio_lab.py"],                          False),
        ("Agente Permutation Test (Aronson — edge real vs suerte)",
         [py, "agents/agente_permutation_test.py"],                       False),
        ("Agente Carver (forecasts continuos + volatility targeting)",
         [py, "agents/agente_carver.py"],                                 False),
        ("Agente Strategy Tracker (forward paper test de estrategias promovidas)",
         [py, "agents/agente_strategy_tracker.py"],                       False),
        ("Agente Custom Lab (estrategias custom / PineScript traducido)",
         [py, "agents/agente_custom_lab.py"],                             False),
        ("Agente Scheduler Info (Task Scheduler)",
         [py, "agents/agente_scheduler_info.py"],                         False),
    ]


def _run_pipeline(strategy: str) -> dict:
    pipeline = _build_pipeline(strategy)
    results = []
    overall_ok = True
    n_warnings = 0
    for label, cmd, fatal in pipeline:
        step_result = {"step": label, "cmd": cmd, "fatal": fatal}
        try:
            proc = subprocess.run(
                cmd,
                cwd=str(ROOT),
                capture_output=True,
                text=True,
                timeout=PIPELINE_TIMEOUT_S,
                encoding="utf-8",
                errors="replace",
            )
            step_result["returncode"] = proc.returncode
            step_result["stdout_tail"] = (proc.stdout or "").splitlines()[-8:]
            step_result["stderr_tail"] = (proc.stderr or "").splitlines()[-5:]
            step_result["ok"] = proc.returncode == 0
        except subprocess.TimeoutExpired:
            step_result["ok"] = False
            step_result["error"] = f"Timeout tras {PIPELINE_TIMEOUT_S}s"
        except FileNotFoundError as exc:
            step_result["ok"] = False
            step_result["error"] = f"Comando no encontrado: {exc}"
        except Exception as exc:  # pragma: no cover
            step_result["ok"] = False
            step_result["error"] = repr(exc)

        results.append(step_result)

        if not step_result["ok"]:
            if fatal:
                # Paso imprescindible fallido → abortamos el pipeline
                overall_ok = False
                step_result["aborted_pipeline"] = True
                break
            else:
                # Paso opcional fallido → warning, seguimos
                n_warnings += 1
                step_result["warning"] = True

    return {
        "ok": overall_ok,
        "strategy": strategy,
        "warnings": n_warnings,
        "results": results,
    }


class LeonexHandler(SimpleHTTPRequestHandler):
    server_version = "Leonex/1.0"
    protocol_version = "HTTP/1.1"

    # Servimos siempre desde la raiz del proyecto (no la cwd, que puede variar)
    def translate_path(self, path: str) -> str:
        rel = urlparse(path).path.lstrip("/")
        # http.server normaliza con os.getcwd; forzamos ROOT
        full = (ROOT / rel).resolve()
        # Evitamos escapes (../) saliendo de ROOT
        try:
            full.relative_to(ROOT)
        except ValueError:
            return str(ROOT)
        return str(full)

    def log_message(self, fmt: str, *args) -> None:  # silenciamos ruido
        sys.stderr.write(f"[{self.log_date_time_string()}] {fmt % args}\n")

    # ── CORS / preflight para entornos donde el navegador sea estricto ─────
    def end_headers(self) -> None:
        self.send_header("Cache-Control", "no-store")
        super().end_headers()

    def do_OPTIONS(self) -> None:
        self.send_response(204)
        self.send_header("Allow", "GET, POST, OPTIONS")
        self.end_headers()

    # ── GET ──────────────────────────────────────────────────────────────
    def _handle_lock_status(self) -> None:
        """GET /api/lock-status — devuelve si el _refresh_lock esta ocupado,
        quien lo tiene (holder name) y desde hace cuanto. Util para el
        dashboard antes de proponer el boton 'Release stuck lock'."""
        body = json.dumps(_refresh_lock.status(), ensure_ascii=False).encode("utf-8")
        self._send_json(200, body)

    def _handle_test_custom_bridge(self, parsed) -> None:
        """GET/POST /api/test-custom-bridge?ticker=X&strategy=Y[&timeframe=1h]
        Ejecuta evaluate_promoted aislada para esa combinacion (ticker, estrategia,
        timeframe) y devuelve el JSON con live/prob_win/trigger_on/note. No toca
        Alpaca ni ninguna BD: pura evaluacion del bridge sobre datos historicos
        en sqlite. Pensado para verificacion post-deploy de las custom strategies.
        """
        from urllib.parse import parse_qs
        qs = parse_qs(parsed.query or "")
        ticker = (qs.get("ticker", [""])[0]).strip().upper()
        strategy = (qs.get("strategy", [""])[0]).strip()
        timeframe = (qs.get("timeframe", ["1h"])[0]).strip() or "1h"
        if not ticker or not strategy:
            body = json.dumps({
                "ok": False,
                "error": "missing_params",
                "message": "Faltan 'ticker' y/o 'strategy' en la query string.",
                "example": "/api/test-custom-bridge?ticker=CEG&strategy=box_prevday_hammer_long&timeframe=1h",
            }, ensure_ascii=False).encode("utf-8")
            self._send_json(400, body)
            return
        try:
            sys.path.insert(0, str(ROOT / "agents"))
            from agente_bridge_promoted import (             # noqa: E402
                evaluate_promoted,
                CUSTOM_STRATEGIES_BY_NAME,
                _CUSTOM_AVAILABLE,
            )
            p = {
                "ticker": ticker, "strategy": strategy,
                "timeframe": timeframe, "dsr": 0.0,
            }
            ev = evaluate_promoted(p)
            # Anadimos contexto util para el dashboard
            resp = {
                "ok": True,
                "input": {"ticker": ticker, "strategy": strategy,
                          "timeframe": timeframe},
                "custom_module_available": bool(_CUSTOM_AVAILABLE),
                "is_registered_custom": strategy in CUSTOM_STRATEGIES_BY_NAME,
                "known_custom_strategies": sorted(CUSTOM_STRATEGIES_BY_NAME.keys()),
                "evaluation": ev,
                "summary": (
                    f"{ticker} / {strategy} @ {timeframe} → "
                    f"live={ev.get('live')}  "
                    f"prob_win={ev.get('prob_win', 0):.2%}  "
                    f"trigger_on={ev.get('trigger_on')}  "
                    f"last_bar={ev.get('last_bar', '?')}"
                ),
            }
        except Exception as exc:
            resp = {
                "ok": False,
                "error": "exception",
                "message": repr(exc),
                "input": {"ticker": ticker, "strategy": strategy,
                          "timeframe": timeframe},
            }
        body = json.dumps(resp, ensure_ascii=False, default=str).encode("utf-8")
        self._send_json(200 if resp.get("ok") else 500, body)

    def _handle_release_lock(self) -> None:
        """POST /api/release-lock — fuerza release del _refresh_lock. Devuelve
        info del holder previo (si lo habia). Si no estaba tomado, devuelve
        was_held=False. Usar SOLO si se detecta atasco real (un job que ya
        terminó pero el lock no se liberó por un crash/timeout)."""
        info = _refresh_lock.force_release()
        info["ok"] = True
        info["message"] = (
            "Lock liberado."
            if info.get("was_held")
            else "El lock no estaba tomado — no habia nada que liberar."
        )
        body = json.dumps(info, ensure_ascii=False).encode("utf-8")
        self._send_json(200, body)

    def _handle_reset_pause(self) -> None:
        """POST /api/reset-pause — Desactiva PAUSE_NEW_TRADES y limpia la
        razon. Uso: tras un reset_paper_total, o cuando el usuario quiere
        reanudar tras revisar. No hace nada mas."""
        try:
            import sqlite3 as _sql
            db_file = ROOT / "data" / "Leonex.sqlite"
            with _sql.connect(str(db_file)) as conn:
                conn.execute(
                    "CREATE TABLE IF NOT EXISTS system_state ("
                    "key TEXT PRIMARY KEY, value TEXT, updated_at TEXT)"
                )
                ts = datetime.now(timezone.utc).isoformat()
                for k, v in (("pause_new_trades", "false"),
                             ("pause_reason", "")):
                    conn.execute(
                        "INSERT INTO system_state(key,value,updated_at) VALUES(?,?,?) "
                        "ON CONFLICT(key) DO UPDATE SET value=excluded.value, "
                        "updated_at=excluded.updated_at",
                        (k, v, ts),
                    )
                conn.commit()
            payload = {"ok": True, "message": "PAUSE_NEW_TRADES desactivado."}
            self._send_json(200, json.dumps(payload).encode("utf-8"))
        except Exception as exc:
            payload = {"ok": False, "error": repr(exc)}
            self._send_json(500, json.dumps(payload).encode("utf-8"))

    def _handle_reset_paper_total(self, parsed) -> None:
        """POST /api/reset-paper-total?execute=1 — RESET TOTAL del paper.
        Cierra TODAS las posiciones en Alpaca, cancela todas las ordenes
        pendientes, marca todas las open_trades como closed_manual_reset y
        limpia el pause. Devuelve JSON con el resumen detallado.

        Por defecto (sin query o ?execute=0) es DRY-RUN: lista que cerraria
        pero NO envia ordenes. Con ?execute=1 lo hace de verdad.

        Toma el _refresh_lock para evitar que un cron intraday se solape con
        el reset y vuelva a abrir posiciones justo mientras cerramos."""
        from urllib.parse import parse_qs
        qs = parse_qs(parsed.query or "")
        execute = (qs.get("execute", ["0"])[0]).strip() in ("1", "true", "yes")

        if not self._acquire_lock_or_send_409(
                f"reset_paper_total:{'execute' if execute else 'dry_run'}"):
            return

        try:
            py = _python_executable()
            cmd = [py, "tools/reset_paper_total.py"]
            if execute:
                cmd.append("--execute")
            proc = subprocess.run(
                cmd, cwd=str(ROOT), capture_output=True, text=True,
                timeout=300, encoding="utf-8", errors="replace",
            )
            payload = {}
            if proc.stdout:
                # El script imprime JSON al stdout; parseamos para exponerlo.
                try:
                    payload = json.loads(proc.stdout)
                except Exception:
                    payload = {"raw_stdout": proc.stdout[-4000:]}
            payload.setdefault("script_returncode", proc.returncode)
            if proc.stderr:
                payload["stderr_tail"] = proc.stderr.splitlines()[-10:]
            payload["ok"] = proc.returncode == 0
            body = json.dumps(payload, ensure_ascii=False, default=str
                              ).encode("utf-8")
            self._send_json(200 if proc.returncode == 0 else 500, body)
        finally:
            _refresh_lock.release()

    def do_GET(self) -> None:
        parsed = urlparse(self.path)
        if parsed.path == "/api/refresh/stream":
            return self._handle_refresh_stream()
        # Disponible por GET tambien para que un cron/scheduler simple (o un
        # web fetch) lo pueda disparar sin POST.
        if parsed.path == "/api/daily-auto-trade":
            return self._handle_daily_auto_trade(parsed)
        if parsed.path == "/api/lock-status":
            return self._handle_lock_status()
        if parsed.path == "/api/test-custom-bridge":
            return self._handle_test_custom_bridge(parsed)
        if parsed.path == "/api/system-info":
            return self._handle_system_info()
        return super().do_GET()

    def _handle_system_info(self) -> None:
        """GET /api/system-info — Info del estado del sistema:
        dry_run_mode global (env var), allow_shorts, versiones, timestamp.
        Usado por el dashboard para pintar el banner amarillo cuando
        LEONEX_DRY_RUN_MODE=1 esta activo."""
        payload = {
            "dry_run_mode_forced": _dry_run_mode_forced(),
            "allow_shorts": os.environ.get(
                "LEONEX_ALLOW_SHORTS", "0").strip().lower() in (
                "1", "true", "yes"),
            "max_positions": int(os.environ.get(
                "LEONEX_MAX_POSITIONS", "25")),
            "generated_at": datetime.now(timezone.utc).isoformat(),
        }
        body = json.dumps(payload).encode("utf-8")
        self._send_json(200, body)

    def _handle_refresh_stream(self) -> None:
        """Stream del pipeline en formato Server-Sent Events.

        Emite eventos:
            start       → {total, strategy}
            step_start  → {step, total, label}
            step_done   → {step, label, ok, returncode, duration, stdout_tail, stderr_tail}
            complete    → {ok, total_duration}
            error       → {error, message}

        Cliente: usar EventSource (solo GET) o fetch + stream reader.
        """
        # Headers SSE
        self.send_response(200)
        self.send_header("Content-Type", "text/event-stream")
        self.send_header("Cache-Control", "no-cache, no-transform")
        self.send_header("Connection", "close")
        self.send_header("X-Accel-Buffering", "no")  # nginx no buffer
        # Para SSE necesitamos cerrar la conexion al final, no keep-alive
        # SimpleHTTPRequestHandler con HTTP/1.1 lo respeta gracias a Connection:close
        self.end_headers()

        def emit(event_name: str, data) -> bool:
            try:
                msg = f"event: {event_name}\ndata: {json.dumps(data, ensure_ascii=False)}\n\n"
                self.wfile.write(msg.encode("utf-8"))
                self.wfile.flush()
                return True
            except (BrokenPipeError, ConnectionResetError, OSError):
                return False

        # Lock: no permitir dos refreshes simultaneos
        if not _refresh_lock.acquire(blocking=False, holder="refresh_stream"):
            s = _refresh_lock.status()
            emit("error", {"error": "refresh_in_progress",
                           "message": (f"Another refresh is already running: "
                                       f"{s.get('holder') or 'unknown'} "
                                       f"({s.get('age_seconds') or '?'}s ago)."),
                           "lock_status": s})
            emit("complete", {"ok": False})
            return

        import time as _time
        try:
            strategy = getattr(self.server, "strategy", DEFAULT_STRATEGY)
            pipeline = _build_pipeline(strategy)
            total = len(pipeline)
            t_overall = _time.time()
            emit("start", {"total": total, "strategy": strategy})
            all_ok = True

            for i, (label, cmd, fatal) in enumerate(pipeline, 1):
                if not emit("step_start", {"step": i, "total": total,
                                           "label": label, "fatal": fatal}):
                    # Cliente desconecto — abortamos sin levantar nada raro
                    return
                t0 = _time.time()
                step_payload = {"step": i, "label": label, "fatal": fatal}
                try:
                    proc = subprocess.run(
                        cmd, cwd=str(ROOT), capture_output=True, text=True,
                        timeout=PIPELINE_TIMEOUT_S,
                        encoding="utf-8", errors="replace",
                    )
                    step_payload.update({
                        "ok": proc.returncode == 0,
                        "returncode": proc.returncode,
                        "duration": round(_time.time() - t0, 2),
                        "stdout_tail": (proc.stdout or "").splitlines()[-5:],
                        "stderr_tail": (proc.stderr or "").splitlines()[-3:],
                    })
                except subprocess.TimeoutExpired:
                    step_payload.update({
                        "ok": False,
                        "duration": round(_time.time() - t0, 2),
                        "error": f"timeout_after_{PIPELINE_TIMEOUT_S}s",
                    })
                except Exception as exc:
                    step_payload.update({
                        "ok": False,
                        "duration": round(_time.time() - t0, 2),
                        "error": repr(exc),
                    })

                # Pasos opcionales fallidos = warning, no rompen el pipeline
                if not step_payload.get("ok"):
                    if fatal:
                        all_ok = False
                        step_payload["aborted_pipeline"] = True
                    else:
                        step_payload["warning"] = True

                if not emit("step_done", step_payload):
                    return
                # Solo paramos si fallo un paso FATAL
                if not step_payload.get("ok") and fatal:
                    break

            emit("complete", {
                "ok": all_ok,
                "total_duration": round(_time.time() - t_overall, 2),
            })
        finally:
            _refresh_lock.release()

    # ── POST ──────────────────────────────────────────────────────────────
    def do_POST(self) -> None:
        parsed = urlparse(self.path)
        if parsed.path == "/api/refresh":
            return self._handle_refresh()
        if parsed.path == "/api/refresh-async":
            return self._handle_refresh_async()
        if parsed.path == "/api/execute":
            return self._handle_execute(parsed)
        if parsed.path == "/api/close-monitor":
            return self._handle_close_monitor(parsed)
        if parsed.path == "/api/run-agent":
            return self._handle_run_agent(parsed)
        if parsed.path == "/api/promote":
            return self._handle_promote(parsed)
        if parsed.path == "/api/save-pinescript":
            return self._handle_save_pinescript(parsed)
        if parsed.path == "/api/alpaca-diagnose":
            return self._handle_alpaca_diagnose()
        if parsed.path == "/api/promoted-diagnose":
            return self._handle_promoted_diagnose()
        if parsed.path == "/api/daily-auto-trade":
            return self._handle_daily_auto_trade(parsed)
        if parsed.path == "/api/refresh-trades":
            return self._handle_refresh_trades()
        if parsed.path == "/api/bridge-intraday":
            return self._handle_bridge_intraday(parsed)
        if parsed.path == "/api/close-position":
            return self._handle_close_position(parsed)
        if parsed.path == "/api/release-lock":
            return self._handle_release_lock()
        if parsed.path == "/api/test-custom-bridge":
            return self._handle_test_custom_bridge(parsed)
        if parsed.path == "/api/reset-paper-total":
            return self._handle_reset_paper_total(parsed)
        if parsed.path == "/api/reset-pause":
            return self._handle_reset_pause()
        self.send_error(404, "Not Found")

    def _handle_refresh(self) -> None:
        strategy = getattr(self.server, "strategy", DEFAULT_STRATEGY)
        # Evitamos refreshes concurrentes
        if not self._acquire_lock_or_send_409("refresh"):
            return
        try:
            payload = _run_pipeline(strategy)
        finally:
            _refresh_lock.release()
        body = json.dumps(payload, ensure_ascii=False).encode("utf-8")
        status = 200 if payload["ok"] else 500
        self._send_json(status, body)

    def _handle_refresh_async(self) -> None:
        """Fire-and-forget: lanza el pipeline en un thread y devuelve YA.

        Pensado para schedulers externos (n8n cron). El cliente recibe la
        respuesta en menos de un segundo y se desconecta — el pipeline sigue
        corriendo en segundo plano dentro del contenedor de Leonex. Con esto:
            - Adios timeouts (el scheduler no espera 40 min, se desentiende).
            - Adios "connection closed unexpectedly" si Leonex se reinicia
              por un deploy a mitad del pipeline — el deploy interrumpe el
              pipeline igualmente, pero n8n ya marco el workflow como OK.
            - Solo se permite UN pipeline a la vez (mismo lock que /api/refresh
              y la SSE; si se solapan, devolvemos 409).
        El dashboard sigue usando /api/refresh/stream (SSE) para ver progreso
        en vivo, y /api/refresh sincrono sigue disponible por compatibilidad.
        """
        strategy = getattr(self.server, "strategy", DEFAULT_STRATEGY)
        if not self._acquire_lock_or_send_409("refresh_async"):
            return

        def _runner() -> None:
            import time as _t
            t0 = _t.time()
            sys.stderr.write(
                f"[refresh-async] pipeline started (strategy={strategy})\n")
            try:
                result = _run_pipeline(strategy)
                elapsed = round(_t.time() - t0, 1)
                outcome = "OK" if result.get("ok") else "FAILED"
                sys.stderr.write(
                    f"[refresh-async] pipeline {outcome} in {elapsed}s "
                    f"(warnings={result.get('warnings', 0)})\n")
            except Exception as exc:                  # pragma: no cover
                sys.stderr.write(
                    f"[refresh-async] pipeline crashed: {exc!r}\n")
            finally:
                _refresh_lock.release()

        threading.Thread(target=_runner, name="refresh-async",
                         daemon=True).start()
        body = json.dumps({
            "ok": True,
            "status": "started",
            "strategy": strategy,
            "message": ("Pipeline lanzado en segundo plano. Esta respuesta no "
                        "espera a que termine — mira los logs del contenedor "
                        "para ver el progreso."),
        }, ensure_ascii=False).encode("utf-8")
        self._send_json(202, body)

    def _handle_execute(self, parsed) -> None:
        """Lanza agente_executor.py --execute --confirm y devuelve resultado."""
        from urllib.parse import parse_qs
        # Permitir override de modo via query: ?mode=dry_run|execute
        qs = parse_qs(parsed.query or "")
        mode = (qs.get("mode", ["execute"])[0]).lower()
        if mode == "execute" and _dry_run_mode_forced():
            mode = "dry_run"
        strategy = getattr(self.server, "strategy", DEFAULT_STRATEGY)
        py = _python_executable()
        if mode == "dry_run":
            cmd = [py, "agents/agente_executor.py", "--strategy", strategy,
                   "--threshold", "0.55", "--dry-run"]
        elif mode == "execute":
            cmd = [py, "agents/agente_executor.py", "--strategy", strategy,
                   "--threshold", "0.55", "--execute", "--confirm"]
        else:
            body = json.dumps({
                "ok": False,
                "error": f"modo no soportado: {mode}. Usa mode=dry_run o mode=execute"
            }).encode("utf-8")
            self._send_json(400, body)
            return

        if not self._acquire_lock_or_send_409("execute"):
            return

        try:
            proc = subprocess.run(
                cmd, cwd=str(ROOT), capture_output=True, text=True,
                timeout=PIPELINE_TIMEOUT_S, encoding="utf-8", errors="replace",
            )
            payload = {
                "ok": proc.returncode == 0,
                "mode": mode,
                "returncode": proc.returncode,
                "stdout_tail": (proc.stdout or "").splitlines()[-30:],
                "stderr_tail": (proc.stderr or "").splitlines()[-10:],
            }
        except Exception as exc:
            payload = {"ok": False, "mode": mode, "error": repr(exc)}
        finally:
            _refresh_lock.release()

        body = json.dumps(payload, ensure_ascii=False).encode("utf-8")
        self._send_json(200 if payload.get("ok") else 500, body)

    def _handle_close_monitor(self, parsed) -> None:
        """Lanza agente_close_monitor.py en dry_run o execute."""
        from urllib.parse import parse_qs
        qs = parse_qs(parsed.query or "")
        mode = (qs.get("mode", ["dry_run"])[0]).lower()
        if mode == "execute" and _dry_run_mode_forced():
            mode = "dry_run"
        py = _python_executable()
        if mode == "dry_run":
            cmd = [py, "agents/agente_close_monitor.py", "--dry-run"]
        elif mode == "execute":
            cmd = [py, "agents/agente_close_monitor.py", "--execute", "--confirm"]
        else:
            body = json.dumps({"ok": False, "error": f"modo no soportado: {mode}"}).encode("utf-8")
            self._send_json(400, body)
            return

        if not self._acquire_lock_or_send_409("close_monitor"):
            return
        try:
            proc = subprocess.run(
                cmd, cwd=str(ROOT), capture_output=True, text=True,
                timeout=PIPELINE_TIMEOUT_S, encoding="utf-8", errors="replace",
            )
            payload = {
                "ok": proc.returncode == 0,
                "mode": mode,
                "returncode": proc.returncode,
                "stdout_tail": (proc.stdout or "").splitlines()[-30:],
                "stderr_tail": (proc.stderr or "").splitlines()[-10:],
            }
        except Exception as exc:
            payload = {"ok": False, "mode": mode, "error": repr(exc)}
        finally:
            _refresh_lock.release()

        body = json.dumps(payload, ensure_ascii=False).encode("utf-8")
        self._send_json(200 if payload.get("ok") else 500, body)

    def _handle_run_agent(self, parsed) -> None:
        """Lanza UN agente de analisis bajo demanda, sin correr el pipeline
        entero. Pensado para los labs que viven al final del pipeline (Custom
        Lab, Strategy Tracker, Strategy Lab, Scalping Lab): el usuario quiere
        refrescarlos sin re-descargar todos los datos del S&P 500. Solo se
        permiten agentes de SOLO LECTURA / analisis — nada que opere ordenes."""
        from urllib.parse import parse_qs

        # Whitelist: clave -> comando relativo. Solo agentes de analisis y de
        # descarga de datos de mercado — nada que opere ordenes ni mueva dinero.
        whitelist = {
            "custom_lab":        ["agents/agente_custom_lab.py"],
            "strategy_tracker":  ["agents/agente_strategy_tracker.py"],
            "strategy_lab":      ["agents/agente_strategy_lab.py"],
            "scalping_lab":      ["agents/agente_strategy_lab_scalping.py"],
            "data_coverage":     ["agents/agente_data_coverage.py"],
            "portfolio_lab":     ["agents/agente_portfolio_lab.py"],
            "permutation_test":  ["agents/agente_permutation_test.py"],
            "carver":            ["agents/agente_carver.py"],
            "daily_data":        ["agents/agente_datos.py"],
            "intraday_data":     ["agents/agente_datos_intraday.py"],
            "intraday_scalping": ["agents/agente_datos_intraday.py",
                                  "--timeframes", "30m,15m,5m"],
            # Limpieza de promociones cuyo trigger/filter/exit ya no exista
            # en el catalogo actual (zombies del catalogo viejo).
            "cleanup_stale":     ["agents/agente_strategy_lab.py",
                                  "--cleanup-stale"],
        }
        qs = parse_qs(parsed.query or "")
        agent = (qs.get("agent", [""])[0]).lower().strip()
        rel = whitelist.get(agent)
        if not rel:
            body = json.dumps({
                "ok": False,
                "error": f"agente no permitido: {agent!r}. "
                         f"Validos: {sorted(whitelist)}",
            }).encode("utf-8")
            self._send_json(400, body)
            return

        cmd = [_python_executable(), *rel]
        if not self._acquire_lock_or_send_409("run_agent"):
            return
        try:
            proc = subprocess.run(
                cmd, cwd=str(ROOT), capture_output=True, text=True,
                timeout=RUN_AGENT_TIMEOUT_S, encoding="utf-8",
                errors="replace",
            )
            payload = {
                "ok": proc.returncode == 0,
                "agent": agent,
                "returncode": proc.returncode,
                "stdout_tail": (proc.stdout or "").splitlines()[-30:],
                "stderr_tail": (proc.stderr or "").splitlines()[-12:],
            }
        except subprocess.TimeoutExpired:
            payload = {"ok": False, "agent": agent,
                       "error": f"timeout tras {RUN_AGENT_TIMEOUT_S}s"}
        except Exception as exc:
            payload = {"ok": False, "agent": agent, "error": repr(exc)}
        finally:
            _refresh_lock.release()

        body = json.dumps(payload, ensure_ascii=False).encode("utf-8")
        self._send_json(200 if payload.get("ok") else 500, body)

    def _handle_promote(self, parsed) -> None:
        """Promueve una estrategia: la registra en asset_strategies para que el
        Strategy Tracker la siga en forward test. Solo registra en la BD — no
        opera ordenes ni mueve dinero."""
        from urllib.parse import parse_qs
        qs = parse_qs(parsed.query or "")
        ticker = (qs.get("ticker", [""])[0]).strip()
        strategy = (qs.get("strategy", [""])[0]).strip()
        if not ticker or not strategy or ":" in ticker:
            body = json.dumps({
                "ok": False,
                "error": "Faltan 'ticker' / 'strategy' validos.",
            }).encode("utf-8")
            self._send_json(400, body)
            return

        cmd = [_python_executable(), "agents/agente_strategy_lab.py",
               "--promote", f"{ticker}:{strategy}"]
        if not self._acquire_lock_or_send_409("promote"):
            return
        try:
            proc = subprocess.run(
                cmd, cwd=str(ROOT), capture_output=True, text=True,
                timeout=PIPELINE_TIMEOUT_S, encoding="utf-8",
                errors="replace",
            )
            payload = {
                "ok": proc.returncode == 0,
                "ticker": ticker,
                "strategy": strategy,
                "returncode": proc.returncode,
                "stdout_tail": (proc.stdout or "").splitlines()[-8:],
                "stderr_tail": (proc.stderr or "").splitlines()[-5:],
            }
        except Exception as exc:
            payload = {"ok": False, "ticker": ticker, "error": repr(exc)}
        finally:
            _refresh_lock.release()

        # Auto-sync del Strategy Tracker en background: si el promote fue OK,
        # disparamos agente_strategy_tracker.py para que el tracker incluya la
        # nueva entrada al instante. Sin esto, la inconsistencia "Bridge dice
        # 8 promovidas / Tracker dice 7 trackeadas" persiste hasta el siguiente
        # cron. Lo lanzamos en thread daemon (fire-and-forget) para no añadir
        # latencia al boton Promote del dashboard.
        if payload.get("ok"):
            def _sync_tracker() -> None:
                try:
                    proc2 = subprocess.run(
                        [_python_executable(),
                         "agents/agente_strategy_tracker.py"],
                        cwd=str(ROOT), capture_output=True, text=True,
                        timeout=RUN_AGENT_TIMEOUT_S, encoding="utf-8",
                        errors="replace",
                    )
                    sys.stderr.write(
                        f"[promote->tracker] sync rc={proc2.returncode} for "
                        f"{ticker}:{strategy}\n"
                    )
                except Exception as exc:                  # pragma: no cover
                    sys.stderr.write(
                        f"[promote->tracker] sync failed for {ticker}: "
                        f"{exc!r}\n")
            threading.Thread(target=_sync_tracker,
                             name="tracker-sync-after-promote",
                             daemon=True).start()
            payload["tracker_sync"] = "started_in_background"

        body = json.dumps(payload, ensure_ascii=False).encode("utf-8")
        self._send_json(200 if payload.get("ok") else 500, body)

    def _handle_save_pinescript(self, parsed) -> None:
        """Guarda en agents/pending_pinescripts/<ts>_<name>.pine el codigo
        recibido en el body como texto plano. NO traduce a Python — solo
        persiste. La traduccion sigue siendo manual via chat con Claude."""
        from urllib.parse import parse_qs
        import re
        qs = parse_qs(parsed.query or "")
        raw_name = (qs.get("name", [""])[0]).strip()
        # Saneamos el nombre: solo alfanumericos/guion/subrayado, max 40 chars
        safe_name = re.sub(r"[^A-Za-z0-9_-]+", "_", raw_name)[:40].strip("_")
        if not safe_name:
            safe_name = "strategy"

        try:
            length = int(self.headers.get("Content-Length", "0") or "0")
        except ValueError:
            length = 0
        if length <= 0:
            body = json.dumps({
                "ok": False, "error": "empty_body",
                "message": "El body esta vacio — pega tu PineScript primero.",
            }, ensure_ascii=False).encode("utf-8")
            self._send_json(400, body)
            return
        # Limitamos a 256 KB por si acaso (un PineScript normal ronda los 5-20 KB)
        if length > 256 * 1024:
            body = json.dumps({
                "ok": False, "error": "too_large",
                "message": "El codigo supera 256 KB. Trocealo o pegalo en chat.",
            }, ensure_ascii=False).encode("utf-8")
            self._send_json(413, body)
            return
        try:
            raw = self.rfile.read(length)
            code = raw.decode("utf-8", errors="replace")
        except Exception as exc:
            body = json.dumps({
                "ok": False, "error": "read_error", "message": repr(exc),
            }, ensure_ascii=False).encode("utf-8")
            self._send_json(400, body)
            return

        ts = datetime.now(timezone.utc).strftime("%Y%m%dT%H%M%SZ")
        out_dir = ROOT / "agents" / "pending_pinescripts"
        out_dir.mkdir(parents=True, exist_ok=True)
        out_path = out_dir / f"{ts}_{safe_name}.pine"
        try:
            out_path.write_text(code, encoding="utf-8")
        except Exception as exc:
            body = json.dumps({
                "ok": False, "error": "write_failed", "message": repr(exc),
            }, ensure_ascii=False).encode("utf-8")
            self._send_json(500, body)
            return

        rel = out_path.relative_to(ROOT).as_posix()
        sys.stderr.write(f"[save-pinescript] saved {rel} ({length} bytes)\n")
        body = json.dumps({
            "ok": True,
            "path": rel,
            "bytes": length,
            "name": safe_name,
            "message": ("PineScript guardado. La traduccion a Python sigue "
                        "siendo manual — pega el mismo codigo en chat con "
                        "Claude para que lo integre en custom_strategies.py."),
        }, ensure_ascii=False).encode("utf-8")
        self._send_json(200, body)

    def _handle_alpaca_diagnose(self) -> None:
        """Devuelve un JSON con el estado de la conexion a Alpaca SIN exponer
        las credenciales en claro. Muestra: presencia (no valor) de cada env
        var conocida, si alpaca-py esta instalado, si el codigo nuevo (helper
        _env) esta deployado, longitudes de key/secret resueltos. Pensado para
        que el usuario diagnostique sin tocar la terminal del contenedor."""
        env_names = [
            "ALPACA_API_KEY", "ALPACA_SECRET_KEY", "ALPACA_API_SECRET",
            "APCA_API_KEY_ID", "APCA_API_SECRET_KEY", "ALPACA_BASE_URL",
        ]
        env_state: dict = {}
        for n in env_names:
            v = os.environ.get(n, "")
            if n == "ALPACA_BASE_URL":
                env_state[n] = v or "(unset — using default)"
            else:
                env_state[n] = {
                    "set": bool(v.strip()),
                    "length": len(v.strip()),
                }

        # SDK instalado? Probamos la importacion REAL — es lo unico que importa
        # (un find_spec falso-negativo me jugo una mala pasada antes).
        try:
            from alpaca.trading.client import TradingClient  # noqa: F401
            sdk_installed = True
            sdk_error = None
        except ImportError as exc:
            sdk_installed = False
            sdk_error = f"ImportError: {exc!s}"
        except Exception as exc:
            sdk_installed = False
            sdk_error = repr(exc)

        # Codigo deployado con el helper _env nuevo?
        try:
            agents_dir = ROOT / "agents"
            ac_path = agents_dir / "alpaca_client.py"
            ac_text = ac_path.read_text(encoding="utf-8") if ac_path.exists() else ""
            has_env_helper = "def _env(" in ac_text
            code_lines = ac_text.count("\n")
        except Exception as exc:
            has_env_helper = False
            code_lines = 0
            ac_text = repr(exc)

        # Resolver con la logica nueva (igual que _get_credentials, sin valores)
        sys.path.insert(0, str(ROOT / "agents"))
        try:
            import alpaca_client as _ac
            try:
                k, s, base_url, is_paper = _ac._get_credentials()
                credentials_resolved = bool(k and s)
                resolved = {
                    "key_present": bool(k),
                    "key_length": len(k) if k else 0,
                    "secret_present": bool(s),
                    "secret_length": len(s) if s else 0,
                    "base_url": base_url,
                    "is_paper": is_paper,
                }
            except Exception as exc:
                credentials_resolved = False
                resolved = {"error": repr(exc)}
        except Exception as exc:
            credentials_resolved = False
            resolved = {"error": f"import_failed: {exc!r}"}

        # Probar una llamada real al cliente, si todo lo anterior cuadra
        live_check = None
        if credentials_resolved and sdk_installed and has_env_helper:
            try:
                client, _, _, err = _ac._build_client()
                if client is None:
                    live_check = {"ok": False, "error": err}
                else:
                    try:
                        acct = client.get_account()
                        live_check = {
                            "ok": True,
                            "account_status": str(acct.status),
                            "is_paper": bool(getattr(acct, "is_paper", False)),
                        }
                    except Exception as exc:
                        live_check = {"ok": False,
                                      "error": f"get_account_failed: {exc!r}"}
            except Exception as exc:
                live_check = {"ok": False, "error": repr(exc)}

        verdict = (
            "OK - all green"
            if (credentials_resolved and sdk_installed and has_env_helper
                and live_check and live_check.get("ok"))
            else "needs_action"
        )

        payload = {
            "ok": True,
            "verdict": verdict,
            "env_vars": env_state,
            "sdk": {
                "alpaca_py_installed": sdk_installed,
                "import_error": sdk_error,
            },
            "code": {
                "has_env_helper": has_env_helper,
                "alpaca_client_lines": code_lines,
                "alpaca_client_path": str(ac_path) if 'ac_path' in locals() else None,
            },
            "credentials_resolved": credentials_resolved,
            "resolved": resolved,
            "live_check": live_check,
        }
        body = json.dumps(payload, ensure_ascii=False, indent=2).encode("utf-8")
        self._send_json(200, body)

    def _handle_promoted_diagnose(self) -> None:
        """Diagnostico del BRIDGE de estrategias promovidas, sin terminal.

        Responde dos preguntas que el dashboard no deja ver de otra forma:
          1) Se deployo el codigo nuevo? (modulo del bridge + cableado en el
             executor). Si esto sale en rojo, el push/deploy no aplico.
          2) Que ve el bridge HOY? Para cada estrategia promovida muestra si
             tiene entrada viva (trigger+regimen disparan en la ultima barra)
             o, si no, por que no. Asi distinguimos 'plano por diseno' de
             'roto'."""
        agents_dir = ROOT / "agents"
        bridge_path = agents_dir / "agente_bridge_promoted.py"
        bridge_deployed = bridge_path.exists()
        try:
            exec_text = (agents_dir / "agente_executor.py").read_text(encoding="utf-8")
            executor_wired = "build_promoted_plan_today" in exec_text
        except Exception:
            executor_wired = False

        entries: list = []
        evaluated: list = []
        note = ""
        run_error = None
        if bridge_deployed:
            sys.path.insert(0, str(agents_dir))
            try:
                import agente_bridge_promoted as _bp
                entries, evaluated, note = _bp.build_promoted_plan_today()
            except Exception as exc:
                run_error = repr(exc)

        n_live = len(entries)
        n_prom = len(evaluated)
        if not bridge_deployed or not executor_wired:
            verdict = "not_deployed"
        elif run_error:
            verdict = "run_error"
        elif n_prom == 0:
            verdict = "no_promoted"
        elif n_live == 0:
            verdict = "flat_today"
        else:
            verdict = "live_signals"

        payload = {
            "ok": True,
            "verdict": verdict,
            "code": {
                "bridge_deployed": bridge_deployed,
                "executor_wired": executor_wired,
                "bridge_path": str(bridge_path),
            },
            "run_error": run_error,
            "n_promoted": n_prom,
            "n_live": n_live,
            "summary": note,
            "live_entries": entries,
            "evaluated": evaluated,
        }
        body = json.dumps(payload, ensure_ascii=False, indent=2).encode("utf-8")
        self._send_json(200, body)

    def _handle_daily_auto_trade(self, parsed) -> None:
        """Flujo diario AUTOMATICO en una sola llamada (fire-and-forget).

        Pensado para que un scheduler (cron / n8n / tarea programada) lo dispare
        una vez cada dia de mercado. En segundo plano:
            1) Refresca datos + senales (pipeline completo) salvo ?skip_refresh=1.
            2) Corre el executor en modo --auto-approve: envia a Alpaca paper
               SOLO las entradas que el Judge marca APPROVE (incluye las
               promovidas vivas que detecta el bridge). Cero ordenes si nada
               dispara.
            3) Corre el Close Monitor en --execute --confirm para cerrar por
               Triple Barrier lo que toque.
        Deja el resultado en dashboard/data/auto_trade_last.json.

        Seguridad: si la env var LEONEX_AUTO_TOKEN esta definida, exige
        ?token=ESE_VALOR. Si no esta definida, se permite sin token (paper).
        """
        from urllib.parse import parse_qs
        qs = parse_qs(parsed.query or "")
        required = os.environ.get("LEONEX_AUTO_TOKEN", "").strip()
        if required:
            given = (qs.get("token", [""])[0]).strip()
            if given != required:
                body = json.dumps({"ok": False, "error": "bad_token"}).encode("utf-8")
                self._send_json(403, body)
                return
        skip_refresh = (qs.get("skip_refresh", ["0"])[0]).lower() in ("1", "true", "yes")
        strategy = getattr(self.server, "strategy", DEFAULT_STRATEGY)
        py = _python_executable()

        if not self._acquire_lock_or_send_409("daily_auto_trade"):
            return

        status_path = ROOT / "dashboard" / "data" / "auto_trade_last.json"

        def _write_status(d: dict) -> None:
            try:
                status_path.parent.mkdir(parents=True, exist_ok=True)
                status_path.write_text(json.dumps(d, ensure_ascii=False, indent=2),
                                       encoding="utf-8")
            except Exception as exc:
                sys.stderr.write(f"[auto-trade] no pude escribir status: {exc!r}\n")

        def _runner() -> None:
            import time as _t
            from datetime import datetime as _dt, timezone as _tz
            t0 = _t.time()
            steps: list[dict] = []
            sys.stderr.write("[auto-trade] flujo diario iniciado\n")
            _write_status({"status": "running", "started_at": _dt.now(_tz.utc).isoformat(),
                           "skip_refresh": skip_refresh})
            try:
                # 1) Refresh LIGERO: solo la cadena minima para decidir hoy.
                # NO corre los Strategy Labs / CPCV / permutation / tracker
                # (eso es investigacion, son horas de computo y no hace falta
                # cada dia). Solo: precios frescos -> regimen -> senales ->
                # triple barrier -> meta. El bridge de promovidas recompila sus
                # triggers desde la tabla prices, asi que con agente_datos ya
                # tiene datos frescos; el resto alimenta el plan legacy.
                if not skip_refresh:
                    lean_pipeline = [
                        ("datos", [py, "agents/agente_datos.py"]),
                        ("regimen", [py, "agents/agente_regimen.py"]),
                        ("senales", [py, "agents/agente_senales.py"]),
                        ("triple_barrier", [py, "agents/agente_triple_barrier.py",
                                            "--strategy", strategy]),
                        # Extended (sin --use-causal-features): mejor OOF AUC.
                        ("meta", [py, "agents/agente_meta.py", "--strategy",
                                  strategy]),
                    ]
                    for _name, _cmd in lean_pipeline:
                        try:
                            _rp = subprocess.run(
                                _cmd, cwd=str(ROOT), capture_output=True, text=True,
                                timeout=PIPELINE_TIMEOUT_S, encoding="utf-8",
                                errors="replace")
                            steps.append({
                                "step": f"refresh_{_name}",
                                "ok": _rp.returncode == 0,
                                "stderr_tail": (_rp.stderr or "").splitlines()[-5:],
                            })
                        except Exception as _exc:
                            # Un paso de refresh fallido no aborta el flujo: el
                            # executor decidira con los datos que haya.
                            steps.append({"step": f"refresh_{_name}", "ok": False,
                                          "error": repr(_exc)})
                # 2) Executor en modo auto-approve (envia paper solo APPROVE).
                # meta-mode=soft: el Meta (AUC ~0.51) no veta senales, solo
                # modula sizing. Cambiar a "gate" cuando el Meta recupere edge.
                # Con LEONEX_DRY_RUN_MODE activo se degrada a --dry-run (los
                # agentes ademas se auto-degradan leyendo la misma env var).
                if _dry_run_mode_forced():
                    exec_cmd = [py, "agents/agente_executor.py", "--strategy",
                                strategy, "--threshold", "0.55",
                                "--meta-mode", "soft", "--dry-run"]
                else:
                    exec_cmd = [py, "agents/agente_executor.py", "--strategy",
                                strategy, "--threshold", "0.55",
                                "--meta-mode", "soft", "--auto-approve"]
                ep = subprocess.run(exec_cmd, cwd=str(ROOT), capture_output=True,
                                    text=True, timeout=PIPELINE_TIMEOUT_S,
                                    encoding="utf-8", errors="replace")
                steps.append({"step": ("executor_dry_run_forced"
                                       if _dry_run_mode_forced()
                                       else "executor_auto_approve"),
                              "ok": ep.returncode == 0,
                              "stdout_tail": (ep.stdout or "").splitlines()[-25:],
                              "stderr_tail": (ep.stderr or "").splitlines()[-8:]})
                # 3) Close Monitor: cerrar por Triple Barrier lo que toque
                if _dry_run_mode_forced():
                    close_cmd = [py, "agents/agente_close_monitor.py", "--dry-run"]
                else:
                    close_cmd = [py, "agents/agente_close_monitor.py",
                                 "--execute", "--confirm"]
                cp = subprocess.run(close_cmd, cwd=str(ROOT), capture_output=True,
                                    text=True, timeout=PIPELINE_TIMEOUT_S,
                                    encoding="utf-8", errors="replace")
                steps.append({"step": "close_monitor_execute",
                              "ok": cp.returncode == 0,
                              "stdout_tail": (cp.stdout or "").splitlines()[-15:],
                              "stderr_tail": (cp.stderr or "").splitlines()[-8:]})
                elapsed = round(_t.time() - t0, 1)
                _write_status({
                    "status": "done",
                    "ok": all(s.get("ok") for s in steps),
                    "finished_at": _dt.now(_tz.utc).isoformat(),
                    "elapsed_s": elapsed,
                    "skip_refresh": skip_refresh,
                    "steps": steps,
                })
                sys.stderr.write(f"[auto-trade] flujo diario terminado en {elapsed}s\n")
            except Exception as exc:
                _write_status({"status": "error", "error": repr(exc),
                               "finished_at": _dt.now(_tz.utc).isoformat(),
                               "steps": steps})
                sys.stderr.write(f"[auto-trade] flujo diario crasheo: {exc!r}\n")
            finally:
                _refresh_lock.release()

        threading.Thread(target=_runner, name="daily-auto-trade",
                         daemon=True).start()
        body = json.dumps({
            "ok": True, "status": "started",
            "strategy": strategy, "skip_refresh": skip_refresh,
            "message": ("Flujo diario automatico lanzado en segundo plano "
                        "(refresh -> executor auto-approve -> close monitor). "
                        "El resultado queda en dashboard/data/auto_trade_last.json."),
        }, ensure_ascii=False).encode("utf-8")
        self._send_json(202, body)

    def _handle_refresh_trades(self) -> None:
        """Refresh RAPIDO solo del estado de posiciones / cierres, sin tocar
        descargas de datos ni labs. Ejecuta en serie:

            1) agente_executor.py --dry-run    (lee Alpaca, escribe paper_report.json)
            2) agente_close_monitor.py --dry-run (re-evalua TP/SL contra quotes
               vivos de Alpaca, escribe close_monitor_report.json)

        Total ~20-30s en vez de los ~50min del pipeline entero. Para el caso de
        uso "tengo posiciones abiertas y quiero ver si han tocado TP/SL ya"
        sin esperar al cron / a la SSE / a re-bajar el S&P 500."""
        strategy = getattr(self.server, "strategy", DEFAULT_STRATEGY)
        if not self._acquire_lock_or_send_409("refresh_trades"):
            return

        py = _python_executable()
        steps: list[dict] = []
        overall_ok = True
        import time as _t
        try:
            for label, cmd in [
                ("executor (dry-run)", [py, "agents/agente_executor.py",
                                        "--strategy", strategy, "--dry-run"]),
                ("close_monitor (dry-run)", [py, "agents/agente_close_monitor.py",
                                             "--dry-run"]),
            ]:
                t0 = _t.time()
                step = {"step": label, "cmd": cmd}
                try:
                    proc = subprocess.run(
                        cmd, cwd=str(ROOT), capture_output=True, text=True,
                        timeout=RUN_AGENT_TIMEOUT_S, encoding="utf-8",
                        errors="replace",
                    )
                    step.update({
                        "ok": proc.returncode == 0,
                        "returncode": proc.returncode,
                        "duration": round(_t.time() - t0, 2),
                        "stdout_tail": (proc.stdout or "").splitlines()[-5:],
                        "stderr_tail": (proc.stderr or "").splitlines()[-3:],
                    })
                except subprocess.TimeoutExpired:
                    step.update({
                        "ok": False, "duration": round(_t.time() - t0, 2),
                        "error": f"timeout_after_{RUN_AGENT_TIMEOUT_S}s",
                    })
                except Exception as exc:
                    step.update({
                        "ok": False, "duration": round(_t.time() - t0, 2),
                        "error": repr(exc),
                    })
                steps.append(step)
                if not step.get("ok"):
                    overall_ok = False
                    break
        finally:
            _refresh_lock.release()

        payload = {
            "ok": overall_ok,
            "scope": "trades_only",
            "steps": steps,
            "total_duration": round(sum(s.get("duration", 0.0) for s in steps), 2),
        }
        body = json.dumps(payload, ensure_ascii=False).encode("utf-8")
        self._send_json(200 if overall_ok else 500, body)

    def _handle_bridge_intraday(self, parsed) -> None:
        """Bridge INTRADIA: re-evalua las estrategias promovidas y, si alguna
        dispara su trigger, manda orden a Alpaca paper. Por defecto NO re-
        descarga datos para que el cron de n8n termine rapido (~30-60s) y no
        bloquee el lock.

        Pensado para que un cron de n8n lo llame cada 15-60 min en horario US
        market (13:30-20:00 UTC L-V). Como el endpoint termina en <90s, el
        timeout HTTP del cron no se dispara y el lock se libera enseguida.

        Query params:
            ?mode=dry_run (default) | execute
                dry_run solo planea; execute envia ordenes a Alpaca paper.
            ?fetch=0 (default) | 1
                fetch=0: SOLO ejecutor + close_monitor. ~30-60s. Usa datos
                  intraday del cache (ultimo cron diario). RECOMENDADO para
                  cron frecuente.
                fetch=1: ademas descarga datos intraday incrementales (1h/4h)
                  antes de evaluar. ~3-8 min. Usa esto cuando los datos del
                  cache estan demasiado viejos.

        Pasos en serie (no se solapan dos ejecuciones a la vez):
            [opcional, solo si fetch=1] intraday data 1h/4h
            executor (con bridge)
            close_monitor

        En el endpoint anterior se incluia tambien una descarga full de 30m/
        15m/5m que tardaba >10 min y colgaba el lock. Eso se elimino: el cron
        diario sigue actualizando los TFs scalping; el bridge no los necesita
        frescos hasta el siguiente cron diario."""
        from urllib.parse import parse_qs
        qs = parse_qs(parsed.query or "")
        mode = (qs.get("mode", ["dry_run"])[0]).lower()
        if mode not in ("dry_run", "execute"):
            body = json.dumps({
                "ok": False, "error": f"mode_invalido: {mode}",
                "message": "Usa mode=dry_run o mode=execute",
            }, ensure_ascii=False).encode("utf-8")
            self._send_json(400, body)
            return
        # LEONEX_DRY_RUN_MODE gana siempre, incluso a un ?mode=execute del
        # cron de n8n: sin esto el modo observacion seria decorativo.
        if mode == "execute" and _dry_run_mode_forced():
            mode = "dry_run"
        fetch = (qs.get("fetch", ["0"])[0]).strip() in ("1", "true", "yes")

        # Auto-promote fetch=0 -> fetch=1 si el cache 1h esta rancio (last bar
        # de antes de hoy UTC). Esto evita que el bridge evalue con datos del
        # dia anterior cuando el daily corrio pre-market (08:00 NY) y aun no
        # tenia data del dia actual. Solo aplica en market hours US (L-V).
        fetch_was_forced = False
        if not fetch and _within_market_hours():
            try:
                import sqlite3 as _sql
                db_file = ROOT / "data" / "Leonex.sqlite"
                if db_file.exists():
                    with _sql.connect(str(db_file)) as conn:
                        row = conn.execute(
                            "SELECT MAX(ts) FROM prices_intraday "
                            "WHERE timeframe = '1h'"
                        ).fetchone()
                        last_ts = row[0] if row and row[0] else None
                if last_ts:
                    # last_ts es ISO string. Si es < hoy UTC, force fetch=1.
                    today_utc = datetime.now(timezone.utc).strftime("%Y-%m-%d")
                    if last_ts[:10] < today_utc:
                        fetch = True
                        fetch_was_forced = True
            except Exception:
                pass  # si falla la deteccion, seguimos con fetch=0

        strategy = getattr(self.server, "strategy", DEFAULT_STRATEGY)
        if not self._acquire_lock_or_send_409("bridge_intraday"):
            return

        py = _python_executable()
        # Mode del executor + close_monitor: ambos respetan el ?mode= del query.
        # Si el cron de n8n llama ?mode=execute -> envia ordenes Y cierra TP/SL.
        # Si llama ?mode=dry_run -> solo simula ambos. Esto evita el bug donde
        # el bridge detectaba SL hit pero nunca cerraba (close_monitor estaba
        # hardcoded a --dry-run independientemente del mode del executor).
        if mode == "execute":
            executor_args = ["--strategy", strategy, "--execute", "--confirm",
                             "--auto-approve"]
            close_args = ["--execute", "--confirm"]
            close_label = "close_monitor (execute)"
        else:
            executor_args = ["--strategy", strategy, "--dry-run"]
            close_args = ["--dry-run"]
            close_label = "close_monitor (dry-run)"

        steps_def: list[tuple[str, list[str]]] = []
        if fetch:
            # Descarga incremental ligera (1h/4h, no scalping pesado)
            steps_def.append((
                "intraday data swing (1h/4h, incremental)",
                [py, "agents/agente_datos_intraday.py"],
            ))
        steps_def.extend([
            (f"executor ({mode}) + bridge",
             [py, "agents/agente_executor.py", *executor_args]),
            (close_label,
             [py, "agents/agente_close_monitor.py", *close_args]),
        ])
        # En modo execute anadimos risk_monitor al final para revalidar y
        # auto-limpiar el pause si los near_stops ya cerraron.
        if mode == "execute":
            steps_def.append((
                "risk_monitor (revalidate pause)",
                [py, "agents/agente_risk_monitor.py"],
            ))

        steps: list[dict] = []
        overall_ok = True
        import time as _t
        try:
            for label, cmd in steps_def:
                t0 = _t.time()
                step = {"step": label, "cmd": cmd}
                try:
                    proc = subprocess.run(
                        cmd, cwd=str(ROOT), capture_output=True, text=True,
                        timeout=PIPELINE_TIMEOUT_S, encoding="utf-8",
                        errors="replace",
                    )
                    step.update({
                        "ok": proc.returncode == 0,
                        "returncode": proc.returncode,
                        "duration": round(_t.time() - t0, 2),
                        "stdout_tail": (proc.stdout or "").splitlines()[-5:],
                        "stderr_tail": (proc.stderr or "").splitlines()[-3:],
                    })
                except subprocess.TimeoutExpired:
                    step.update({
                        "ok": False, "duration": round(_t.time() - t0, 2),
                        "error": f"timeout_after_{PIPELINE_TIMEOUT_S}s",
                    })
                except Exception as exc:
                    step.update({
                        "ok": False, "duration": round(_t.time() - t0, 2),
                        "error": repr(exc),
                    })
                steps.append(step)
                # Si falla la descarga de datos no seguimos al bridge (datos
                # rotos = bridge se evalua contra basura). Si falla el executor
                # paramos para no seguir al close_monitor sin estado fresco.
                if not step.get("ok"):
                    overall_ok = False
                    break
        finally:
            _refresh_lock.release()

        payload = {
            "ok": overall_ok,
            "scope": "intraday_bridge",
            "mode": mode,
            "strategy": strategy,
            "fetch": fetch,
            "fetch_was_forced": fetch_was_forced,
            "steps": steps,
            "total_duration": round(sum(s.get("duration", 0.0) for s in steps), 2),
        }
        body = json.dumps(payload, ensure_ascii=False).encode("utf-8")
        self._send_json(200 if overall_ok else 500, body)

    def _handle_close_position(self, parsed) -> None:
        """Cierra UNA posicion abierta especifica en Alpaca paper sin tocar
        las demas. Util cuando tienes varias y quieres cerrar solo una.

        Pasos:
            1. Resolver ticker → simbolo Alpaca (cripto: BTC-USD → BTC/USD).
            2. Confirmar que existe posicion abierta para ese simbolo.
            3. client.close_position(symbol) — Alpaca cierra TODA la cantidad
               de esa posicion con una orden market en el lado contrario.
            4. Marcar todos los open_trades de ese ticker en BD como
               status=closed_manual para que el Close Monitor no intente
               re-cerrar algo que ya esta cerrado.

        Pensado para invocarse desde un boton "Close" por posicion en el
        dashboard."""
        from urllib.parse import parse_qs
        qs = parse_qs(parsed.query or "")
        ticker = (qs.get("ticker", [""])[0]).strip()
        if not ticker:
            body = json.dumps({
                "ok": False, "error": "missing_ticker",
                "message": "Pasa ?ticker=SYM en la query (e.g. ?ticker=STT).",
            }, ensure_ascii=False).encode("utf-8")
            self._send_json(400, body)
            return

        if not self._acquire_lock_or_send_409("close_position"):
            return

        try:
            sys.path.insert(0, str(ROOT / "agents"))
            try:
                from alpaca_client import _build_client, to_alpaca_symbol
            except Exception as exc:
                body = json.dumps({
                    "ok": False, "error": "import_failed",
                    "message": repr(exc),
                }, ensure_ascii=False).encode("utf-8")
                self._send_json(500, body)
                return

            client, _, _, err = _build_client()
            if client is None:
                body = json.dumps({
                    "ok": False, "error": "alpaca_not_configured",
                    "message": str(err or "_build_client returned None"),
                }, ensure_ascii=False).encode("utf-8")
                self._send_json(500, body)
                return

            alpaca_sym = to_alpaca_symbol(ticker)

            # 1) Verificar posicion
            try:
                pos = client.get_open_position(alpaca_sym)
                qty = float(getattr(pos, "qty", 0) or 0)
            except Exception as exc:
                body = json.dumps({
                    "ok": False, "error": "position_not_found",
                    "ticker": ticker,
                    "message": f"Sin posicion abierta para {ticker} en Alpaca: {exc!r}",
                }, ensure_ascii=False).encode("utf-8")
                self._send_json(404, body)
                return

            # 2) Cerrar via Alpaca
            try:
                close_order = client.close_position(alpaca_sym)
            except Exception as exc:
                body = json.dumps({
                    "ok": False, "error": "close_failed",
                    "ticker": ticker, "message": repr(exc),
                }, ensure_ascii=False).encode("utf-8")
                self._send_json(500, body)
                return

            order_id = str(getattr(close_order, "id", "") or "") or None
            order_status = str(getattr(close_order, "status", "") or "") or None

            # 3) Marcar open_trades en BD como closed_manual
            import sqlite3
            from datetime import datetime as _dt
            marked = 0
            db_file = ROOT / "data" / "Leonex.sqlite"
            if db_file.exists():
                try:
                    with sqlite3.connect(str(db_file)) as conn:
                        now_iso = _dt.now(timezone.utc).isoformat()
                        cur = conn.execute(
                            "UPDATE open_trades "
                            "   SET status = 'closed_manual', "
                            "       exit_date = ?, "
                            "       alpaca_sell_order_id = ?, "
                            "       closed_at = ? "
                            " WHERE ticker = ? AND status = 'open'",
                            (now_iso, order_id, now_iso, ticker)
                        )
                        marked = cur.rowcount
                        conn.commit()
                except Exception as exc:
                    sys.stderr.write(
                        f"[close-position] no pude marcar open_trades "
                        f"para {ticker}: {exc!r}\n"
                    )

            sys.stderr.write(
                f"[close-position] {ticker} qty={qty} order_id={order_id} "
                f"status={order_status} open_trades_marked={marked}\n"
            )
            body = json.dumps({
                "ok": True,
                "ticker": ticker,
                "alpaca_symbol": alpaca_sym,
                "qty_closed": qty,
                "order_id": order_id,
                "order_status": order_status,
                "open_trades_marked_closed_manual": marked,
                "message": (
                    f"Posicion {ticker} cerrada en Alpaca paper "
                    f"(qty={qty}). {marked} open_trades marcado(s) en BD."
                ),
            }, ensure_ascii=False).encode("utf-8")
            self._send_json(200, body)
        finally:
            _refresh_lock.release()

    def _send_json(self, status: int, body: bytes) -> None:
        self.send_response(status)
        self.send_header("Content-Type", "application/json; charset=utf-8")
        self.send_header("Content-Length", str(len(body)))
        self.end_headers()
        self.wfile.write(body)

    def _acquire_lock_or_send_409(self, holder: str) -> bool:
        """Intenta tomar el _refresh_lock para una accion. Si esta ocupado,
        envia respuesta 409 con info del holder actual + edad y devuelve
        False. Asi cada handler responde uniforme y el dashboard ve quien
        bloquea (y desde hace cuanto) cuando ofrece liberar el lock."""
        if _refresh_lock.acquire(blocking=False, holder=holder):
            return True
        s = _refresh_lock.status()
        msg_holder = s.get("holder") or "unknown"
        age = s.get("age_seconds")
        age_txt = f" (hace {age}s)" if age is not None else ""
        body = json.dumps({
            "ok": False, "error": "refresh_in_progress",
            "message": (f"Hay otra accion en curso: {msg_holder}{age_txt}. "
                        "Espera a que termine o, si crees que esta atascado, "
                        "usa el boton 'Release stuck lock' del dashboard "
                        "(POST /api/release-lock)."),
            "lock_status": s,
        }, ensure_ascii=False).encode("utf-8")
        self._send_json(409, body)
        return False


class LeonexServer(ThreadingHTTPServer):
    daemon_threads = True
    allow_reuse_address = True

    def __init__(self, addr, handler, strategy: str):
        super().__init__(addr, handler)
        self.strategy = strategy


def _port_in_use(port: int, host: str = "127.0.0.1") -> bool:
    with socket.socket(socket.AF_INET, socket.SOCK_STREAM) as s:
        return s.connect_ex((host, port)) == 0


def main() -> int:
    parser = argparse.ArgumentParser(description="Servidor local de Leonex")
    parser.add_argument("--port", type=int, default=DEFAULT_PORT)
    parser.add_argument("--host", type=str, default="127.0.0.1",
                        help="Bind address. Usa 0.0.0.0 dentro de contenedor "
                             "para que el proxy externo pueda alcanzarlo.")
    parser.add_argument("--strategy", type=str, default=DEFAULT_STRATEGY)
    args = parser.parse_args()

    # Cambiar al directorio raiz del proyecto para que rutas relativas en los
    # handlers (dashboard/, data/, etc.) funcionen.
    os.chdir(ROOT)

    print("=" * 78)
    print(f"Leonex Server escuchando en {args.host}:{args.port} (all interfaces)")
    print(f"Dashboard local      : http://localhost:{args.port}/dashboard/index.html")
    print(f"Estrategia pipeline  : {args.strategy}")
    print(f"Endpoint refresh     : POST /api/refresh  (sincrono, bloquea ~40min)")
    print(f"Endpoint async       : POST /api/refresh-async  (para n8n cron)")
    print(f"Endpoint lock-status : GET  /api/lock-status")
    print(f"Endpoint release-lock: POST /api/release-lock")
    print(f"Endpoint test-bridge : GET  /api/test-custom-bridge")
    _sched_mode = ("dry_run (FORZADO por LEONEX_DRY_RUN_MODE)"
                   if _dry_run_mode_forced() else "execute")
    print(f"Scheduler intradia   : bridge cada {INTRADAY_BRIDGE_INTERVAL_S // 60} min "
          f"en L-V {MARKET_OPEN_UTC[0]:02d}:{MARKET_OPEN_UTC[1]:02d}-"
          f"{MARKET_CLOSE_UTC[0]:02d}:{MARKET_CLOSE_UTC[1]:02d} UTC (mode={_sched_mode})")
    print(f"Python interprete    : {sys.executable}")
    if _dry_run_mode_forced():
        print("=" * 78)
        print("MODO OBSERVACION GLOBAL ACTIVO (LEONEX_DRY_RUN_MODE=1)")
        print("Pipeline y journal funcionan, pero NINGUNA orden automatica saldra")
        print("hacia Alpaca. Cierres manuales (boton Close / reset) siguen operativos.")
        print("Quita la env var y redeploya para volver a execute.")
    print("=" * 78)
    sys.stdout.flush()

    # Scheduler interno: dispara el bridge intradia cada 15 min en horario de
    # mercado. Daemon → muere con el proceso; no bloquea el shutdown.
    threading.Thread(
        target=_intraday_bridge_scheduler, args=(args.strategy,),
        name="intraday-bridge-scheduler", daemon=True,
    ).start()

    server = LeonexServer((args.host, args.port), LeonexHandler, args.strategy)
    try:
        server.serve_forever()
    except KeyboardInterrupt:
        print("\nServer detenido por KeyboardInterrupt.")
    finally:
        server.server_close()
    return 0


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