import asyncio
import functools
import hmac
import html
import logging
import os
import secrets
import signal
import time
import urllib.parse
from datetime import timedelta

import psutil
from flask import Flask, g, jsonify, redirect, render_template, request, session

import config
from core import ForwarderCore
from db import Database

# ── Logging ─────────────────────────────────────────────────────────────────
log_format = "%(asctime)s %(levelname)s %(name)s: %(message)s"
logging.basicConfig(level=logging.INFO, format=log_format)

log_file = os.path.join(
    os.path.dirname(os.path.abspath(__file__)), "logs", "app.log"
)
try:
    os.makedirs(os.path.dirname(log_file), exist_ok=True)
    from logging.handlers import RotatingFileHandler

    fh = RotatingFileHandler(log_file, maxBytes=1 * 1024 * 1024, backupCount=2)
    fh.setFormatter(logging.Formatter(log_format))
    logging.getLogger().addHandler(fh)
except Exception:
    pass

logger = logging.getLogger(__name__)

# ── Flask app ───────────────────────────────────────────────────────────────
app = Flask(__name__)
app.secret_key = config.FLASK_SECRET_KEY
app.config.update(
    SESSION_COOKIE_HTTPONLY=True,
    SESSION_COOKIE_SECURE=True,
    SESSION_COOKIE_SAMESITE="Lax",
    PERMANENT_SESSION_LIFETIME=timedelta(hours=12),
)

# Global singletons
db = Database()
core = ForwarderCore(db)

# Override env defaults with DB-persisted config.
# Import-time load is best-effort: if MariaDB is briefly down (shared-hosting
# blips, restart race), boot on .env defaults rather than crashing the whole
# app. ensure_core() retries initialization on later requests, and other
# handlers re-patch config once the DB is reachable again.
try:
    config.patch_from_db(db)
except Exception as exc:
    logger.warning("DB config load failed at startup; using .env defaults: %s", exc)

_core_started = False


# ── Security headers ────────────────────────────────────────────────────────
@app.after_request
def add_security_headers(response):
    response.headers["X-Content-Type-Options"] = "nosniff"
    response.headers["X-Frame-Options"] = "DENY"
    response.headers["Content-Security-Policy"] = (
        "default-src 'self'; "
        "script-src 'self' 'unsafe-inline' 'unsafe-eval' https://unpkg.com https://cdn.jsdelivr.net; "
        "style-src 'self' 'unsafe-inline' https://fonts.googleapis.com; "
        "font-src 'self' https://fonts.gstatic.com; "
        "img-src 'self' data:; "
        "connect-src 'self'"
    )
    response.headers["Referrer-Policy"] = "strict-origin-when-cross-origin"
    response.headers["Permissions-Policy"] = "geolocation=(), microphone=(), camera=()"
    if request.is_secure or request.headers.get("X-Forwarded-Proto") == "https":
        response.headers["Strict-Transport-Security"] = "max-age=31536000; includeSubDomains"
    return response


def _is_https() -> bool:
    return request.is_secure or request.headers.get("X-Forwarded-Proto") == "https"


@app.before_request
def force_https():
    """Redirect plain-HTTP requests to HTTPS so secure cookies work."""
    if _is_https():
        return None
    host = request.host.split(":")[0]
    if host in ("127.0.0.1", "localhost", "[::1]"):
        return None
    # Health checks may still be probed over HTTP by some monitors.
    if request.path.rstrip("/") == "/api/health":
        return None
    url = request.url.replace("http://", "https://", 1)
    return redirect(url, code=301)


# ── CSRF helpers ────────────────────────────────────────────────────────────
def generate_csrf() -> str:
    if "_csrf" not in session:
        session["_csrf"] = secrets.token_urlsafe(32)
    return session["_csrf"]


def csrf_required(f):
    @functools.wraps(f)
    def wrapper(*args, **kwargs):
        if request.method in ("GET", "HEAD", "OPTIONS"):
            return f(*args, **kwargs)
        token = request.headers.get("X-CSRF-Token")
        if not token and request.is_json:
            token = (request.get_json(silent=True) or {}).get("csrf_token")
        if not token or not hmac.compare_digest(token, session.get("_csrf", "")):
            return jsonify({"error": "Invalid or missing CSRF token"}), 403
        return f(*args, **kwargs)

    return wrapper


# ── Rate limiting (in-memory per-IP) ────────────────────────────────────────
_rate_limit_store: dict = {}


def rate_limit(max_requests: int = 5, window: int = 60):
    def decorator(f):
        @functools.wraps(f)
        def wrapper(*args, **kwargs):
            ip = request.remote_addr or "unknown"
            now = time.time()
            entries = [t for t in _rate_limit_store.get(ip, []) if now - t < window]
            if len(entries) >= max_requests:
                return jsonify({"error": "Rate limit exceeded"}), 429
            entries.append(now)
            _rate_limit_store[ip] = entries
            return f(*args, **kwargs)

        return wrapper

    return decorator


# ── PIN authentication ──────────────────────────────────────────────────────
def pin_required(f):
    @functools.wraps(f)
    def wrapper(*args, **kwargs):
        if not config.WEB_PIN:
            return f(*args, **kwargs)
        if not session.get("pin_authenticated"):
            return jsonify({"error": "PIN required"}), 401
        return f(*args, **kwargs)

    return wrapper


# ── Core lifecycle ──────────────────────────────────────────────────────────
@app.before_request
def ensure_core():
    global _core_started
    if not config.AUTO_START_CORE:
        if not getattr(ensure_core, "_logged_skip", False):
            ensure_core._logged_skip = True
            logger.info(
                "AUTO_START_CORE is disabled — Telegram core will not start automatically"
            )
        return
    if not _core_started:
        core.start()
        _core_started = True
        return
    # After core started: if no clients loaded after a grace period and
    # there are authorized accounts in the DB, retry initialization.
    # The initial _init_accounts can fail if Telegram is unreachable at
    # startup; this gives it another chance on subsequent requests.
    if not core.clients and not getattr(core, "_retrying", False) and not getattr(core, "_init_started", False):
        import time as _time
        if not hasattr(core, "_first_request_time"):
            core._first_request_time = _time.time()
        # Init is now bounded by asyncio.wait_for timeouts (~45s worst case
        # per account), so a 15s grace is enough for the first attempt to
        # either succeed or bail out before we schedule a retry.
        if _time.time() - core._first_request_time < 15:
            return
        try:
            accounts = core._run_db_sync(db.get_accounts, timeout=5)
            if any(a.get("is_authorized") for a in accounts):
                core._retrying = True
                logger.info("Retrying account initialization (clients=0 but authorized accounts in DB)")
                asyncio.run_coroutine_threadsafe(core._init_accounts(), core.loop)
        except Exception:
            pass


def _sigterm_handler(signum, frame):
    logger.info("SIGTERM received, shutting down...")
    try:
        core.shutdown()
    except Exception as e:
        logger.warning("Shutdown error: %s", e)
    os._exit(0)


signal.signal(signal.SIGTERM, _sigterm_handler)


# ── Web UI ──────────────────────────────────────────────────────────────────
@app.route("/")
def index():
    status = core.status()
    try:
        rss = psutil.Process().memory_info().rss
        pairs_count = len(status.get("pairs", []))
    except Exception:
        rss = 0
        pairs_count = 0
    health = {"pairs_count": pairs_count, "rss_mb": round(rss / 1024 / 1024, 2)}
    pin_required = bool(config.WEB_PIN) and not session.get("pin_authenticated")
    return render_template(
        "index.html",
        csrf=generate_csrf(),
        status=status,
        health=health,
        pin_required=pin_required,
        core_started=bool(_core_started),
        auto_start_core=bool(config.AUTO_START_CORE),
    )


def _ctrader_status() -> dict:
    client = getattr(core, "_ctrader_client", None)
    if not client:
        return {"started": False}
    tick = client.get_tick(config.CTRADER_SYMBOL)
    latest_tick = None
    if tick:
        latest_tick = {
            "symbol_id": tick.symbol_id,
            "symbol_name": tick.symbol_name,
            "bid": tick.bid,
            "ask": tick.ask,
            "timestamp_ms": tick.timestamp_ms,
        }
    return {
        "started": getattr(core, "_ctrader_started", False),
        "connected": getattr(client, "_connected", False),
        "symbol": config.CTRADER_SYMBOL,
        "live": config.CTRADER_USE_LIVE,
        "latest_tick": latest_tick,
    }


# ── Favicon ─────────────────────────────────────────────────────────────────
@app.route("/favicon.ico")
def favicon():
    """Silence browser requests for a missing favicon."""
    return "", 204


# ── Public health ───────────────────────────────────────────────────────────
@app.route("/api/health")
def api_health():
    status = core.status()
    rss = psutil.Process().memory_info().rss
    return jsonify(
        {
            "status": "ok",
            "authorized": status.get("authorized"),
            "forwarding": status.get("forwarding"),
            "pairs_count": len(status.get("pairs", [])),
            "rss_mb": round(rss / 1024 / 1024, 2),
            "ctrader": _ctrader_status(),
        }
    )


# ── Real-time price (cached tick) ───────────────────────────────────────────
@app.route("/api/price")
def api_price():
    """Return the latest cached tick from the cTrader spot stream.

    Query params:
        symbol — optional, defaults to CTRADER_SYMBOL
    """
    symbol = request.args.get("symbol", config.CTRADER_SYMBOL).upper()
    client = getattr(core, "_ctrader_client", None)
    if not client:
        return jsonify({"error": "cTrader spot client not started"}), 503
    tick = client.get_tick(symbol)
    if not tick:
        return jsonify({"error": "No tick cached yet", "symbol": symbol}), 503
    return jsonify(
        {
            "status": "ok",
            "symbol": tick.symbol_name,
            "symbol_id": tick.symbol_id,
            "bid": tick.bid,
            "ask": tick.ask,
            "spread": round(tick.ask - tick.bid, 5),
            "timestamp_ms": tick.timestamp_ms,
            "source": "cache",
        }
    )


# ── Signal price-augmentation API ───────────────────────────────────────────
@app.route("/api/signals", methods=["GET"])
@pin_required
def api_signals():
    source_chat_id = request.args.get("source_chat_id", type=int)
    since_id = request.args.get("since_id", type=int)
    if source_chat_id:
        rows = db.get_signal_snapshots(source_chat_id, limit=50, since_id=since_id)
    else:
        rows = db.get_all_signal_snapshots(limit=50, since_id=since_id)
    return jsonify({"status": "ok", "signals": rows})


# ── PIN auth routes ─────────────────────────────────────────────────────────
@app.route("/api/auth/status")
def api_auth_status():
    return jsonify(
        {
            "authenticated": bool(session.get("pin_authenticated"))
            or not config.WEB_PIN,
            "pin_required": bool(config.WEB_PIN),
            "csrf_token": generate_csrf()
            if (config.WEB_PIN and session.get("pin_authenticated"))
            else None,
        }
    )


@app.route("/api/auth/login", methods=["POST"])
@rate_limit(max_requests=5, window=60)
def api_login():
    if not config.WEB_PIN:
        return jsonify({"status": "ok", "csrf_token": generate_csrf()})
    data = request.get_json(force=True, silent=True) or {}
    pin = str(data.get("pin", "")).strip()
    if not pin:
        return jsonify({"error": "PIN required"}), 400
    if pin != config.WEB_PIN:
        return jsonify({"error": "Invalid PIN"}), 403
    session["pin_authenticated"] = True
    return jsonify({"status": "ok", "csrf_token": generate_csrf()})


@app.route("/api/auth/logout", methods=["POST"])
@csrf_required
def api_logout():
    session.clear()
    return jsonify({"status": "ok"})


# ── Telegram authentication (account-scoped; see /api/accounts/<id>/... ) ──
@app.route("/api/config")
@pin_required
def api_config():
    return jsonify(
        {
            "phone": config.TG_PHONE or "",
            "has_password": bool(config.TG_PASSWORD),
            "api_configured": bool(config.TG_API_ID and config.TG_API_HASH),
        }
    )


@app.route("/api/config/telegram/test", methods=["POST"])
@pin_required
@csrf_required
def api_test_telegram():
    """Validate TG_API_ID / TG_API_HASH by issuing a lightweight Telegram call."""
    try:
        return jsonify(core.test_telegram_api())
    except Exception as e:
        logger.exception("test_telegram failed")
        return jsonify({"ok": False, "error": str(e)}), 500


@app.route("/api/config/signals", methods=["GET"])
@pin_required
def api_signals_config():
    """Read current signal-processing config (webhook, bot, etc.)."""
    try:
        url = core._run_db_sync(db.get_config, "SIGNAL_WEBHOOK_URL")
        secret = core._run_db_sync(db.get_config, "SIGNAL_WEBHOOK_SECRET")
        chat_id = core._run_db_sync(db.get_config, "SIGNAL_ADMIN_CHAT_ID")
        return jsonify(
            {
                "status": "ok",
                "webhook_url": url or config.SIGNAL_WEBHOOK_URL,
                "webhook_secret": bool(secret or config.SIGNAL_WEBHOOK_SECRET),
                "admin_chat_id": int(chat_id) if chat_id else config.SIGNAL_ADMIN_CHAT_ID,
            }
        )
    except Exception as exc:
        logger.exception("signals_config GET failed")
        return jsonify({"error": str(exc)}), 500


@app.route("/api/config/signals", methods=["POST"])
@pin_required
@csrf_required
def api_signals_config_update():
    """Save signal-processing config to DB and update runtime."""
    data = request.get_json(force=True, silent=True) or {}
    url = data.get("webhook_url", "").strip()
    secret = data.get("webhook_secret", "").strip()
    chat_id_raw = data.get("admin_chat_id")

    try:
        core._run_db_sync(db.set_config, "SIGNAL_WEBHOOK_URL", url)
        core._run_db_sync(db.set_secret_config, "SIGNAL_WEBHOOK_SECRET", secret)
        chat_id_val = str(int(chat_id_raw)) if chat_id_raw is not None else ""
        core._run_db_sync(db.set_config, "SIGNAL_ADMIN_CHAT_ID", chat_id_val)
    except Exception as exc:
        logger.exception("signals_config DB save failed")
        return jsonify({"error": str(exc)}), 500

    # Update module-level defaults so next restart picks them up
    config.SIGNAL_WEBHOOK_URL = url
    config.SIGNAL_WEBHOOK_SECRET = secret
    config.SIGNAL_ADMIN_CHAT_ID = int(chat_id_val) if chat_id_val else 0

    # Live-update the running SignalProcessor if it exists
    sp = getattr(core, "_signal_processor", None)
    if sp:
        sp.update_config(
            webhook_url=url,
            webhook_secret=secret,
            admin_chat_id=int(chat_id_val) if chat_id_val else None,
        )
        logger.info("SignalProcessor config updated live")
    else:
        # Trigger _ensure_ctrader to create processor if a pair now needs it
        try:
            asyncio.run_coroutine_threadsafe(
                core._ensure_ctrader(), core.loop
            )
        except Exception as exc:
            logger.warning("_ensure_ctrader trigger failed: %s", exc)

    return jsonify({"status": "ok"})


# ── Optional OpenAI-compatible LLM classification config ───────────────────
@app.route("/api/config/llm", methods=["GET"])
@pin_required
def api_llm_config():
    """Read current LLM config. Never returns the plaintext API key/secret."""
    try:
        st = core.llm_status()
        webhook_url = core._run_db_sync(db.get_config, "LLM_WEBHOOK_URL") or ""
        webhook_secret = core._run_db_sync(db.get_secret_config, "LLM_WEBHOOK_SECRET")
        return jsonify(
            {
                "status": "ok",
                "enabled": bool(st.get("enabled")),
                "provider": st.get("provider", "openai-compatible"),
                "base_url": st.get("base_url", "https://api.cerebras.ai/v1"),
                "model": st.get("model", "llama3.1-8b"),
                "api_key_set": bool(st.get("api_key_set")),
                "client_available": bool(st.get("client_available", True)),
                "available": bool(st.get("available")),
                "min_confidence": st.get("min_confidence", 0.80),
                "timeout": st.get("timeout", 12.0),
                "max_tokens": st.get("max_tokens", 120),
                "webhook_url": webhook_url,
                "webhook_secret_set": bool(webhook_secret),
            }
        )
    except Exception as exc:
        logger.exception("llm_config GET failed")
        return jsonify({"error": str(exc)}), 500


@app.route("/api/config/llm", methods=["POST"])
@pin_required
@csrf_required
def api_llm_config_update():
    """Save LLM config to DB (encrypting the API key) and reload the engine."""
    data = request.get_json(force=True, silent=True) or {}
    enabled = bool(data.get("enabled", False))
    provider = str(data.get("provider", "openai-compatible") or "openai-compatible").strip().lower()
    base_url = str(data.get("base_url", "") or "").strip()
    model = str(data.get("model", "llama3.1-8b") or "llama3.1-8b").strip()
    api_key = str(data.get("api_key", "") or "").strip()
    webhook_url = str(data.get("webhook_url", "") or "").strip()
    webhook_secret = str(data.get("webhook_secret", "") or "").strip()
    try:
        min_confidence = float(data.get("min_confidence", 0.80))
    except (TypeError, ValueError):
        min_confidence = 0.80
    try:
        timeout = float(data.get("timeout", 12.0))
    except (TypeError, ValueError):
        timeout = 12.0
    try:
        max_tokens = int(data.get("max_tokens", 120))
    except (TypeError, ValueError):
        max_tokens = 120

    try:
        core._run_db_sync(db.set_config, "LLM_ENABLED", "1" if enabled else "0")
        core._run_db_sync(db.set_config, "LLM_PROVIDER", provider)
        core._run_db_sync(db.set_config, "LLM_BASE_URL", base_url or "https://api.cerebras.ai/v1")
        core._run_db_sync(db.set_config, "LLM_MODEL", model)
        core._run_db_sync(db.set_config, "LLM_MIN_CONFIDENCE", str(min_confidence))
        core._run_db_sync(db.set_config, "LLM_TIMEOUT", str(timeout))
        core._run_db_sync(db.set_config, "LLM_MAX_TOKENS", str(max_tokens))
        core._run_db_sync(db.set_config, "LLM_WEBHOOK_URL", webhook_url)
        # Only overwrite secrets when a new value is supplied (empty = keep current).
        if api_key:
            core._run_db_sync(db.set_secret_config, "LLM_API_KEY", api_key)
        if webhook_secret:
            core._run_db_sync(db.set_secret_config, "LLM_WEBHOOK_SECRET", webhook_secret)
    except Exception as exc:
        logger.exception("llm_config DB save failed")
        return jsonify({"error": str(exc)}), 500

    try:
        status = core.reload_llm_config()
    except Exception as exc:
        logger.warning("llm reload failed: %s", exc)
        status = core.llm_status()
    return jsonify({"status": "ok", "engine": status})


@app.route("/api/config/llm/test", methods=["POST"])
@pin_required
@csrf_required
def api_llm_test():
    """Ping the configured (or override) OpenAI-compatible endpoint.

    Empty override fields fall back to the saved config, so the settings panel
    can test with unsaved form values (masked secrets are sent as empty).
    """
    data = request.get_json(force=True, silent=True) or {}
    base_url = str(data.get("base_url", "") or "").strip()
    api_key = str(data.get("api_key", "") or "").strip()
    model = str(data.get("model", "") or "").strip()
    try:
        return jsonify(core.llm_test_connection(base_url, api_key, model))
    except Exception as exc:
        logger.exception("llm test failed")
        return jsonify({"ok": False, "error": str(exc)}), 500


@app.route("/api/status")
@pin_required
def api_status():
    return jsonify(core.status())


@app.route("/api/auth/send_code", methods=["POST"])
@csrf_required
@rate_limit(max_requests=3, window=300)
def api_send_code():
    """Legacy single-account endpoint — superseded by /api/accounts/<id>/send_code."""
    return jsonify({"error": "Use the account-scoped endpoint"}), 410


@app.route("/api/auth/sign_in", methods=["POST"])
@csrf_required
def api_sign_in():
    """Legacy single-account endpoint — superseded by /api/accounts/<id>/sign_in."""
    return jsonify({"error": "Use the account-scoped endpoint"}), 410


@app.route("/api/dialogs")
@pin_required
def api_dialogs():
    """Legacy endpoint — superseded by /api/accounts/<id>/dialogs."""
    return jsonify({"error": "Use the account-scoped endpoint"}), 410


@app.route("/api/dialogs/refresh", methods=["POST"])
@pin_required
@csrf_required
def api_refresh_dialogs():
    """Legacy endpoint — superseded by /api/accounts/<id>/dialogs/refresh."""
    return jsonify({"error": "Use the account-scoped endpoint"}), 410


# ── Pair management ─────────────────────────────────────────────────────────
@app.route("/api/pairs", methods=["GET"])
@pin_required
def api_get_pairs():
    return jsonify(core.get_pairs())


@app.route("/api/pairs", methods=["POST"])
@pin_required
@csrf_required
def api_create_pair():
    data = request.get_json(force=True, silent=True) or {}
    source = data.get("source_id")
    dest = data.get("dest_id")
    if source is None or dest is None:
        return jsonify({"error": "source_id and dest_id required"}), 400
    try:
        source_id = int(source)
        dest_id = int(dest)
    except (TypeError, ValueError):
        return jsonify({"error": "source_id and dest_id must be integers"}), 400
    if source_id == dest_id:
        return jsonify({"error": "Source and destination must differ"}), 400
    fields = {
        "source_title": data.get("source_title", ""),
        "dest_title": data.get("dest_title", ""),
        "enabled": bool(data.get("enabled", True)),
        "include_text": bool(data.get("include_text", True)),
        "include_media": bool(data.get("include_media", True)),
        "skip_standalone_media": bool(data.get("skip_standalone_media", False)),
        "forward_as_link": bool(data.get("forward_as_link", False)),
        "price_augment": bool(data.get("price_augment", False)),
        "forward_via_bot": bool(data.get("forward_via_bot", False)),
        "augment_symbol": data.get("augment_symbol") or None,
        "filter_type": data.get("filter_type", "none"),
        "llm_enabled": bool(data.get("llm_enabled", False)),
        "account_id": data.get("account_id") or None,
        "bot_token_id": data.get("bot_token_id") or None,
    }
    try:
        return jsonify(core.create_pair(source_id, dest_id, **fields))
    except Exception as e:
        logger.exception("create_pair failed")
        return jsonify({"error": str(e)}), 500


@app.route("/api/pairs/<int:pair_id>", methods=["PATCH"])
@pin_required
@csrf_required
def api_update_pair(pair_id):
    data = request.get_json(force=True, silent=True) or {}
    allowed = {
        "source_title",
        "dest_title",
        "enabled",
        "include_text",
        "include_media",
        "skip_standalone_media",
        "forward_as_link",
        "price_augment",
        "forward_via_bot",
        "augment_symbol",
        "filter_type",
        "llm_enabled",
        "account_id",
        "bot_token_id",
    }
    fields = {k: v for k, v in data.items() if k in allowed}
    try:
        return jsonify(core.update_pair(pair_id, **fields))
    except Exception as e:
        logger.exception("update_pair failed")
        return jsonify({"error": str(e)}), 500


@app.route("/api/pairs/<int:pair_id>/toggle", methods=["POST"])
@pin_required
@csrf_required
def api_toggle_pair(pair_id):
    try:
        return jsonify(core.toggle_pair(pair_id))
    except Exception as e:
        logger.exception("toggle_pair failed")
        return jsonify({"error": str(e)}), 500


@app.route("/api/pairs/<int:pair_id>", methods=["DELETE"])
@pin_required
@csrf_required
def api_delete_pair(pair_id):
    try:
        return jsonify(core.delete_pair(pair_id))
    except Exception as e:
        logger.exception("delete_pair failed")
        return jsonify({"error": str(e)}), 500


@app.route("/api/pairs/<int:pair_id>/test_emit", methods=["POST"])
@pin_required
@csrf_required
def api_test_emit(pair_id):
    data = request.get_json(force=True, silent=True) or {}
    text = str(data.get("text", "")).strip()
    if not text:
        return jsonify({"error": "text required"}), 400
    try:
        return jsonify(core.emit_test_signal(pair_id, text))
    except Exception as e:
        logger.exception("test_emit failed")
        return jsonify({"error": str(e)}), 500


# ── Bot webhook (optional inbound endpoint) ─────────────────────────────────
@app.route("/api/bot/webhook/<token>", methods=["POST"])
def api_bot_webhook(token):
    """Receive Telegram Bot API updates when registered as a webhook."""
    if token != config.SIGNAL_BOT_TOKEN:
        return jsonify({"error": "Invalid token"}), 403
    data = request.get_json(force=True, silent=True) or {}
    logger.info(
        "Bot webhook update: message_id=%s",
        data.get("message", {}).get("message_id"),
    )
    return jsonify({"status": "ok"})


# ── Forwarder control ───────────────────────────────────────────────────────
@app.route("/api/forwarder/start", methods=["POST"])
@pin_required
@csrf_required
def api_start():
    return jsonify(core.start_forwarding())


@app.route("/api/forwarder/stop", methods=["POST"])
@pin_required
@csrf_required
def api_stop():
    return jsonify(core.stop_forwarding())


# ── Account management ──────────────────────────────────────────────────────
@app.route("/api/accounts", methods=["GET"])
@pin_required
def api_accounts():
    return jsonify(core._run_db_sync(db.get_accounts))


@app.route("/api/accounts", methods=["POST"])
@pin_required
@csrf_required
def api_create_account():
    data = request.get_json(force=True, silent=True) or {}
    name = str(data.get("name", "")).strip()
    phone = str(data.get("phone", "")).strip()
    if not name:
        return jsonify({"error": "Account name required"}), 400
    accounts = core._run_db_sync(db.get_accounts)
    if len(accounts) >= config.MAX_TG_ACCOUNTS:
        return jsonify(
            {"error": f"Maximum {config.MAX_TG_ACCOUNTS} accounts allowed"}
        ), 400
    try:
        is_primary = len(accounts) == 0
        acct_id = core._run_db_sync(db.save_account, name, phone, is_primary)
        return jsonify({"status": "ok", "id": acct_id})
    except Exception as e:
        logger.exception("create_account failed")
        return jsonify({"error": str(e)}), 500


@app.route("/api/accounts/<int:account_id>", methods=["PATCH"])
@pin_required
@csrf_required
def api_update_account(account_id):
    data = request.get_json(force=True, silent=True) or {}
    allowed = {"name", "phone", "is_primary"}
    fields = {k: v for k, v in data.items() if k in allowed}
    try:
        ok = core._run_db_sync(db.update_account, account_id, **fields)
        return jsonify({"status": "ok" if ok else "not_found"})
    except Exception as e:
        logger.exception("update_account failed")
        return jsonify({"error": str(e)}), 500


@app.route("/api/accounts/<int:account_id>", methods=["DELETE"])
@pin_required
@csrf_required
def api_delete_account(account_id):
    try:
        # Disconnect client first if connected
        core._run(core._stop_account_client(account_id), timeout=10)
        ok = core._run_db_sync(db.delete_account, account_id)
        return jsonify({"status": "ok" if ok else "not_found"})
    except Exception as e:
        logger.exception("delete_account failed")
        return jsonify({"error": str(e)}), 500


@app.route("/api/accounts/<int:account_id>/reset", methods=["POST"])
@pin_required
@csrf_required
def api_reset_account(account_id):
    """Clear all stale auth state for an account so it can start fresh."""
    try:
        # Disconnect any in-memory client for this account
        core._run(core._stop_account_client(account_id), timeout=10)
        # Clear in-memory auth state for this account only
        core._auth_state.pop(account_id, None)
        core._qr_state.pop(account_id, None)
        ok = core._run_db_sync(db.reset_account, account_id)
        if ok:
            return jsonify({"status": "ok", "message": "Account reset — request a new code."})
        return jsonify({"error": "Account not found"}), 404
    except Exception as e:
        logger.exception("reset_account failed")
        return jsonify({"error": str(e)}), 500


@app.route("/api/accounts/reset-stale", methods=["POST"])
@pin_required
@csrf_required
def api_reset_all_stale():
    """Clear stale auth state for all pending (unauthorized) accounts."""
    try:
        core._auth_state.clear()
        core._qr_state.clear()
        count = core._run_db_sync(db.reset_all_stale)
        return jsonify({"status": "ok", "reset": count})
    except Exception as e:
        logger.exception("reset_all_stale failed")
        return jsonify({"error": str(e)}), 500


# ── Account auth (phone code) ───────────────────────────────────────────────
@app.route("/api/accounts/<int:account_id>/send_code", methods=["POST"])
@pin_required
@csrf_required
@rate_limit(max_requests=3, window=300)
def api_account_send_code(account_id):
    data = request.get_json(force=True, silent=True) or {}
    acct = core._run_db_sync(db.get_account, account_id)
    if not acct:
        return jsonify({"error": "Account not found"}), 404
    phone = str(data.get("phone", acct.get("phone", "") or config.TG_PHONE or "")).strip()
    if not phone:
        return jsonify({"error": "Phone number required"}), 400
    try:
        return jsonify(core.send_code(phone, account_id))
    except RuntimeError as e:
        return jsonify({"error": str(e)}), 400
    except Exception as e:
        logger.exception("account_send_code failed")
        return jsonify({"error": str(e)}), 500


@app.route("/api/accounts/<int:account_id>/sign_in", methods=["POST"])
@pin_required
@csrf_required
def api_account_sign_in(account_id):
    data = request.get_json(force=True, silent=True) or {}
    code = str(data.get("code", "")).strip()
    password = data.get("password")
    password = str(password).strip() if password else None
    if not code:
        return jsonify({"error": "Code required"}), 400
    try:
        return jsonify(core.sign_in(code, password, account_id))
    except RuntimeError as e:
        msg = str(e)
        if msg == "2FA_PASSWORD_REQUIRED":
            return jsonify({"status": "2fa_required", "error": "2FA password required"}), 403
        if msg == "INVALID_CODE":
            return jsonify({"error": "Invalid code"}), 400
        return jsonify({"error": msg}), 400
    except Exception as e:
        logger.exception("account_sign_in failed")
        return jsonify({"error": str(e)}), 500


# ── Account auth (QR code) ─────────────────────────────────────────────────
def _qr_data_uri(url: str) -> str:
    """Render a login URL as a base64 PNG data URI for an <img> tag."""
    import base64
    import io
    import qrcode

    img = qrcode.make(url)
    buf = io.BytesIO()
    img.save(buf, format="PNG")
    buf.seek(0)
    return "data:image/png;base64," + base64.b64encode(buf.read()).decode()


@app.route("/api/accounts/<int:account_id>/qr", methods=["POST"])
@pin_required
@csrf_required
def api_start_qr(account_id):
    try:
        result = core.start_qr_login(account_id)
        if result.get("qr_url"):
            try:
                result["qr_data_uri"] = _qr_data_uri(result["qr_url"])
            except Exception as e:
                logger.warning("QR image render failed: %s", e)
        return jsonify(result)
    except RuntimeError as e:
        return jsonify({"error": str(e)}), 400
    except Exception as e:
        logger.exception("start_qr failed")
        return jsonify({"error": str(e)}), 500


@app.route("/api/accounts/<int:account_id>/qr/check", methods=["POST"])
@pin_required
@csrf_required
def api_check_qr(account_id):
    try:
        result = core.check_qr_login(account_id)
        # When the token was refreshed, hand back a fresh image to render.
        if result.get("status") == "refresh" and result.get("qr_url"):
            try:
                result["qr_data_uri"] = _qr_data_uri(result["qr_url"])
            except Exception as e:
                logger.warning("QR refresh image render failed: %s", e)
        return jsonify(result)
    except RuntimeError as e:
        return jsonify({"error": str(e)}), 400
    except Exception as e:
        logger.exception("check_qr failed")
        return jsonify({"error": str(e)}), 500


@app.route("/api/accounts/<int:account_id>/qr/cancel", methods=["POST"])
@pin_required
@csrf_required
def api_cancel_qr(account_id):
    try:
        return jsonify(core.cancel_qr_login(account_id))
    except Exception as e:
        logger.exception("cancel_qr failed")
        return jsonify({"error": str(e)}), 500


@app.route("/api/accounts/<int:account_id>/qr/2fa", methods=["POST"])
@pin_required
@csrf_required
def api_complete_qr_2fa(account_id):
    data = request.get_json(force=True, silent=True) or {}
    password = str(data.get("password", "")).strip()
    if not password:
        return jsonify({"error": "Password required"}), 400
    try:
        return jsonify(core.complete_qr_2fa(account_id, password))
    except RuntimeError as e:
        return jsonify({"error": str(e)}), 400
    except Exception as e:
        logger.exception("complete_qr_2fa failed")
        return jsonify({"error": str(e)}), 500


# ── Account-specific dialogs ────────────────────────────────────────────────
@app.route("/api/accounts/<int:account_id>/dialogs")
@pin_required
def api_account_dialogs(account_id):
    try:
        return jsonify(core.get_dialogs(account_id=account_id))
    except RuntimeError as e:
        return jsonify({"error": str(e)}), 400
    except Exception as e:
        logger.exception("account_dialogs failed")
        return jsonify({"error": str(e)}), 500


@app.route("/api/accounts/<int:account_id>/dialogs/refresh", methods=["POST"])
@pin_required
@csrf_required
def api_account_refresh_dialogs(account_id):
    try:
        return jsonify(core.refresh_dialogs(account_id=account_id))
    except RuntimeError as e:
        return jsonify({"error": str(e)}), 400
    except Exception as e:
        logger.exception("account_refresh_dialogs failed")
        return jsonify({"error": str(e)}), 500


# ── Bot token management ────────────────────────────────────────────────────
@app.route("/api/bots", methods=["GET"])
@pin_required
def api_bot_tokens():
    return jsonify(core._run_db_sync(db.get_bot_tokens))


@app.route("/api/bots", methods=["POST"])
@pin_required
@csrf_required
def api_add_bot():
    data = request.get_json(force=True, silent=True) or {}
    token = str(data.get("token", "")).strip()
    if not token:
        return jsonify({"error": "Bot token required"}), 400
    if not token.count(":") == 1 or not token.split(":")[0].isdigit():
        return jsonify({"error": "Invalid bot token format"}), 400
    try:
        bot_id = core._run_db_sync(
            db.save_bot_token,
            token,
            bot_name=data.get("bot_name", ""),
            bot_username=data.get("bot_username", ""),
        )
        return jsonify({"status": "ok", "id": bot_id})
    except Exception as e:
        logger.exception("add_bot failed")
        return jsonify({"error": str(e)}), 500


@app.route("/api/bots/<int:bot_id>", methods=["DELETE"])
@pin_required
@csrf_required
def api_delete_bot(bot_id):
    try:
        ok = core._run_db_sync(db.delete_bot_token, bot_id)
        return jsonify({"status": "ok" if ok else "not_found"})
    except Exception as e:
        logger.exception("delete_bot failed")
        return jsonify({"error": str(e)}), 500


@app.route("/api/bots/<int:bot_id>/verify", methods=["POST"])
@pin_required
@csrf_required
def api_verify_bot(bot_id):
    try:
        token = core._run_db_sync(db.get_decrypted_bot_token, bot_id)
        if not token:
            return jsonify({"error": "Bot token not found"}), 404
        # Verify with Telegram
        import urllib.request

        url = f"https://api.telegram.org/bot{token}/getMe"
        req = urllib.request.Request(url)
        with urllib.request.urlopen(req, timeout=10) as resp:
            import json as _json

            data = _json.loads(resp.read().decode())
        if data.get("ok"):
            bot_info = data["result"]
            core._run_db_sync(
                db.update_bot_token,
                bot_id,
                bot_name=bot_info.get("first_name", ""),
                bot_username=bot_info.get("username", ""),
                verified=True,
            )
            return jsonify(
                {
                    "status": "ok",
                    "bot_name": bot_info.get("first_name", ""),
                    "bot_username": bot_info.get("username", ""),
                }
            )
        return jsonify({"error": "Invalid bot token"}), 400
    except Exception as e:
        logger.exception("verify_bot failed")
        return jsonify({"error": str(e)}), 500


# ── Backup ──────────────────────────────────────────────────────────────────
@app.route("/api/backup", methods=["GET"])
@pin_required
def api_backup():
    try:
        import json as _json
        from datetime import datetime, timezone

        backup = {
            "version": 2,
            "created_at": datetime.now(timezone.utc).isoformat(),
            "tables": {},
        }

        # Dump all relevant tables
        tables = [
            "channel_pairs",
            "tg_accounts",
            "bot_tokens",
            "config_vars",
            "signal_log",
            "auth_state",
            "message_map",
        ]
        for table in tables:
            rows = core._run_db_sync(db._dump_table, table)
            if rows is not None:
                backup["tables"][table] = rows

        # Session strings are Telegram auth keys — never include them in a backup.
        acct_rows = backup.get("tables", {}).get("tg_accounts")
        if acct_rows:
            for row in acct_rows:
                row.pop("session_string", None)
                row.pop("pending_session", None)
                row.pop("phone_code_hash", None)

        # Include .env content (with sensitive fields masked)
        env_path = os.path.join(os.path.dirname(os.path.abspath(__file__)), ".env")
        if os.path.exists(env_path):
            with open(env_path, "r") as f:
                env_content = f.read()
            # Mask sensitive values
            sensitive_keys = [
                "API_KEY", "PASSWORD", "SECRET", "TOKEN", "PIN",
                "MYSQL_PASSWORD", "TG_API_HASH", "WEB_PIN",
            ]
            masked_lines = []
            for line in env_content.split("\n"):
                line = line.strip()
                if not line or line.startswith("#"):
                    masked_lines.append(line)
                else:
                    key = line.split("=")[0].strip() if "=" in line else ""
                    if any(s in key.upper() for s in sensitive_keys):
                        masked_lines.append(f"{key}=***MASKED***")
                    else:
                        masked_lines.append(line)
            backup["env"] = "\n".join(masked_lines)

        response = app.response_class(
            response=_json.dumps(backup, indent=2, default=str),
            status=200,
            mimetype="application/json",
            headers={
                "Content-Disposition": "attachment; filename=tg-forwarder-backup.json"
            },
        )
        return response
    except Exception as e:
        logger.exception("backup failed")
        return jsonify({"error": str(e)}), 500


# ── cTrader OAuth Callback + Token Vault ────────────────────────────────────

# In-memory OAuth state storage with TTL (5 minutes)
# Maps state_token -> created_at_timestamp
_oauth_states: dict[str, float] = {}


def _clean_oauth_states():
    """Prune expired OAuth states (older than 5 min)."""
    now = time.time()
    expired = [s for s, t in _oauth_states.items() if now - t > 300]
    for s in expired:
        _oauth_states.pop(s, None)


def _build_consent_url() -> str:
    """Build the cTrader OAuth consent URL with a fresh CSRF state."""
    if not config.CTRADER_CLIENT_ID:
        raise RuntimeError("CTRADER_CLIENT_ID not configured")

    state = secrets.token_urlsafe(32)
    _clean_oauth_states()
    _oauth_states[state] = time.time()

    params = {
        "client_id": config.CTRADER_CLIENT_ID,
        "redirect_uri": config.CTRADER_REDIRECT_URI,
        "scope": "trading",
        "product": "web",
        "state": state,
    }
    return (
        "https://id.ctrader.com/my/settings/openapi/grantingaccess/?"
        + urllib.parse.urlencode(params)
    )


def _exchange_code(code: str) -> dict:
    """Exchange an authorization code for access/refresh tokens.

    Matches the v0 implementation: GET with query parameters to
    https://openapi.ctrader.com/apps/token
    """
    import json as _json
    import urllib.parse
    import urllib.request

    params = urllib.parse.urlencode({
        "grant_type": "authorization_code",
        "code": code,
        "redirect_uri": config.CTRADER_REDIRECT_URI,
        "client_id": config.CTRADER_CLIENT_ID,
        "client_secret": config.CTRADER_CLIENT_SECRET,
    })
    url = f"https://openapi.ctrader.com/apps/token?{params}"

    req = urllib.request.Request(
        url,
        headers={"Accept": "application/json"},
        method="GET",
    )

    with urllib.request.urlopen(req, timeout=30) as resp:
        body = _json.loads(resp.read().decode("utf-8"))

    access_token = body.get("accessToken") or body.get("access_token")
    if not access_token:
        raise RuntimeError(f"Token exchange returned no access_token: {body}")

    return {
        "access_token": access_token,
        "refresh_token": body.get("refreshToken") or body.get("refresh_token", ""),
        "expires_at": time.time() + body.get("expiresIn", body.get("expires_in", 2_628_000)),
    }


def _fetch_ctrader_profile(access_token: str) -> dict | None:
    """Fetch cTID profile via Spotware REST (only works during OAuth window)."""
    import json as _json
    import urllib.request

    try:
        url = f"https://api.spotware.com/connect/profile?access_token={access_token}"
        req = urllib.request.Request(url, headers={"Accept": "application/json"})
        with urllib.request.urlopen(req, timeout=15) as resp:
            data = _json.loads(resp.read().decode("utf-8"))
        return data.get("data") if isinstance(data, dict) else None
    except Exception as exc:
        logger.warning("Profile fetch failed: %s", exc)
        return None


def _fetch_ctrader_accounts(access_token: str) -> list[dict]:
    """Fetch trading accounts via Spotware REST.

    Spotware wraps the account list in a ``{"data": [...]}`` envelope.
    """
    import json as _json
    import urllib.request

    try:
        url = f"https://api.spotware.com/connect/tradingaccounts?access_token={access_token}"
        req = urllib.request.Request(url, headers={"Accept": "application/json"})
        with urllib.request.urlopen(req, timeout=15) as resp:
            body = _json.loads(resp.read().decode("utf-8"))
        # Spotware returns {"data": [ {...}, ... ]}
        data = body.get("data") if isinstance(body, dict) else body
        return data if isinstance(data, list) else []
    except Exception as exc:
        logger.warning("Account discovery failed: %s", exc)
        return []


@app.route("/callback")
def ctrader_callback():
    """OAuth callback handler — public, protected by state parameter."""
    code = request.args.get("code")
    state = request.args.get("state")
    error = request.args.get("error")

    if error:
        logger.error("OAuth error from cTrader: %s", error)
        return render_template(
            "callback.html", ok=False,
            title="cTrader refused the connection",
            message=f"cTrader returned an error: <code>{html.escape(error)}</code>.",
            hint="Start again from the dashboard → cTrader tab.",
        ), 400

    if not code:
        return render_template(
            "callback.html", ok=False,
            title="No authorization code",
            message="The redirect from cTrader did not include an authorization code, so nothing was linked.",
            hint="Start again from the dashboard → cTrader tab.",
        ), 400

    _clean_oauth_states()
    if not state or state not in _oauth_states:
        logger.warning("Invalid or expired OAuth state: %s", state)
        return render_template(
            "callback.html", ok=False,
            title="This sign-in link has expired",
            message="The authentication session is no longer valid — sign-in links are single-use and short-lived.",
            hint="Start again from the dashboard → cTrader tab.",
        ), 403

    # State is valid — consume it
    _oauth_states.pop(state, None)

    try:
        tokens = _exchange_code(code)
    except Exception as exc:
        logger.exception("Token exchange failed")
        return render_template(
            "callback.html", ok=False,
            title="Token exchange failed",
            message=f"The authorization code could not be exchanged for tokens: <code>{html.escape(str(exc))}</code>.",
            hint="Start again from the dashboard → cTrader tab.",
        ), 400

    # Discover accounts
    profile = _fetch_ctrader_profile(tokens["access_token"])
    accounts = _fetch_ctrader_accounts(tokens["access_token"])
    user_id = profile.get("userId") if profile else None

    if not config.CTRADER_TOKEN_SECRET:
        logger.error("CTRADER_TOKEN_SECRET not configured — cannot encrypt tokens")
        return render_template(
            "callback.html", ok=False,
            title="Server is missing its token secret",
            message="CTRADER_TOKEN_SECRET is not configured, so the tokens could not be stored.",
            hint="Add CTRADER_TOKEN_SECRET to the server configuration, then reconnect.",
        ), 500

    from crypto_vault import encrypt_payload

    for acc in accounts:
        account_id = acc.get("accountId")
        if not account_id:
            continue

        payload = {
            "access_token": tokens["access_token"],
            "refresh_token": tokens["refresh_token"],
            "expires_at": tokens["expires_at"],
            "user_id": user_id,
        }

        try:
            encrypted = encrypt_payload(config.CTRADER_TOKEN_SECRET, payload)
        except Exception as exc:
            logger.error("Encryption failed for account %s: %s", account_id, exc)
            continue

        try:
            core._run_db_sync(
                db.save_ctrader_account,
                account_id=account_id,
                encrypted_payload=encrypted,
                user_id=user_id,
                broker_name=acc.get("brokerName", ""),
                broker_title=acc.get("brokerTitle", ""),
                account_number=acc.get("accountNumber"),
                is_live=bool(acc.get("live", False)),
                deposit_currency=acc.get("depositCurrency", ""),
                balance_cents=acc.get("balance"),
                money_digits=acc.get("moneyDigits", 2),
                leverage=acc.get("leverage"),
                account_type=acc.get("traderAccountType", ""),
                account_status=acc.get("accountStatus", ""),
            )
            logger.info("Stored encrypted token for account %s", account_id)
        except Exception as exc:
            logger.exception("DB save failed for account %s", account_id)

    # Render the success page
    view_accounts = [
        {
            "account_id": a.get("accountId"),
            "broker_title": a.get("brokerTitle") or a.get("brokerName"),
            "is_live": bool(a.get("live", False)),
            "deposit_currency": a.get("depositCurrency", ""),
        }
        for a in accounts
    ]
    if accounts:
        return render_template(
            "callback.html", ok=True, accounts=view_accounts,
            title="cTrader is now linked",
            message="Tokens are encrypted and stored. The forwarder can read live prices for these accounts.",
        ), 200
    return render_template(
        "callback.html", ok=True, accounts=[],
        title="Connected, but no accounts found",
        message="The OAuth flow completed, but this cTrader ID has no trading accounts attached to it.",
        hint="Create a demo or live account with your broker, then reconnect.",
    ), 200


@app.route("/api/ctrader/connect", methods=["POST"])
@pin_required
@csrf_required
def api_ctrader_connect():
    """Return the cTrader OAuth consent URL for the user to visit."""
    try:
        url = _build_consent_url()
        return jsonify({"status": "ok", "url": url})
    except RuntimeError as exc:
        return jsonify({"error": str(exc)}), 500


@app.route("/api/ctrader/accounts", methods=["GET"])
@pin_required
def api_ctrader_accounts():
    """List stored cTrader accounts (metadata only, unless ?encrypted=true)."""
    include_encrypted = request.args.get("encrypted", "").lower() in ("1", "true", "yes")
    try:
        rows = core._run_db_sync(db.list_ctrader_accounts)
        if include_encrypted:
            # Include the encrypted payload for consumers
            full_rows = core._run_db_sync(db.list_ctrader_accounts)
            result = []
            for row in full_rows:
                acct = core._run_db_sync(db.get_ctrader_account, row["account_id"])
                if acct:
                    row["encrypted_payload"] = acct.get("encrypted_payload")
                result.append(row)
            return jsonify({"status": "ok", "accounts": result})
        return jsonify({"status": "ok", "accounts": rows})
    except Exception as exc:
        logger.exception("ctrader_accounts failed")
        return jsonify({"error": str(exc)}), 500


@app.route("/api/ctrader/accounts/<int:account_id>", methods=["DELETE"])
@pin_required
@csrf_required
def api_ctrader_account_delete(account_id):
    """Remove a stored cTrader account from the vault."""
    try:
        ok = core._run_db_sync(db.delete_ctrader_account, account_id)
        return jsonify({"status": "ok" if ok else "not_found"})
    except Exception as exc:
        logger.exception("ctrader_account_delete failed")
        return jsonify({"error": str(exc)}), 500


@app.route("/api/ctrader/accounts/<int:account_id>/refresh", methods=["POST"])
@pin_required
@csrf_required
def api_ctrader_account_refresh(account_id):
    """Push a refreshed encrypted payload back into the vault.

    Expects JSON: {"encrypted_payload": "base64..."}
    """
    data = request.get_json(force=True, silent=True) or {}
    encrypted_payload = data.get("encrypted_payload")
    if not encrypted_payload:
        return jsonify({"error": "encrypted_payload required"}), 400
    try:
        ok = core._run_db_sync(
            db.update_ctrader_account_payload, account_id, encrypted_payload
        )
        if ok:
            logger.info("Refreshed payload saved for account %d", account_id)
        return jsonify({"status": "ok" if ok else "not_found"})
    except Exception as exc:
        logger.exception("ctrader_account_refresh failed")
        return jsonify({"error": str(exc)}), 500


# ── HTMX Component Routes ───────────────────────────────────────────────────
# These return HTML partials for HTMX-driven UI.

@app.route("/_components/status")
@pin_required
def htmx_status():
    status = core.status()
    # Build health metrics similar to /api/health
    try:
        rss = psutil.Process().memory_info().rss
        pairs_count = len(status.get("pairs", []))
    except Exception:
        rss = 0
        pairs_count = 0
    health = {"pairs_count": pairs_count, "rss_mb": round(rss / 1024 / 1024, 2)}
    return render_template(
        "components/_status_panel.html",
        status=status,
        health=health,
        csrf=generate_csrf(),
        auto_poll=True,
        core_started=bool(_core_started),
        auto_start_core=bool(config.AUTO_START_CORE),
    )


@app.route("/_components/accounts")
@pin_required
def htmx_accounts():
    accounts = core._run_db_sync(db.get_accounts)
    return render_template(
        "components/_accounts_panel.html",
        accounts=accounts,
        csrf=generate_csrf(),
        api_id=config.TG_API_ID,
        config_phone=config.TG_PHONE,
    )


@app.route("/_components/pairs")
@pin_required
def htmx_pairs():
    q, status_filter = _pairs_filter()
    pairs = core.get_pairs()
    total = len(pairs)
    if q:
        pairs = [p for p in pairs if _pair_searchable(p, q)]
    if status_filter == "enabled":
        pairs = [p for p in pairs if p.get("enabled")]
    elif status_filter == "disabled":
        pairs = [p for p in pairs if not p.get("enabled")]
    accounts = core._run_db_sync(db.get_accounts)
    bots = core._run_db_sync(db.get_bot_tokens)
    # Default the channel picker to a connected (prefer primary) authorized
    # account so the initial dropdown never mixes channels from arbitrary
    # Telegram accounts and keeps working when the primary is offline.
    default_account_id = _default_account_id(core, accounts)
    try:
        dialogs = core.get_dialogs(account_id=default_account_id)
    except Exception:
        dialogs = []
    return render_template(
        "components/_pairs_panel.html",
        pairs=pairs,
        accounts=accounts or [],
        bots=bots or [],
        dialogs=dialogs or [],
        default_account_id=default_account_id,
        csrf=generate_csrf(),
        q=q,
        status=status_filter,
        total_count=total,
        filtered_count=len(pairs),
    )


def _default_account_id(core, accounts):
    """Pick the account the channel picker should default to.

    Prefer a *connected* authorized account (the primary among the connected
    set, else the first connected), falling back to authorized-but-disconnected
    accounts. This keeps the New Rule / channel picker working even when the DB
    primary account's Telegram client is momentarily offline, which is the most
    common cause of "no channels load" in a multi-account setup.
    """
    authorized_accounts = [a for a in (accounts or []) if a.get("is_authorized")]
    if not authorized_accounts:
        return None
    connected_ids = set(getattr(core, "clients", {}).keys())
    pool = (
        [a for a in authorized_accounts if a["id"] in connected_ids]
        or authorized_accounts
    )
    if len(pool) == 1:
        return pool[0]["id"]
    primary = next((a for a in pool if a.get("is_primary")), None)
    return primary["id"] if primary else pool[0]["id"]


def _pairs_filter():
    """Resolve the active search/filter for a full-panel reload.

    Reads explicit ?q=/&status= args first, then falls back to the
    HX-Current-URL (HTMX sends it on every request), so toggling/deleting/
    creating a rule keeps the current search + On/Off filter instead of
    dumping the user back to the unfiltered list.
    """
    q = request.args.get("q")
    status_filter = request.args.get("status")
    if q is None or status_filter is None:
        # POST bodies (create form) carry hidden q/status fields.
        q = q if q is not None else request.form.get("q")
        status_filter = (
            status_filter if status_filter is not None else request.form.get("status")
        )
    if q is None or status_filter is None:
        cur = request.headers.get("HX-Current-URL", "")
        if cur.startswith("/") or "://" in cur:
            from urllib.parse import parse_qs, urlsplit
            qs = parse_qs(urlsplit(cur).query)
            if q is None:
                q = (qs.get("q") or [""])[0]
            if status_filter is None:
                status_filter = (qs.get("status") or ["all"])[0]
    return (q or "").strip().lower(), (status_filter or "all").strip().lower()


def _pair_searchable(p: dict, q: str) -> bool:
    """Token-aware match: short queries match whole tokens, longer match substrings."""
    import re as _re
    parts = [
        str(p.get("id", "")),
        p.get("source_title", ""),
        p.get("dest_title", ""),
        p.get("account_name", ""),
        p.get("bot_username", ""),
        p.get("augment_symbol", ""),
        p.get("filter_type", ""),
        p.get("category", ""),
        "on" if p.get("enabled") else "off",
    ]
    if p.get("price_augment"):
        parts.append("price+")
    if p.get("forward_via_bot"):
        parts.append("bot")
    if p.get("forward_as_link"):
        parts.append("link")
    if p.get("skip_standalone_media"):
        parts.append("nomedia")
    if p.get("llm_enabled"):
        parts.append("llm")
    hay = " ".join(part for part in parts if part).lower()
    if len(q) <= 2:
        return bool(_re.search(r"\b" + _re.escape(q), hay))
    return q in hay


@app.route("/_components/pairs/search")
@pin_required
def htmx_pairs_search():
    """Return filtered pairs table.  Supports ?q= (full-text) and ?status= (enabled|disabled|all)."""
    q, status_filter = _pairs_filter()
    pairs = core.get_pairs()
    total = len(pairs)
    if q:
        pairs = [p for p in pairs if _pair_searchable(p, q)]
    if status_filter == "enabled":
        pairs = [p for p in pairs if p.get("enabled")]
    elif status_filter == "disabled":
        pairs = [p for p in pairs if not p.get("enabled")]
    filtered = len(pairs)
    return render_template(
        "components/_pairs_table.html",
        pairs=pairs,
        csrf=generate_csrf(),
        q=q,
        status=status_filter,
        total_count=total,
        filtered_count=filtered,
    )


@app.route("/_components/pairs/channel_options")
@pin_required
def htmx_channel_options():
    q = request.args.get("q", "")
    account_id = request.args.get("account_id", type=int) or None
    if account_id is None:
        # "Auto" mode still needs a channel list.  Prefer a connected account
        # (primary among connected) so the picker is stable and never mixes
        # channels from arbitrary Telegram accounts.
        accounts = core._run_db_sync(db.get_accounts)
        account_id = _default_account_id(core, accounts)
    try:
        dialogs = core.get_dialogs(account_id=account_id)
    except Exception:
        dialogs = []
    if q:
        ql = q.lower().strip()
        dialogs = [
            d
            for d in (dialogs or [])
            if ql
            in " ".join(
                str(d.get(field) or "")
                for field in ("title", "type", "category", "username", "id")
            ).lower()
        ]
    # Stable sort so the same q always renders in the same order
    dialogs = sorted(
        (dialogs or []), key=lambda d: d.get("title", "").lower()
    )
    opts = render_template(
        "components/_channel_options.html",
        dialogs=dialogs or [],
    )
    # Destination search: update only the destination list (an OOB outerHTML
    # swap would replace the select element and lose the current selection).
    if request.headers.get("HX-Target") == "dest-options":
        return opts
    # Account change: refresh BOTH dropdowns via OOB (different account =
    # different channel list).
    if request.args.get("reload"):
        oob = (
            '<select id="dest-options" hx-swap-oob="true" name="dest_ids" multiple '
            'size="4" style="min-height:90px;" required '
            'x-ref="dest" @change="checkConflict()">{}</select>'
        ).format(opts)
        return opts + oob
    # Plain source search: only the source list changes.
    return opts


@app.route("/_components/bots")
@pin_required
def htmx_bots():
    bots = core._run_db_sync(db.get_bot_tokens)
    url = core._run_db_sync(db.get_config, "SIGNAL_WEBHOOK_URL") or config.SIGNAL_WEBHOOK_URL
    chat_id = core._run_db_sync(db.get_config, "SIGNAL_ADMIN_CHAT_ID") or config.SIGNAL_ADMIN_CHAT_ID
    return render_template(
        "components/_bots_panel.html",
        bots=bots or [],
        sig_webhook_url=url,
        sig_admin_chat=chat_id,
        csrf=generate_csrf(),
    )


@app.route("/_components/ctrader")
@pin_required
def htmx_ctrader():
    rows = core._run_db_sync(db.list_ctrader_accounts) or []
    return render_template(
        "components/_ctrader_panel.html",
        accounts=rows,
        csrf=generate_csrf(),
    )


@app.route("/_components/signals")
@pin_required
def htmx_signals():
    rows = core._run_db_sync(db.get_all_signal_snapshots, limit=50)
    return render_template(
        "components/_signals_panel.html",
        signals=rows or [],
    )


@app.route("/_components/llm")
@pin_required
def htmx_llm():
    st = core.llm_status()
    webhook_url = core._run_db_sync(db.get_config, "LLM_WEBHOOK_URL") or ""
    webhook_secret_set = bool(core._run_db_sync(db.get_secret_config, "LLM_WEBHOOK_SECRET"))
    pairs = core._run_db_sync(db.get_pairs) or []
    llm_pairs = sum(1 for p in pairs if p.get("llm_enabled") and p.get("enabled"))
    return render_template(
        "components/_llm_panel.html",
        llm={
            "enabled": bool(st.get("enabled")),
            "provider": st.get("provider", "openai-compatible"),
            "base_url": st.get("base_url", "https://api.cerebras.ai/v1"),
            "model": st.get("model", "llama3.1-8b"),
            "api_key_set": bool(st.get("api_key_set")),
            "client_available": bool(st.get("client_available", True)),
            "available": bool(st.get("available")),
            "min_confidence": st.get("min_confidence", 0.80),
            "timeout": st.get("timeout", 12.0),
            "max_tokens": st.get("max_tokens", 120),
            "webhook_url": webhook_url,
            "webhook_secret_set": webhook_secret_set,
            "pairs_count": llm_pairs,
        },
    )


@app.route("/_components/auth_modal/<int:account_id>")
@pin_required
def htmx_auth_modal(account_id):
    acct = core._run_db_sync(db.get_account, account_id)
    if not acct:
        return '<div class="overlay"><div class="overlay-box"><h2>Error</h2><p>Account not found</p><button onclick="closeModal()">Close</button></div></div>', 404
    return render_template(
        "components/_auth_modal.html",
        account_id=account_id,
        account_name=acct.get("name", ""),
        csrf=generate_csrf(),
        api_configured=bool(config.TG_API_ID and config.TG_API_HASH),
        config_phone=config.TG_PHONE,
    )


@app.route("/_components/pair_editor/<source_id>")
@pin_required
def htmx_pair_editor(source_id):
    try:
        source_id = int(source_id)
    except (ValueError, TypeError):
        return '<div class="overlay"><div class="overlay-box"><h2>Error</h2><p>Invalid source ID</p><button onclick="closeModal()">Close</button></div></div>', 404
    all_pairs = core.get_pairs()
    group = [p for p in all_pairs if p.get("source_chat_id") == source_id]
    if not group:
        # Try treating source_id as the pair id directly
        group = [p for p in all_pairs if p.get("id") == source_id]
    if not group:
        return '<div class="overlay"><div class="overlay-box"><h2>Error</h2><p>Rule not found</p><button onclick="closeModal()">Close</button></div></div>', 404

    first = group[0]
    # Scope the dialog list to the account the pair is bound to, so the
    # destination picker shows channels visible to the right account.
    acct_id = first.get("account_id") or None
    try:
        dialogs = core.get_dialogs(account_id=acct_id)
    except Exception:
        dialogs = []
    accounts = core._run_db_sync(db.get_accounts)
    bots = core._run_db_sync(db.get_bot_tokens)
    return render_template(
        "components/_pair_editor.html",
        source_id=source_id,
        source_title=first.get("source_title", first.get("source_chat_id", "Unknown")),
        pairs=group,
        filter_type=first.get("filter_type", "none"),
        augment_symbol=first.get("augment_symbol", ""),
        include_text=bool(first.get("include_text", True)),
        include_media=bool(first.get("include_media", True)),
        skip_standalone_media=bool(first.get("skip_standalone_media", False)),
        forward_as_link=bool(first.get("forward_as_link", False)),
        price_augment=bool(first.get("price_augment", False)),
        forward_via_bot=bool(first.get("forward_via_bot", False)),
        llm_enabled=bool(first.get("llm_enabled", False)),
        original_dest_ids=[p.get("dest_chat_id") for p in group],
        dialogs=dialogs or [],
        accounts=accounts or [],
        bots=bots or [],
        csrf=generate_csrf(),
    )


@app.route("/_components/pair_editor/channel_options/<source_id>")
@pin_required
def htmx_pair_editor_channel_options(source_id):
    """Return filtered channel options for the pair editor, excluding already-used dests."""
    try:
        source_id = int(source_id)
    except (ValueError, TypeError):
        return "", 404
    q = request.args.get("q", "")
    all_pairs = core.get_pairs()
    group = [p for p in all_pairs if p.get("source_chat_id") == source_id]
    used_ids = set(p.get("dest_chat_id") for p in group if p.get("dest_chat_id"))
    acct_id = (group[0].get("account_id") or None) if group else None
    try:
        dialogs = core.get_dialogs(account_id=acct_id)
    except Exception:
        dialogs = []
    # Filter out used destinations
    dialogs = [d for d in (dialogs or []) if d.get("id") not in used_ids]
    if q:
        ql = q.lower().strip()
        dialogs = [d for d in dialogs if ql in f"{d.get('title','')} {d.get('type','')}".lower()]
    return render_template(
        "components/_channel_options.html",
        dialogs=dialogs or [],
    )


@app.route("/_components/add_bot_form")
@pin_required
def htmx_add_bot_form():
    return render_template("components/_add_bot_form.html")


# ── HTMX Action Routes (mutations returning HTML partials) ──────────────────

@app.route("/_actions/forwarder/start", methods=["POST"])
@pin_required
@csrf_required
def htmx_start():
    core.start_forwarding()
    # Return refreshed status panel
    return htmx_status()


@app.route("/_actions/forwarder/stop", methods=["POST"])
@pin_required
@csrf_required
def htmx_stop():
    core.stop_forwarding()
    return htmx_status()


@app.route("/_actions/core/start", methods=["POST"])
@pin_required
@csrf_required
def htmx_core_start():
    """Manual core start when AUTO_START_CORE is disabled."""
    global _core_started
    if not _core_started:
        core.start()
        _core_started = True
    return htmx_status()


@app.route("/_actions/dialogs/refresh", methods=["POST"])
@pin_required
@csrf_required
def htmx_refresh_dialogs():
    """Refetch dialogs for the selected account and re-render both dropdowns."""
    account_id = request.form.get("account_id", type=int)
    if account_id is None:
        account_id = request.args.get("account_id", type=int)
    try:
        core.refresh_dialogs(account_id=account_id)
    except Exception:
        pass
    try:
        dialogs = core.get_dialogs(account_id=account_id)
    except Exception:
        dialogs = []
    opts = render_template("components/_channel_options.html", dialogs=dialogs or [])
    # OOB-refresh the destination dropdown too, keeping the Alpine
    # conflict-check attributes on the target select.
    oob = (
        '<select id="dest-options" hx-swap-oob="true" name="dest_ids" multiple '
        'size="4" style="min-height:90px;" required '
        'x-ref="dest" @change="checkConflict()">{}</select>'
    ).format(opts)
    return opts + oob


@app.route("/_actions/pairs", methods=["POST"])
@pin_required
@csrf_required
def htmx_create_pair():
    data = request.form
    source_id = data.get("source_id", type=int)
    # dest_ids comes from the multiple select — Flask sends multiple values
    dest_ids = data.getlist("dest_ids", type=int)

    if not source_id or not dest_ids:
        return htmx_pairs()
    if source_id in dest_ids:
        # UI prevents this; enforce server-side too (would forward a channel
        # to itself otherwise).
        return htmx_pairs()

    fields = {
        "source_title": data.get("source_title", ""),
        "enabled": True,
        "include_text": bool(data.get("include_text")),
        "include_media": bool(data.get("include_media")),
        "skip_standalone_media": bool(data.get("skip_standalone_media")),
        "forward_as_link": bool(data.get("forward_as_link")),
        "price_augment": bool(data.get("price_augment")),
        "forward_via_bot": bool(data.get("forward_via_bot")),
        "llm_enabled": bool(data.get("llm_enabled")),
        "augment_symbol": data.get("augment_symbol") or None,
        "filter_type": data.get("filter_type", "none"),
        "account_id": data.get("account_id", type=int) or None,
        "bot_token_id": data.get("bot_token_id", type=int) or None,
    }

    # Resolve source/dest titles from the SELECTED account's dialogs so the
    # names shown in the rules table match the account the pair is bound to.
    try:
        dialogs = core.get_dialogs(account_id=fields["account_id"])
    except Exception:
        dialogs = []
    src = next((d for d in (dialogs or []) if d.get("id") == source_id), None)
    fields["source_title"] = src.get("title", "") if src else data.get("source_title", "")

    for dest_id in dest_ids:
        try:
            dst = next((d for d in (dialogs or []) if d.get("id") == dest_id), None)
            core.create_pair(source_id, dest_id, dest_title=dst.get("title", "") if dst else "", **fields)
        except Exception as e:
            logger.warning("create_pair(%s, %s) failed: %s", source_id, dest_id, e)

    return htmx_pairs()


@app.route("/_actions/pairs/<int:pair_id>/toggle", methods=["POST"])
@pin_required
@csrf_required
def htmx_toggle_pair(pair_id):
    try:
        core.toggle_pair(pair_id)
    except Exception as e:
        logger.warning("toggle_pair(%s) failed: %s", pair_id, e)
    return htmx_pairs()


@app.route("/_actions/pairs/<int:pair_id>/delete", methods=["POST"])
@pin_required
@csrf_required
def htmx_delete_pair(pair_id):
    try:
        core.delete_pair(pair_id)
    except Exception as e:
        logger.warning("delete_pair(%s) failed: %s", pair_id, e)
    return htmx_pairs()


@app.route("/_actions/pairs/refresh_titles", methods=["POST"])
@pin_required
@csrf_required
def htmx_refresh_pair_titles():
    """Re-resolve source_title and dest_title for all pairs from dialog data.

    Fixes pairs whose titles are empty because the dialog cache wasn't
    available at creation time.  Dialogs are cached per account to avoid
    repeated lookups when many pairs share the same account.
    """
    pairs = core.get_pairs()
    updated = 0
    dialog_cache = {}
    for p in pairs:
        acct_id = p.get("account_id") or None
        if acct_id not in dialog_cache:
            try:
                dialog_cache[acct_id] = core.get_dialogs(account_id=acct_id)
            except Exception:
                dialog_cache[acct_id] = []
        dialogs = dialog_cache[acct_id]
        updates = {}
        src_id = p.get("source_chat_id")
        dst_id = p.get("dest_chat_id")
        if src_id is not None and not p.get("source_title"):
            d = next((x for x in dialogs if x.get("id") == src_id), None)
            if d and d.get("title"):
                updates["source_title"] = d["title"]
        if dst_id is not None and not p.get("dest_title"):
            d = next((x for x in dialogs if x.get("id") == dst_id), None)
            if d and d.get("title"):
                updates["dest_title"] = d["title"]
        if updates:
            try:
                core.update_pair(p["id"], **updates)
                updated += 1
            except Exception as e:
                logger.warning("refresh_titles pair %s failed: %s", p["id"], e)
    logger.info("refresh_titles: updated %d pair(s)", updated)
    return htmx_pairs()


@app.route("/_actions/accounts/<int:account_id>/delete", methods=["POST"])
@pin_required
@csrf_required
def htmx_delete_account(account_id):
    try:
        core._run(core._stop_account_client(account_id), timeout=10)
        core._run_db_sync(db.delete_account, account_id)
    except Exception as e:
        logger.warning("delete_account(%s) failed: %s", account_id, e)
    return htmx_accounts()


@app.route("/_actions/accounts/<int:account_id>/reset", methods=["POST"])
@pin_required
@csrf_required
def htmx_reset_account(account_id):
    try:
        core._run(core._stop_account_client(account_id), timeout=10)
        core._auth_state.pop(account_id, None)
        core._qr_state.pop(account_id, None)
        core._run_db_sync(db.reset_account, account_id)
    except Exception as e:
        logger.warning("reset_account(%s) failed: %s", account_id, e)
    return htmx_accounts()


@app.route("/_actions/bots/<int:bot_id>/delete", methods=["POST"])
@pin_required
@csrf_required
def htmx_delete_bot(bot_id):
    try:
        core._run_db_sync(db.delete_bot_token, bot_id)
    except Exception as e:
        logger.warning("delete_bot(%s) failed: %s", bot_id, e)
    return htmx_bots()


@app.route("/_actions/ctrader/<int:account_id>/delete", methods=["POST"])
@pin_required
@csrf_required
def htmx_delete_ctrader(account_id):
    try:
        core._run_db_sync(db.delete_ctrader_account, account_id)
    except Exception as e:
        logger.warning("delete_ctrader(%s) failed: %s", account_id, e)
    return htmx_ctrader()


@app.route("/_actions/ctrader/refresh", methods=["POST"])
@pin_required
@csrf_required
def htmx_ctrader_refresh():
    """Re-fetch account snapshots (balance, leverage, status) from Spotware."""
    refresh_errors = []
    try:
        rows = core._run_db_sync(db.list_ctrader_accounts) or []
        for row in rows:
            try:
                full = core._run_db_sync(db.get_ctrader_account, row["account_id"])
                if not full or not full.get("encrypted_payload") or not config.CTRADER_TOKEN_SECRET:
                    continue
                from crypto_vault import decrypt_payload
                payload = decrypt_payload(config.CTRADER_TOKEN_SECRET, full["encrypted_payload"])
                accounts = _fetch_ctrader_accounts(payload["access_token"])
                if not accounts:
                    refresh_errors.append(f"Account {row['account_id']}: no accounts returned (token may be expired)")
                    continue
                for acc in accounts:
                    account_id = acc.get("accountId")
                    if not account_id:
                        continue
                    core._run_db_sync(
                        db.update_ctrader_account_snapshot,
                        account_id,
                        broker_name=acc.get("brokerName", ""),
                        broker_title=acc.get("brokerTitle", ""),
                        account_number=acc.get("accountNumber"),
                        is_live=bool(acc.get("live", False)),
                        deposit_currency=acc.get("depositCurrency", ""),
                        balance_cents=acc.get("balance"),
                        money_digits=acc.get("moneyDigits", 2),
                        leverage=acc.get("leverage"),
                        account_type=acc.get("traderAccountType", ""),
                        account_status=acc.get("accountStatus", ""),
                    )
            except Exception as exc:
                logger.warning("ctrader refresh for account %s failed: %s", row.get("account_id"), exc)
                refresh_errors.append(f"Account {row.get('account_id')}: {exc}")
    except Exception as e:
        logger.warning("ctrader refresh failed: %s", e)
        refresh_errors.append(str(e))
    if refresh_errors:
        # Stash errors so the template can display them on next render
        app.logger.warning("ctrader refresh errors: %s", "; ".join(refresh_errors))
    return htmx_ctrader()


@app.route("/_actions/bots", methods=["POST"])
@pin_required
@csrf_required
def htmx_add_bot():
    token = request.form.get("token", "").strip()
    if token:
        try:
            core._run_db_sync(db.save_bot_token, token)
        except Exception as e:
            logger.warning("add_bot failed: %s", e)
    return htmx_bots()


@app.route("/_actions/bots/<int:bot_id>/verify", methods=["POST"])
@pin_required
@csrf_required
def htmx_verify_bot(bot_id):
    try:
        token = core._run_db_sync(db.get_decrypted_bot_token, bot_id)
        if token:
            import urllib.request, json as _json
            url = f"https://api.telegram.org/bot{token}/getMe"
            req = urllib.request.Request(url)
            with urllib.request.urlopen(req, timeout=10) as resp:
                data = _json.loads(resp.read().decode())
            if data.get("ok"):
                bot_info = data["result"]
                core._run_db_sync(
                    db.update_bot_token, bot_id,
                    bot_name=bot_info.get("first_name", ""),
                    bot_username=bot_info.get("username", ""),
                    verified=True,
                )
    except Exception as e:
        logger.warning("verify_bot(%s) failed: %s", bot_id, e)
    return htmx_bots()


@app.route("/_actions/config/signals", methods=["POST"])
@pin_required
@csrf_required
def htmx_save_signal_config():
    url = request.form.get("webhook_url", "").strip()
    secret = request.form.get("webhook_secret", "").strip()
    chat_id_raw = request.form.get("admin_chat_id", "").strip()
    try:
        core._run_db_sync(db.set_config, "SIGNAL_WEBHOOK_URL", url)
        if secret:
            core._run_db_sync(db.set_secret_config, "SIGNAL_WEBHOOK_SECRET", secret)
        if chat_id_raw:
            core._run_db_sync(db.set_config, "SIGNAL_ADMIN_CHAT_ID", chat_id_raw)
        config.SIGNAL_WEBHOOK_URL = url
        if secret:
            config.SIGNAL_WEBHOOK_SECRET = secret
        config.SIGNAL_ADMIN_CHAT_ID = int(chat_id_raw) if chat_id_raw else 0
        sp = getattr(core, "_signal_processor", None)
        if sp:
            sp.update_config(
                webhook_url=url,
                webhook_secret=secret or None,
                admin_chat_id=int(chat_id_raw) if chat_id_raw else None,
            )
    except Exception as e:
        logger.warning("save_signal_config failed: %s", e)
    return htmx_bots()


@app.route("/_actions/config/llm", methods=["POST"])
@pin_required
@csrf_required
def htmx_save_llm_config():
    data = request.form
    try:
        enabled = bool(data.get("enabled"))
        provider = data.get("provider", "openai-compatible").strip().lower()
        base_url = data.get("base_url", "").strip()
        model = data.get("model", "llama3.1-8b").strip()
        api_key = data.get("api_key", "").strip()
        webhook_url = data.get("webhook_url", "").strip()
        webhook_secret = data.get("webhook_secret", "").strip()
        min_confidence = float(data.get("min_confidence", 0.80))
        timeout = float(data.get("timeout", 12.0))
        max_tokens = int(data.get("max_tokens", 120))

        core._run_db_sync(db.set_config, "LLM_ENABLED", "1" if enabled else "0")
        core._run_db_sync(db.set_config, "LLM_PROVIDER", provider)
        core._run_db_sync(db.set_config, "LLM_BASE_URL", base_url or "https://api.cerebras.ai/v1")
        core._run_db_sync(db.set_config, "LLM_MODEL", model)
        core._run_db_sync(db.set_config, "LLM_MIN_CONFIDENCE", str(min_confidence))
        core._run_db_sync(db.set_config, "LLM_TIMEOUT", str(timeout))
        core._run_db_sync(db.set_config, "LLM_MAX_TOKENS", str(max_tokens))
        core._run_db_sync(db.set_config, "LLM_WEBHOOK_URL", webhook_url)
        if api_key:
            core._run_db_sync(db.set_secret_config, "LLM_API_KEY", api_key)
        if webhook_secret:
            core._run_db_sync(db.set_secret_config, "LLM_WEBHOOK_SECRET", webhook_secret)
        try:
            core.reload_llm_config()
        except Exception:
            pass
    except Exception as e:
        logger.warning("save_llm_config failed: %s", e)
    return htmx_llm()


@app.route("/_actions/llm/models", methods=["POST"])
@pin_required
@csrf_required
def htmx_llm_models():
    """Fill the model-field datalist from GET {base_url}/models.

    Uses the form's (possibly unsaved) base_url/api_key; empty fields fall
    back to the saved config. API key travels in the POST body, never a URL.
    On failure the reason lands in the test-result callout via OOB swap.
    """
    f = request.form
    res = core.llm_list_models(f.get("base_url", ""), f.get("api_key", ""))
    if not res.get("ok"):
        err = html.escape(res.get("error", ""))
        return (
            '<div id="llm-test-result" hx-swap-oob="innerHTML">'
            '<div class="test-err"><span class="badge err">Models unavailable</span>'
            f'<span class="hint">{err}</span>'
            '<span class="hint">Check the API key and base URL, then Test Connection.</span>'
            '</div></div>'
        )
    opts = "".join(f'<option value="{html.escape(m)}"></option>' for m in res["models"])
    return opts or '<option value=""></option>'


@app.route("/_actions/llm/test", methods=["POST"])
@pin_required
@csrf_required
def htmx_llm_test():
    """Ping the configured provider/model and show a structured result."""
    f = request.form
    res = core.llm_test_connection(
        f.get("base_url", ""), f.get("api_key", ""), f.get("model", "")
    )
    if res.get("ok"):
        return (
            '<div class="test-ok"><span class="badge ok">Connected</span>'
            f'<span class="hint">Model <code>{html.escape(res["model"])}</code> responded in {res["latency_ms"]}ms.</span></div>'
        )
    return (
        '<div class="test-err"><span class="badge err">Connection failed</span>'
        f'<span class="hint">{html.escape(res.get("error", ""))}</span>'
        '<span class="hint">Check the API key, base URL and model name, then Save and test again.</span></div>'
    )


@app.route("/_actions/ctrader/connect", methods=["POST"])
@pin_required
@csrf_required
def htmx_ctrader_connect():
    """Open cTrader OAuth consent URL in a new window."""
    try:
        url = _build_consent_url()
    except RuntimeError as exc:
        return f'<div class="toast err">{html.escape(str(exc))}</div>', 400
    return f'<script>window.open("{url}","_blank");</script>'


@app.route("/_actions/auth/login", methods=["POST"])
@rate_limit(max_requests=5, window=60)
def htmx_pin_login():
    if not config.WEB_PIN:
        return ""
    pin = request.form.get("pin", "").strip()
    if pin != config.WEB_PIN:
        # 200, not 403: htmx does not swap error responses by default, so an
        # error status would leave the user staring at an unchanged form.
        # The form must also carry the hx-on reload hook — after a failed
        # attempt this overlay replaces the original one from _pin_modal.html.
        return """<div id="pin-overlay" class="overlay">
  <div class="overlay-box">
    <h2>Enter PIN</h2>
    <form hx-post="/_actions/auth/login" hx-target="#pin-overlay" hx-swap="outerHTML" hx-on::after-request="if(event.detail.successful && !event.detail.xhr.responseText){ location.reload(); }">
      <div class="form-group">
        <label>Dashboard PIN</label>
        <input type="password" name="pin" placeholder="1234" autofocus required />
      </div>
      <p style="color:var(--err);font-size:0.85rem;">Invalid PIN</p>
      <button type="submit">Unlock</button>
    </form>
  </div>
</div>""", 200
    session["pin_authenticated"] = True
    return ""  # Empty 200 — HTMX does nothing, we reload via hx-on


# ── Settings (DB-backed config via web UI) ─────────────────────────────────────

@app.route("/_components/settings")
@pin_required
def htmx_settings():
    """Render the full settings panel with all config values."""
    values = config.get_all_env(db)
    return render_template(
        "components/_settings_panel.html",
        values=values,
        csrf=generate_csrf(),
    )


@app.route("/_actions/settings/save", methods=["POST"])
@pin_required
@csrf_required
def htmx_save_settings():
    """Save a single config var (sent as key + value)."""
    key = request.form.get("key", "").strip()
    value = request.form.get("value", "").strip()

    # Determine if it's a secret (Fernet-encrypted) or plain config.
    # Keep this in sync with config._db_secret_overrides so the web UI stores
    # secrets encrypted (tg_api_hash is the Telegram "app secret").
    secret_keys = {
        "tg_api_hash",
        "ctrader_client_secret", "ctrader_access_token", "ctrader_token_secret",
        "llm_api_key",
        "signal_bot_token",
        "signal_webhook_secret", "tg_password", "llm_webhook_secret",
    }

    if not key:
        return render_template(
            "components/_toast.html",
            msg="Missing key",
            ok=False,
        )

    # Skip if masked (unchanged secret). get_all_env() masks as "abcd…wxyz"
    # or, for short secrets, "••••••••" — both mean "keep current".
    if value and ("…" in value or "•" in value) and len(value) < 24:
        # User didn't modify the secret — keep existing
        values = config.get_all_env(db)
        return render_template(
            "components/_settings_panel.html",
            values=values,
            csrf=generate_csrf(),
        )

    try:
        if key in secret_keys:
            db.set_secret_config(key, value) if value else db.delete_config(key)
        else:
            db.set_config(key, value) if value else db.delete_config(key)

        # Reload config from DB
        config.patch_from_db(db)

        # Apply runtime changes for live systems
        if key in ("signal_webhook_url", "signal_webhook_secret"):
            sp = getattr(core, "_signal_processor", None)
            if sp:
                sp.update_config(
                    webhook_url=config.SIGNAL_WEBHOOK_URL,
                    webhook_secret=config.SIGNAL_WEBHOOK_SECRET or None,
                )
        if key.startswith("llm_"):
            try:
                core.reload_llm_config()
            except Exception:
                pass
    except Exception as e:
        logger.warning("save_setting(%s) failed: %s", key, e)

    values = config.get_all_env(db)
    return render_template(
        "components/_settings_panel.html",
        values=values,
        csrf=generate_csrf(),
    )


@app.route("/_actions/settings/reset", methods=["POST"])
@pin_required
@csrf_required
def htmx_reset_settings():
    """Reset a single config var to its env default."""
    key = request.form.get("key", "").strip()
    if key:
        db.delete_config(key)
        config.patch_from_db(db)
    values = config.get_all_env(db)
    return render_template(
        "components/_settings_panel.html",
        values=values,
        csrf=generate_csrf(),
    )


# ── Entry point ─────────────────────────────────────────────────────────────
if __name__ == "__main__":
    core.start()
    _core_started = True
    logger.info("Starting Flask on port %s", config.PORT)
    # Reverse proxies may pass traffic over IPv6 localhost (::1),
    # so bind to all addresses. (Only used when run directly; on cPanel
    # the app is served by Passenger via passenger_wsgi.py.)
    app.run(host="::", port=config.PORT, threaded=True)
