"""Панель OZPAY: API девайсов + статика мини-аппа.

Отдельный процесс от notify_server.py.
Проверки (check_balance / check_turnover / check_cards / full_check) и выпуск карты (add_card) идут в thread pool,
чтобы не блокировать event loop на ADB.
"""

from __future__ import annotations

import asyncio
import hmac
import html
import ipaddress
import json
import queue
import re
import sqlite3
import sys
import threading
import time
import urllib.request
from collections import deque
from urllib.parse import unquote
from pathlib import Path
from typing import Optional, Union

from aiogram import Bot, Dispatcher
from aiogram.client.default import DefaultBotProperties
from aiogram.enums import ParseMode
from aiogram.filters import Command, CommandStart
from aiogram.types import Message, ReplyKeyboardRemove
from fastapi import APIRouter, FastAPI, HTTPException, Request
from fastapi.middleware.cors import CORSMiddleware
from fastapi.responses import FileResponse, JSONResponse, Response
from fastapi.staticfiles import StaticFiles

from db_api import add_panel_worker, add_proxies, card_flag_keys, create_device, delete_device, delete_proxy, device_slot_value, find_device_by_ip_port, find_device_by_public_id, get_card_flag_defs, get_device, get_watch_settings, is_panel_worker, list_banned_ip_rows, list_banned_ips, list_devices, list_devices_by_ip, list_panel_workers, list_proxies, next_device_slot, parse_card_flags, parse_cards, proxy_counts, remove_banned_ip, remove_panel_worker, rename_device, restore_proxy, save_banned_ip, save_card_flag_defs, save_watch_settings, take_next_proxy, update_card_flags, update_device, update_password, used_device_ports
from config import ADMIN_TG_ID, BOT_TOKEN, DEVICE_CHAT_MAP, NOTIFY_BOT_TOKEN, PANEL_BOT_TOKEN, PANEL_LINK_TOKEN, SSL_CERTFILE, SSL_KEYFILE
from telegram_auth import (
    extract_user_from_init_data,
    make_panel_token,
    parse_panel_token,
    validate_init_data_any,
)
from main import ActionCancelled, DEFAULT_ADB_TCP_PORT, DOCKER_PS_END, DOCKER_PS_START, SLOT_ADB_PORT_STEP, add_card, add_device, cancel_login, capture_screen, check_balance, check_cards, check_login_state, check_turnover, clear_lk_session, full_check, install_device_apps, logout_lk, press_device_back, probe_adb, provision_new_device, rebuild_device, reboot_device, redroid_layout, screencap_png, set_action_cancel_hook
from ssh_cmd import (
    DEFAULT_PORT as SSH_DEFAULT_PORT,
    DEFAULT_USER as SSH_DEFAULT_USER,
    delete_host,
    get_host,
    list_hosts,
    rename_host,
    run_ssh_commands,
    upsert_host,
)
from userbot_service import runtime as userbot_runtime

WEBAPP_DIR = Path(__file__).parent / "webapp"
PORT = 5001
API_FLOOD_WINDOW_SEC = 5.0
API_FLOOD_MAX_REQUESTS = 10
SCREEN_FLOOD_WINDOW_SEC = 5.0
SCREEN_FLOOD_MAX_REQUESTS = 4
_LOOPBACK_IPS = {"127.0.0.1", "::1", "unknown", ""}

CHECKERS = {
    "balance": check_balance,
    "turnover": check_turnover,
    "cards": check_cards,
    "all": full_check,
}

app = FastAPI(title="OZPAY Panel")
app.add_middleware(
    CORSMiddleware,
    allow_origins=["*"],
    allow_methods=["*"],
    allow_headers=["*"],
)

api = APIRouter(prefix="/api")

_PANEL_BOT_MESSAGE_LIMIT = 4000
panel_bot: Optional[Bot] = None
panel_dp = Dispatcher()

_device_locks: dict[str, asyncio.Lock] = {}
_checking: set[str] = set()
_busy_kind: dict[str, str] = {}
_bg_tasks: set[asyncio.Task] = set()
_cancel_flags: dict[str, threading.Event] = {}
_tls = threading.local()
_orig_sleep = time.sleep

_log_lock = threading.Lock()
_log_lines: deque[dict] = deque(maxlen=5000)
_log_seq = 0
_log_partial = ""
_log_docker_ps = False

_rate_lock = threading.Lock()
_ip_hits: dict[str, deque[float]] = {}
_screen_hits: dict[str, deque[float]] = {}
_banned_ips: set[str] = set()


def _skip_captured_log(line: str) -> bool:
    if " /api/logs" in line:
        return True
    if "/screen" in line and " /api/" in line:
        return True
    return False


def _ingest_log(text: str) -> None:
    global _log_partial, _log_seq, _log_docker_ps
    if not text:
        return
    chunk = _log_partial + text
    parts = chunk.split("\n")
    _log_partial = parts.pop()
    if not parts:
        return
    with _log_lock:
        for line in parts:
            if _skip_captured_log(line):
                continue
            stripped = line.strip()
            hl = None
            if DOCKER_PS_START in stripped:
                _log_docker_ps = True
                hl = "docker_mark"
            elif DOCKER_PS_END in stripped:
                hl = "docker_mark"
                _log_docker_ps = False
            elif _log_docker_ps:
                hl = "docker"
            _log_seq += 1
            _log_lines.append({
                "id": _log_seq,
                "text": line,
                "device": getattr(_tls, "device_id", None),
                "hl": hl,
            })


class _StdTee:
    def __init__(self, original):
        self._original = original

    def write(self, data):
        if data is None:
            return 0
        if isinstance(data, bytes):
            try:
                text = data.decode("utf-8", "replace")
            except Exception:
                text = repr(data)
        else:
            text = str(data)
        try:
            self._original.write(text)
        except Exception:
            pass
        _ingest_log(text)
        return len(data) if not isinstance(data, int) else data

    def flush(self):
        try:
            self._original.flush()
        except Exception:
            pass

    def isatty(self):
        try:
            return self._original.isatty()
        except Exception:
            return False

    def __getattr__(self, name):
        return getattr(self._original, name)


if not isinstance(sys.stdout, _StdTee):
    sys.stdout = _StdTee(sys.stdout)
if not isinstance(sys.stderr, _StdTee):
    sys.stderr = _StdTee(sys.stderr)


def _action_cancel_hook() -> None:
    device_id = getattr(_tls, "device_id", None)
    if not device_id:
        return
    flag = _cancel_flags.get(device_id)
    if flag is None or not flag.is_set():
        return
    if not getattr(_tls, "cancel_logged", False):
        _tls.cancel_logged = True
        print(f"действие отменено ({device_id})")
    raise ActionCancelled("Действие отменено")


def _interruptible_sleep(seconds) -> None:
    device_id = getattr(_tls, "device_id", None)
    if not device_id:
        _orig_sleep(seconds)
        return
    end = time.monotonic() + max(0.0, float(seconds or 0))
    while True:
        _action_cancel_hook()
        remaining = end - time.monotonic()
        if remaining <= 0:
            return
        _orig_sleep(min(0.2, remaining))


time.sleep = _interruptible_sleep
set_action_cancel_hook(_action_cancel_hook)


def _begin_action(device_id: str, kind: str = "all") -> None:
    _checking.add(device_id)
    _busy_kind[device_id] = kind or "all"
    _cancel_flags[device_id] = threading.Event()


def _end_action(device_id: str) -> None:
    _checking.discard(device_id)
    _busy_kind.pop(device_id, None)
    _cancel_flags.pop(device_id, None)


def _call_device_action(device_id: str, fn, *args):
    _tls.device_id = device_id
    _tls.cancel_logged = False
    try:
        return fn(*args)
    except ActionCancelled:
        print(f"[{device_id}] действие отменено")
        raise
    except Exception as exc:
        print(f"[{device_id}] ошибка: {exc}")
        raise
    finally:
        _tls.device_id = None


# Активные сессии входа в ЛК (device_id -> LoginSession).
_login_sessions: dict[str, "LoginSession"] = {}
_ACTIVE_LOGIN_STATES = {"running", "awaiting_code", "verifying", "cancelling"}
CODE_WAIT_TIMEOUT = 300.0


class LoginCancelled(Exception):
    """Пользователь отменил вход из панели."""


class LoginSession:
    """Интерактивная сессия входа: add_device выполняется в отдельном потоке и на
    шаге ввода кода блокируется, ожидая действие из мини-аппа (код или повторную
    отправку). Пока ждём код, поток сам опрашивает состояние кнопки 'Получить новый
    код' на устройстве и кладёт его в `resend_available`, чтобы кнопка в панели
    была активна ровно тогда же, когда она активна в Ozon."""

    def __init__(self, device_id: str, number: str, password: str):
        self.device_id = device_id
        self.number = number
        self.password = password
        self.status = "running"  # running | awaiting_code | verifying | done | error
        self.method: Optional[str] = None
        self.target: Optional[str] = None
        self.error: Optional[str] = None
        self.device: Optional[dict] = None
        self.resend_available = False
        self.cancel_requested = False
        self._device = None  # ppadb device, выдаётся add_device на шаге кода
        self._action_q: "queue.Queue[tuple]" = queue.Queue()
        self.thread: Optional[threading.Thread] = None

    def _update_resend_available(self):
        if self._device is None:
            return
        try:
            from main import _get_new_code_button_enabled
            state = _get_new_code_button_enabled(self._device)
            if state is not None:
                self.resend_available = bool(state)
        except Exception:
            pass

    def _perform_resend(self):
        if self._device is None:
            return
        try:
            from main import detect_code_screen, dismiss_permission_dialog, wait_and_tap_get_new_code
            self.resend_available = False
            wait_and_tap_get_new_code(self._device, timeout=10.0)
            dismiss_permission_dialog(self._device)
            time.sleep(1.0)
            hint = detect_code_screen(self._device) or {}
            self.method = hint.get("method")
            self.target = hint.get("target")
        except Exception:
            pass

    def code_provider(self, ctx=None):
        ctx = ctx or {}
        self._device = ctx.get("device")
        hint = ctx.get("hint")
        if hint is None and ("method" in ctx or "target" in ctx):
            hint = ctx  # совместимость: ctx уже является хинтом
        hint = hint or {}
        self.method = hint.get("method")
        self.target = hint.get("target")
        if self.cancel_requested:
            raise LoginCancelled()
        self.status = "awaiting_code"

        deadline = time.time() + CODE_WAIT_TIMEOUT
        while time.time() < deadline:
            self._update_resend_available()
            try:
                action, payload = self._action_q.get(timeout=2.0)
            except queue.Empty:
                continue
            if action == "cancel":
                raise LoginCancelled()
            if action == "code":
                self.status = "verifying"
                return payload
            if action == "resend":
                self._perform_resend()
                deadline = time.time() + CODE_WAIT_TIMEOUT
        raise RuntimeError("Код не был введён вовремя")

    def submit_code(self, code: str):
        self._action_q.put(("code", code))

    def request_resend(self):
        self._action_q.put(("resend", None))

    def request_cancel(self):
        """Пометить сессию как отменяемую и разбудить поток, ждущий код."""
        self.cancel_requested = True
        self.status = "cancelling"
        self._action_q.put(("cancel", None))


def _run_login(session: LoginSession):
    try:
        add_device(
            session.device_id,
            session.number,
            session.password,
            code_provider=session.code_provider,
            press_get_new_code=False,
        )
        session.device = serialize_device(_require_device(session.device_id))
        session.status = "done"
    except LoginCancelled:
        session.status = "cancelling"
        try:
            cancel_login(session.device_id)
        except Exception as exc:  # noqa: BLE001 — навигация назад не должна ронять поток
            print(f"cancel_login({session.device_id}) failed: {exc}")
        session.status = "cancelled"
    except Exception as exc:  # noqa: BLE001 — прокидываем текст ошибки в UI
        if session.cancel_requested:
            try:
                cancel_login(session.device_id)
            except Exception as nav_exc:  # noqa: BLE001
                print(f"cancel_login({session.device_id}) failed: {nav_exc}")
            session.status = "cancelled"
        else:
            session.error = str(exc)
            session.status = "error"
    finally:
        _checking.discard(session.device_id)


def _lock_for(device_id: str) -> asyncio.Lock:
    lock = _device_locks.get(device_id)
    if lock is None:
        lock = asyncio.Lock()
        _device_locks[device_id] = lock
    return lock


def _spawn_device_action(
    device_id: str,
    fn,
    *args,
    kind: str = "all",
    notify_setup: bool = False,
    created: bool = False,
) -> None:
    _begin_action(device_id, kind)

    async def runner():
        cancelled = False
        error = None
        try:
            async with _lock_for(device_id):
                await asyncio.to_thread(_call_device_action, device_id, fn, *args)
        except ActionCancelled:
            cancelled = True
        except Exception as exc:
            error = exc
        finally:
            _end_action(device_id)
        if notify_setup and error is None and not cancelled:
            row = get_device(device_id)
            if row and row.get("needs_setup"):
                await asyncio.to_thread(_notify_setup_needed, device_id, created)

    task = asyncio.create_task(runner())
    _bg_tasks.add(task)
    task.add_done_callback(_bg_tasks.discard)


def _to_number(value) -> float:
    if value is None or value == "":
        return 0.0
    if isinstance(value, (int, float)):
        return float(value)
    text = str(value).replace("\xa0", " ").replace(" ", "").replace(",", ".")
    text = re.sub(r"[^\d.]", "", text)
    try:
        return float(text) if text else 0.0
    except ValueError:
        return 0.0


_LK_COLOR_RE = re.compile(r"^#[0-9a-f]{6}$")


def _lk_color(value) -> str:
    text = str(value or "").strip().lower()
    if _LK_COLOR_RE.fullmatch(text):
        return text
    return ""


def serialize_device(row: dict, host: Optional[dict] = None) -> dict:
    device_id = row.get("device") or ""
    ip = row.get("ip")
    port = row.get("port")
    linked = bool(row.get("number"))
    blocked = bool(row.get("blocked"))
    needs_setup = bool(row.get("needs_setup"))
    if device_id in _checking:
        status = "busy"
    elif needs_setup:
        status = "setup"
    elif not linked:
        status = "new"
    elif blocked:
        status = "blocked"
    elif ip and port:
        status = "online"
    else:
        status = "offline"

    if host is None:
        host = get_host(device_id)
    host = host or {}
    ssh_user = str(host.get("username") or "").strip() or SSH_DEFAULT_USER
    ssh_port_raw = host.get("port")
    try:
        ssh_port = int(ssh_port_raw) if ssh_port_raw not in (None, "") else SSH_DEFAULT_PORT
    except (TypeError, ValueError):
        ssh_port = SSH_DEFAULT_PORT

    slot = device_slot_value(row)

    return {
        "id": device_id,
        "name": (row.get("name") or "").strip(),
        "number": row.get("number") or "",
        "ip": ip or "",
        "slot": slot,
        "status": status,
        "linked": linked,
        "blocked": blocked,
        "needs_setup": needs_setup,
        "checking": device_id in _checking,
        "busy_kind": _busy_kind.get(device_id) or "",
        "balance": _to_number(row.get("balance")),
        "income": _to_number(row.get("income")),
        "outcome": _to_number(row.get("outcome")),
        "cards": parse_cards(row.get("cards"), row.get("card_flags")),
        "ssh_user": ssh_user,
        "ssh_port": ssh_port,
        "from_public": bool(row.get("from_public")),
        "public_id": str(row.get("public_id") or "").strip(),
        "color": _lk_color(row.get("color")),
    }


def _suggested_adb_port(ip: str, slot: int, exclude: str = "", *, strict: bool = True) -> int:
    used = used_device_ports(ip, exclude=exclude)
    try:
        slot = int(slot)
    except (TypeError, ValueError):
        slot = 1
    if slot < 1:
        slot = 1
    candidate = int(DEFAULT_ADB_TCP_PORT) + (slot - 1) * int(SLOT_ADB_PORT_STEP)
    if not 1 <= candidate <= 65535:
        candidate = int(DEFAULT_ADB_TCP_PORT)
    while candidate in used and candidate < 65535:
        candidate += 1
    if candidate in used and strict:
        raise HTTPException(status_code=409, detail="Нет свободного порта ADB на этом сервере")
    return candidate


def _server_summaries(rows: list[dict]) -> list[dict]:
    by_ip: dict[str, list] = {}
    for row in rows:
        ip = str(row.get("ip") or "").strip()
        if not ip:
            continue
        by_ip.setdefault(ip, []).append(row)
    result = []
    for ip in sorted(by_ip.keys()):
        devices_on_ip = by_ip[ip]
        next_slot = next_device_slot(ip, min_slot=2 if devices_on_ip else 1)
        result.append({
            "ip": ip,
            "count": len(devices_on_ip),
            "next_slot": next_slot,
            "suggested_port": _suggested_adb_port(ip, next_slot, strict=False),
        })
    return result


_DEVICE_ID_RE = re.compile(r"^[A-Za-z0-9][A-Za-z0-9._-]{0,63}$")
_DEVICE_ENDPOINT_RE = re.compile(r"(\d{1,3}(?:\.\d{1,3}){3}):(\d{4,5})\b")
_PROXY_HOSTPORT_RE = re.compile(r"^(\d{1,3}(?:\.\d{1,3}){3}|[A-Za-z0-9.-]+):(\d{1,5})$")
SSH_COMMAND_TIMEOUT = 60.0


def _without_device_port(text) -> str:
    return _DEVICE_ENDPOINT_RE.sub(r"\1", str(text or ""))


def _parse_device_id(value) -> str:
    device_id = str(value or "").strip()
    if not device_id:
        raise HTTPException(status_code=400, detail="Укажите имя девайса")
    if not _DEVICE_ID_RE.fullmatch(device_id):
        raise HTTPException(status_code=400, detail="Имя девайса: латиница, цифры, . _ -")
    return device_id


def _parse_ip(value) -> str:
    ip = str(value or "").replace(",", ".").strip()
    if not ip:
        raise HTTPException(status_code=400, detail="Укажите IP")
    return ip


def _parse_tcp_port(value, *, label: str, required: bool = True, default: Optional[int] = None) -> int:
    if value in (None, ""):
        if required:
            raise HTTPException(status_code=400, detail=f"Укажите {label.lower()}")
        if default is None:
            raise HTTPException(status_code=400, detail=f"Укажите {label.lower()}")
        return int(default)
    try:
        port = int(str(value).strip())
    except (TypeError, ValueError):
        raise HTTPException(status_code=400, detail=f"{label} должен быть числом")
    if not 1 <= port <= 65535:
        raise HTTPException(status_code=400, detail=f"{label}: 1–65535")
    return port


def _parse_ready_flag(value, default: bool = True) -> bool:
    if value in (None, ""):
        return default
    if isinstance(value, bool):
        return value
    if isinstance(value, (int, float)):
        return bool(value)
    text = str(value).strip().lower()
    if text in {"1", "true", "yes", "y", "da", "on", "ready", "да"}:
        return True
    if text in {"0", "false", "no", "n", "off", "not_ready", "нет"}:
        return False
    raise HTTPException(status_code=400, detail="Укажите, готов ли девайс")


def _parse_proxy_credentials(body: dict, *, port_fallback: bool = True) -> dict:
    parsed = _parse_proxy_url(str(body.get("proxy_url") or body.get("proxy") or ""))
    if not parsed:
        parsed = _parse_proxy_url(str(body.get("proxy_ip") or ""))
    if not parsed and port_fallback:
        parsed = _parse_proxy_url(str(body.get("ip") or ""))
    login = str(body.get("proxy_login") or body.get("login") or "").strip()
    password = body.get("proxy_password")
    if password is None and port_fallback:
        password = body.get("password")
    password = "" if password is None else str(password)
    ip_raw = body.get("proxy_ip")
    if ip_raw in (None, "") and port_fallback:
        ip_raw = body.get("ip")
    proxy_port_raw = body.get("proxy_port")
    if proxy_port_raw in (None, "") and port_fallback:
        proxy_port_raw = body.get("port")
    if parsed:
        login = login or parsed["username"]
        if password == "":
            password = parsed["password"]
        if not str(ip_raw or "").strip() or _parse_proxy_url(str(ip_raw or "")):
            ip_raw = parsed["ip"]
        if proxy_port_raw in (None, ""):
            proxy_port_raw = parsed["port"]
    ip = _parse_ip(ip_raw)
    proxy_port = _parse_tcp_port(proxy_port_raw, label="Порт прокси")
    if not login:
        raise HTTPException(status_code=400, detail="Укажите логин прокси")
    if not password:
        raise HTTPException(status_code=400, detail="Укажите пароль прокси")
    return {"login": login, "password": password, "ip": ip, "proxy_port": proxy_port}


def _notify_setup_needed(device_id: str, created: bool) -> None:
    action = "создан" if created else "пересобран"
    _notify_admin(
        (
            f"Девайс <b>{html.escape(device_id)}</b> {action}.\n\n"
            "Его нужно настроить. После настройки в панели нажмите «Готов к работе»."
        )
    )


def _parse_proxy_url(value: str) -> Optional[dict]:
    text = (value or "").strip()
    if not text or "@" not in text:
        return None
    rest = text
    scheme, sep, remainder = rest.partition("://")
    if sep and scheme and scheme.replace("+", "").replace(".", "").replace("-", "").isalnum():
        rest = remainder
    rest = rest.split("/", 1)[0].split("?", 1)[0].split("#", 1)[0]
    at = rest.rfind("@")
    if at <= 0:
        return None
    userinfo, hostport = rest[:at], rest[at + 1 :]
    colon = userinfo.find(":")
    if colon < 0:
        return None
    user = unquote(userinfo[:colon]).strip()
    password = unquote(userinfo[colon + 1 :])
    match = _PROXY_HOSTPORT_RE.fullmatch(hostport.strip())
    if not match:
        return None
    try:
        port = int(match.group(2))
    except (TypeError, ValueError):
        return None
    if not user or password == "" or not (1 <= port <= 65535):
        return None
    return {
        "ip": match.group(1),
        "port": port,
        "username": user,
        "password": password,
    }


def _proxy_dict_from_url(url: str) -> Optional[dict]:
    """Превращает ссылку прокси в {login, password, ip, proxy_port} для сборки."""
    parsed = _parse_proxy_url(str(url or ""))
    if not parsed:
        return None
    return {
        "login": parsed["username"],
        "password": parsed["password"],
        "ip": parsed["ip"],
        "proxy_port": parsed["port"],
    }


def _mask_proxy_url(url: str) -> str:
    """Прячет пароль в ссылке прокси для отдачи в панель."""
    text = str(url or "")
    parsed = _parse_proxy_url(text)
    if not parsed:
        return text
    scheme = ""
    rest = text.strip()
    head, sep, remainder = rest.partition("://")
    if sep:
        scheme = f"{head}://"
    return f"{scheme}{parsed['username']}:***@{parsed['ip']}:{parsed['port']}"


def _proxy_public(row: dict) -> dict:
    return {
        "id": row.get("id"),
        "url": _mask_proxy_url(row.get("url")),
        "used": bool(row.get("used")),
        "added_at": row.get("added_at"),
        "used_at": row.get("used_at"),
        "used_by": row.get("used_by"),
    }


def _proxies_payload() -> dict:
    rows = list_proxies()
    available = [_proxy_public(r) for r in rows if not r.get("used")]
    used = [_proxy_public(r) for r in rows if r.get("used")]
    return {
        "available": available,
        "used": used,
        "counts": proxy_counts(),
        "ok": True,
    }


def _wants_pool_proxy(body: dict) -> bool:
    val = body.get("use_pool")
    if val is None:
        val = body.get("proxy_pool")
    if val is None:
        return False
    if isinstance(val, bool):
        return val
    if isinstance(val, (int, float)):
        return bool(val)
    return str(val).strip().lower() in {"1", "true", "yes", "on", "да", "pool"}


def _take_pool_proxy(device_id: str) -> dict:
    """Берёт свободную прокси из пула, парсит и помечает использованной."""
    while True:
        row = take_next_proxy(device_id)
        if not row:
            raise HTTPException(status_code=409, detail="Нет свободных прокси в пуле")
        proxy = _proxy_dict_from_url(row.get("url"))
        if proxy:
            return proxy
        # Битую ссылку удаляем из пула и берём следующую.
        delete_proxy(row.get("id"))


def _retarget_runtime(old_id: str, new_id: str) -> None:
    if old_id == new_id:
        return
    if old_id in _device_locks:
        _device_locks[new_id] = _device_locks.pop(old_id)
    if old_id in _checking:
        _checking.discard(old_id)
        _checking.add(new_id)
    if old_id in _busy_kind:
        _busy_kind[new_id] = _busy_kind.pop(old_id)
    if old_id in _cancel_flags:
        _cancel_flags[new_id] = _cancel_flags.pop(old_id)


def _require_device(device_id: str) -> dict:
    row = get_device(device_id)
    if not row:
        raise HTTPException(status_code=404, detail=f"Устройство '{device_id}' не найдено")
    return row


def _normalize_ip(ip: str) -> str:
    ip = (ip or "").strip()
    if ip.startswith("::ffff:"):
        ip = ip[7:]
    return ip


def _client_ip(request: Request) -> str:
    peer = _normalize_ip(request.client.host if request.client else "")
    if peer in _LOOPBACK_IPS:
        forwarded = (
            request.headers.get("x-real-ip")
            or request.headers.get("x-forwarded-for")
            or ""
        ).strip()
        if forwarded:
            return _normalize_ip(forwarded.split(",")[0])
    return peer or "unknown"


def _cors_headers(request: Request) -> dict[str, str]:
    origin = request.headers.get("origin") or "*"
    return {
        "Access-Control-Allow-Origin": origin,
        "Access-Control-Allow-Methods": "GET, POST, PATCH, DELETE, OPTIONS",
        "Access-Control-Allow-Headers": "Accept, Content-Type, Authorization, X-Telegram-Init-Data, X-Panel-Link-Token",
        "Vary": "Origin",
    }


def _banned_response(request: Request) -> JSONResponse:
    return JSONResponse(
        status_code=403,
        content={"detail": "IP заблокирован"},
        headers=_cors_headers(request),
    )


def _is_screen_api(path: str) -> bool:
    return path.startswith("/api/") and path.endswith("/screen")


def _is_counted_api(path: str) -> bool:
    if path == "/api":
        return True
    if not path.startswith("/api/"):
        return False
    if path == "/api/logs":
        return False
    if path.startswith("/api/userbot"):
        return False
    if "/check/" in path:
        return False
    if _is_screen_api(path):
        return False
    return True


def _load_banned_ips() -> None:
    try:
        _banned_ips.update(list_banned_ips())
    except Exception as exc:
        print(f"не удалось загрузить banned_ips: {exc}")


def _ban_ip(ip: str, reason: str) -> None:
    if ip in _LOOPBACK_IPS:
        return
    print(f"автобан IP {ip}: {reason}")
    _set_ip_banned(ip, True)


def _set_ip_banned(ip: str, banned: bool) -> bool:
    ip = (ip or "").strip()
    if not ip:
        return False
    with _rate_lock:
        if banned:
            _banned_ips.add(ip)
            _ip_hits.pop(ip, None)
            _screen_hits.pop(ip, None)
        else:
            _banned_ips.discard(ip)
    try:
        if banned:
            save_banned_ip(ip)
            return True
        return remove_banned_ip(ip)
    except Exception as exc:
        print(f"не удалось {'сохранить' if banned else 'снять'} бан {ip}: {exc}")
        return banned


def _parse_ban_ip(raw) -> str:
    text = _normalize_ip(str(raw or "").strip())
    if not text:
        raise HTTPException(status_code=400, detail="Укажите IP")
    if "/" in text:
        raise HTTPException(status_code=400, detail="Укажите конкретный IP, без маски")
    try:
        parsed = ipaddress.ip_address(text)
    except ValueError as exc:
        raise HTTPException(status_code=400, detail="Некорректный IP") from exc
    if parsed.is_loopback or text in _LOOPBACK_IPS:
        raise HTTPException(status_code=400, detail="localhost нельзя заблокировать")
    return parsed.compressed


def _banned_ips_payload() -> dict:
    return {"ips": list_banned_ip_rows(), "ok": True}


def _hit_and_maybe_ban(
    ip: str,
    *,
    hits_map: dict[str, deque[float]],
    max_requests: int,
    window_sec: float,
    reason: str,
) -> bool:
    """True, если IP уже в бане или только что превысил лимит."""
    if ip in _LOOPBACK_IPS:
        return False
    now = time.monotonic()
    just_banned = False
    with _rate_lock:
        if ip in _banned_ips:
            return True
        hits = hits_map.get(ip)
        if hits is None:
            hits = deque()
            hits_map[ip] = hits
        cutoff = now - window_sec
        while hits and hits[0] <= cutoff:
            hits.popleft()
        hits.append(now)
        if len(hits) > max_requests:
            _banned_ips.add(ip)
            hits_map.pop(ip, None)
            just_banned = True
    if just_banned:
        _ban_ip(ip, reason)
        return True
    return False


_RESERVED_API = {"health", "logs", "devices", "userbot", "me", "workers", "auth", "banned-ips", "card-flags"}
_WORKER_POST_TAILS = {
    ("reboot",),
    ("cancel",),
    ("back",),
    ("ready",),
    ("check", "balance"),
    ("check", "turnover"),
    ("check", "cards"),
    ("cards", "add"),
    ("cards", "flags"),
    ("color",),
}


def _as_tg_id(value) -> Optional[int]:
    try:
        return int(value)
    except (TypeError, ValueError):
        return None


def _panel_role(user_id) -> Optional[str]:
    uid = _as_tg_id(user_id)
    if uid is None:
        return None
    admin_id = _as_tg_id(ADMIN_TG_ID)
    if admin_id is not None and uid == admin_id:
        return "admin"
    if is_panel_worker(uid):
        return "worker"
    return None


def _format_lk_password_messages() -> list[str]:
    rows = list_devices()
    if not rows:
        return ["Личных кабинетов пока нет."]

    header = "<b>Пароли от ЛК</b>"
    blocks = []
    for row in rows:
        device = html.escape(str(row.get("device") or "—"))
        name = html.escape((str(row.get("name") or "")).strip() or "—")
        number = html.escape((str(row.get("number") or "")).strip() or "—")
        password = html.escape((str(row.get("password") or "")).strip() or "—")
        blocks.append(
            f"<b>{device}</b>\n"
            f"ЛК: {name}\n"
            f"Номер: <code>{number}</code>\n"
            f"Пароль: <code>{password}</code>"
        )

    messages: list[str] = []
    current = header
    for block in blocks:
        candidate = f"{current}\n\n{block}"
        if len(candidate) > _PANEL_BOT_MESSAGE_LIMIT and current != header:
            messages.append(current)
            current = f"{header}\n\n{block}"
        else:
            current = candidate
    messages.append(current)
    return messages


def _staff_message(message: Message) -> bool:
    if message.chat.type != "private" or message.from_user is None:
        return False
    return _panel_role(message.from_user.id) in {"admin", "worker"}


@panel_dp.message(CommandStart())
async def on_panel_start(message: Message) -> None:
    if not _staff_message(message):
        return
    await message.answer(
        "Команда /pass покажет пароли всех личных кабинетов.",
        reply_markup=ReplyKeyboardRemove(),
    )


@panel_dp.message(Command("pass"))
async def on_lk_passwords(message: Message) -> None:
    if not _staff_message(message):
        return
    try:
        parts = _format_lk_password_messages()
    except Exception as exc:
        await message.answer(f"Не удалось получить список: {html.escape(str(exc))}")
        return
    for index, part in enumerate(parts):
        await message.answer(part, reply_markup=ReplyKeyboardRemove() if index == 0 else None)


async def _run_panel_bot() -> None:
    if panel_bot is None:
        return
    try:
        await panel_dp.start_polling(panel_bot, allowed_updates=["message"])
    except Exception as exc:
        print(f"panel bot polling failed: {exc}")


def _worker_can(method: str, path: str) -> bool:
    method = (method or "").upper()
    parts = [p for p in (path or "").split("/") if p]
    if not parts or parts[0] != "api":
        return False
    rest = tuple(parts[1:])
    if method == "GET":
        if rest in ((), ("devices",), ("me",), ("card-flags",)):
            return True
        if len(rest) == 1 and rest[0] not in _RESERVED_API:
            return True
        return len(rest) == 2 and rest[0] not in _RESERVED_API and rest[1] == "screen"
    if method != "POST" or len(rest) < 2 or rest[0] in _RESERVED_API:
        return False
    return rest[1:] in _WORKER_POST_TAILS


def _bot_tokens() -> list[str]:
    return [PANEL_BOT_TOKEN, BOT_TOKEN, NOTIFY_BOT_TOKEN]


def _session_secret() -> str:
    return (BOT_TOKEN or NOTIFY_BOT_TOKEN or "ozpay-panel").strip()


def _is_public_api(path: str) -> bool:
    normalized = (path or "").rstrip("/") or "/"
    return normalized in {"/api/health", "/api/auth"}


def _is_import_public_api(path: str) -> bool:
    return (path or "").rstrip("/") == "/api/devices/import-public"


def _extract_link_token(request: Request) -> str:
    header = (request.headers.get("x-panel-link-token") or "").strip()
    if header:
        return header
    return _extract_bearer(request)


def _link_token_ok(request: Request) -> bool:
    expected = str(PANEL_LINK_TOKEN or "").strip()
    got = _extract_link_token(request)
    if not expected or not got or len(expected) != len(got):
        return False
    return hmac.compare_digest(expected, got)


def _extract_init_data(request: Request) -> str:
    header = (request.headers.get("x-telegram-init-data") or "").strip()
    if header:
        return header
    auth = (request.headers.get("authorization") or "").strip()
    if auth.lower().startswith("tma "):
        return auth[4:].strip()
    return ""


def _extract_bearer(request: Request) -> str:
    auth = (request.headers.get("authorization") or "").strip()
    if auth.lower().startswith("bearer "):
        return auth[7:].strip()
    return ""


def _resolve_user_from_init_data(init_data: str, fallback_id=None) -> dict:
    raw = str(init_data or "").strip()
    fallback = _as_tg_id(fallback_id)
    if not raw:
        raise ValueError("Откройте мини-приложение из бота")
    try:
        return validate_init_data_any(raw, _bot_tokens())
    except ValueError as hmac_error:
        try:
            user = extract_user_from_init_data(raw)
        except ValueError:
            raise hmac_error
        if fallback is not None and fallback != user["id"]:
            raise ValueError("ID пользователя не совпадает с данными Telegram")
        print(f"panel auth: подпись не сошлась, берём ID из initData {user['id']}: {hmac_error}")
        return user


def _auth_payload(user: dict, role: str) -> dict:
    uid = int(user["id"])
    return {
        "ok": True,
        "allowed": True,
        "role": role,
        "admin": role == "admin",
        "token": make_panel_token(uid, _session_secret()),
        "user": {
            "id": uid,
            "first_name": user.get("first_name") or "",
            "username": user.get("username") or "",
        },
    }


def _auth_denied(request: Request) -> Optional[JSONResponse]:
    path = request.url.path
    if _is_public_api(path) or not path.startswith("/api"):
        return None
    if _is_import_public_api(path):
        if request.method == "OPTIONS":
            return None
        if not _link_token_ok(request):
            return JSONResponse(
                status_code=401,
                content={"detail": "Нет доступа"},
                headers=_cors_headers(request),
            )
        request.state.panel_user = {"id": 0, "first_name": "public-panel"}
        request.state.panel_role = "admin"
        return None
    user = None
    bearer = _extract_bearer(request)
    if bearer:
        try:
            uid = parse_panel_token(bearer, _session_secret())
            user = {"id": uid}
        except ValueError as exc:
            return JSONResponse(
                status_code=401,
                content={"detail": str(exc)},
                headers=_cors_headers(request),
            )
    else:
        try:
            user = _resolve_user_from_init_data(_extract_init_data(request))
        except ValueError as exc:
            print(f"panel auth failed: {exc}")
            return JSONResponse(
                status_code=401,
                content={"detail": str(exc)},
                headers=_cors_headers(request),
            )
    role = _panel_role(user["id"])
    if role is None:
        print(f"panel access denied tg_id={user['id']}")
        return JSONResponse(
            status_code=403,
            content={"detail": f"Нет доступа к панели (ID {user['id']})"},
            headers=_cors_headers(request),
        )
    request.state.panel_user = user
    request.state.panel_role = role
    if role != "admin" and not _worker_can(request.method, path):
        return JSONResponse(
            status_code=403,
            content={"detail": "Недостаточно прав"},
            headers=_cors_headers(request),
        )
    return None


_load_banned_ips()


@app.on_event("startup")
async def _userbot_startup() -> None:
    global panel_bot
    await userbot_runtime.boot()
    token = (PANEL_BOT_TOKEN or "").strip()
    if not token:
        print("panel bot skipped: нет PANEL_BOT_TOKEN")
        return
    panel_bot = Bot(token=token, default=DefaultBotProperties(parse_mode=ParseMode.HTML))
    asyncio.create_task(_run_panel_bot())
    print("panel bot polling started")


@app.on_event("shutdown")
async def _userbot_shutdown() -> None:
    await userbot_runtime.shutdown()
    if panel_bot is not None:
        await panel_bot.session.close()


@app.middleware("http")
async def log_requests(request: Request, call_next):
    ip = _client_ip(request)
    path = request.url.path
    if request.method == "OPTIONS":
        return Response(
            status_code=204,
            headers={
                **_cors_headers(request),
                "Access-Control-Max-Age": "86400",
            },
        )
    denied = _auth_denied(request)
    if denied is not None:
        return denied
    if _is_public_api(path):
        response = await call_next(request)
        response.headers.update(_cors_headers(request))
        print(f"{request.method} {path} -> {response.status_code}")
        return response
    authed = getattr(request.state, "panel_role", None) in {"admin", "worker"}
    if not authed:
        if ip in _banned_ips:
            return _banned_response(request)
        if _is_screen_api(path) and _hit_and_maybe_ban(
            ip,
            hits_map=_screen_hits,
            max_requests=SCREEN_FLOOD_MAX_REQUESTS,
            window_sec=SCREEN_FLOOD_WINDOW_SEC,
            reason=f"больше {SCREEN_FLOOD_MAX_REQUESTS} запросов скринов за {int(SCREEN_FLOOD_WINDOW_SEC)} с",
        ):
            return _banned_response(request)
        if _is_counted_api(path) and _hit_and_maybe_ban(
            ip,
            hits_map=_ip_hits,
            max_requests=API_FLOOD_MAX_REQUESTS,
            window_sec=API_FLOOD_WINDOW_SEC,
            reason=f"больше {API_FLOOD_MAX_REQUESTS} запросов за {int(API_FLOOD_WINDOW_SEC)} с",
        ):
            return _banned_response(request)
    elif ip in _banned_ips and getattr(request.state, "panel_role", None) != "admin":
        return _banned_response(request)
    response = await call_next(request)
    response.headers.update(_cors_headers(request))
    path = request.url.path
    if path != "/api/logs" and not path.endswith("/screen"):
        print(f"{request.method} {path} -> {response.status_code}")
    return response


@api.get("/health")
async def health() -> dict:
    return {"status": "ok"}


@api.post("/auth")
async def api_auth(request: Request) -> dict:
    try:
        body = await request.json()
    except Exception:
        body = {}
    if not isinstance(body, dict):
        body = {}
    init_data = str(body.get("init_data") or body.get("initData") or _extract_init_data(request) or "").strip()
    fallback_id = body.get("user_id") if body.get("user_id") not in (None, "") else body.get("id")
    try:
        user = _resolve_user_from_init_data(init_data, fallback_id)
    except ValueError as exc:
        print(f"panel /auth failed: {exc}")
        raise HTTPException(status_code=401, detail=str(exc)) from exc

    role = _panel_role(user["id"])
    if role is None:
        print(f"panel /auth deny tg_id={user['id']}")
        return {
            "ok": False,
            "allowed": False,
            "role": None,
            "admin": False,
            "user": {"id": user["id"]},
            "detail": f"Нет доступа к панели (ID {user['id']})",
        }

    if body.get("first_name") and not user.get("first_name"):
        user["first_name"] = str(body.get("first_name") or "")
    if body.get("username") and not user.get("username"):
        user["username"] = str(body.get("username") or "")
    print(f"panel /auth allow tg_id={user['id']} role={role}")
    return _auth_payload(user, role)


@api.get("/me")
async def api_me(request: Request) -> dict:
    user = getattr(request.state, "panel_user", None) or {}
    role = getattr(request.state, "panel_role", "worker")
    return {
        "user": {
            "id": user.get("id"),
            "first_name": user.get("first_name") or "",
            "username": user.get("username") or "",
        },
        "role": role,
        "admin": role == "admin",
    }


def _require_admin(request: Request) -> None:
    if getattr(request.state, "panel_role", None) != "admin":
        raise HTTPException(status_code=403, detail="Недостаточно прав")


def _notify_admin(text: str) -> None:
    admin_id = _as_tg_id(ADMIN_TG_ID)
    token = (BOT_TOKEN or NOTIFY_BOT_TOKEN or "").strip()
    if admin_id is None or not token:
        print("admin notify skipped: нет ADMIN_TG_ID или токена бота")
        return
    payload = json.dumps(
        {
            "chat_id": admin_id,
            "text": text,
            "parse_mode": "HTML",
            "disable_web_page_preview": True,
        },
        ensure_ascii=False,
    ).encode("utf-8")
    req = urllib.request.Request(
        f"https://api.telegram.org/bot{token}/sendMessage",
        data=payload,
        headers={"Content-Type": "application/json"},
        method="POST",
    )
    try:
        with urllib.request.urlopen(req, timeout=15) as resp:
            resp.read()
    except Exception as exc:
        print(f"admin notify failed: {exc}")


@api.get("/workers")
async def api_list_workers(request: Request) -> dict:
    _require_admin(request)
    return {"workers": list_panel_workers()}


@api.post("/workers")
async def api_add_worker(request: Request) -> dict:
    _require_admin(request)
    body = await request.json()
    tg_id = _as_tg_id(body.get("tg_id") if isinstance(body, dict) else None)
    if tg_id is None:
        tg_id = _as_tg_id(body.get("id") if isinstance(body, dict) else None)
    if tg_id is None or tg_id <= 0:
        raise HTTPException(status_code=400, detail="Укажите Telegram ID")
    admin_id = _as_tg_id(ADMIN_TG_ID)
    if admin_id is not None and tg_id == admin_id:
        raise HTTPException(status_code=400, detail="Админ и так имеет доступ")
    if not add_panel_worker(tg_id):
        raise HTTPException(status_code=400, detail="Некорректный Telegram ID")
    return {"workers": list_panel_workers(), "ok": True}


@api.delete("/workers/{tg_id}")
async def api_delete_worker(tg_id: int, request: Request) -> dict:
    _require_admin(request)
    if not remove_panel_worker(tg_id):
        raise HTTPException(status_code=404, detail="Работник не найден")
    return {"workers": list_panel_workers(), "ok": True}


@api.get("/banned-ips")
async def api_list_banned_ips(request: Request) -> dict:
    _require_admin(request)
    return _banned_ips_payload()


@api.post("/banned-ips")
async def api_add_banned_ip(request: Request) -> dict:
    _require_admin(request)
    try:
        body = await request.json()
    except Exception:
        body = {}
    ip = _parse_ban_ip(body.get("ip") if isinstance(body, dict) else None)
    _set_ip_banned(ip, True)
    return _banned_ips_payload()


@api.delete("/banned-ips")
async def api_delete_banned_ip(request: Request) -> dict:
    _require_admin(request)
    raw = str(request.query_params.get("ip") or "").strip()
    parsed = _parse_ban_ip(raw)
    candidates = []
    for item in (raw, _normalize_ip(raw), parsed):
        if item and item not in candidates:
            candidates.append(item)
    removed = False
    for ip in candidates:
        if _set_ip_banned(ip, False):
            removed = True
    if not removed:
        raise HTTPException(status_code=404, detail="IP не в банлисте")
    return _banned_ips_payload()


@api.get("/proxies")
async def api_list_proxies(request: Request) -> dict:
    _require_admin(request)
    return _proxies_payload()


@api.post("/proxies")
async def api_add_proxies(request: Request) -> dict:
    _require_admin(request)
    try:
        body = await request.json()
    except Exception:
        body = {}
    if not isinstance(body, dict):
        body = {}
    raw = body.get("urls")
    if raw is None:
        raw = body.get("list")
    if raw is None:
        raw = body.get("text")
    if isinstance(raw, str):
        candidates = raw.splitlines()
    elif isinstance(raw, (list, tuple)):
        candidates = list(raw)
    else:
        candidates = []
    valid = []
    invalid = []
    for item in candidates:
        url = str(item or "").strip()
        if not url:
            continue
        if _parse_proxy_url(url):
            valid.append(url)
        else:
            invalid.append(url)
    if not valid and not invalid:
        raise HTTPException(status_code=400, detail="Добавьте хотя бы одну ссылку прокси")
    if not valid:
        raise HTTPException(status_code=400, detail="Формат ссылки: http://login:password@ip:port")
    result = add_proxies(valid)
    payload = _proxies_payload()
    payload.update({
        "added": result.get("added", 0),
        "skipped": result.get("skipped", 0),
        "invalid": len(invalid),
    })
    return payload


@api.delete("/proxies")
async def api_delete_proxy(request: Request) -> dict:
    _require_admin(request)
    raw = request.query_params.get("id")
    if not delete_proxy(raw):
        raise HTTPException(status_code=404, detail="Прокси не найдена")
    return _proxies_payload()


@api.post("/proxies/restore")
async def api_restore_proxy(request: Request) -> dict:
    _require_admin(request)
    try:
        body = await request.json()
    except Exception:
        body = {}
    proxy_id = body.get("id") if isinstance(body, dict) else None
    if proxy_id is None:
        proxy_id = request.query_params.get("id")
    if not restore_proxy(proxy_id):
        raise HTTPException(status_code=404, detail="Прокси не найдена")
    return _proxies_payload()


@api.get("/logs")
async def api_logs(after: int = 0, device: str = "") -> dict:
    device_id = (device or "").strip()
    with _log_lock:
        scanned = [item for item in _log_lines if item["id"] > after]
        nxt = scanned[-1]["id"] if scanned else after
        items = scanned
        if device_id:
            items = [item for item in scanned if item.get("device") == device_id]
    return {"lines": [_client_log_line(item) for item in items], "next": nxt}


def _client_log_line(item: dict) -> Union[dict, str]:
    hl = item.get("hl")
    text = item.get("text") or ""
    if hl not in {"docker", "docker_mark"}:
        text = _without_device_port(text)
    if not hl:
        return text
    return {"text": text, "hl": hl}


def _userbot_http_error(exc: BaseException) -> HTTPException:
    if isinstance(exc, ValueError):
        return HTTPException(status_code=400, detail=str(exc))
    return HTTPException(status_code=400, detail=str(exc) or "Ошибка userbot")


@api.get("/userbot")
async def api_userbot_status() -> dict:
    return userbot_runtime.snapshot()


@api.post("/userbot/login")
async def api_userbot_login(request: Request) -> dict:
    body = await request.json()
    try:
        return await userbot_runtime.request_code(str(body.get("phone") or ""))
    except (ValueError, RuntimeError) as exc:
        raise _userbot_http_error(exc) from exc


@api.post("/userbot/code")
async def api_userbot_code(request: Request) -> dict:
    body = await request.json()
    try:
        return await userbot_runtime.submit_code(str(body.get("code") or ""))
    except (ValueError, RuntimeError) as exc:
        raise _userbot_http_error(exc) from exc


@api.post("/userbot/password")
async def api_userbot_password(request: Request) -> dict:
    body = await request.json()
    try:
        return await userbot_runtime.submit_password(str(body.get("password") or ""))
    except (ValueError, RuntimeError) as exc:
        raise _userbot_http_error(exc) from exc


@api.post("/userbot/cancel")
async def api_userbot_cancel() -> dict:
    return await userbot_runtime.cancel_login()


@api.post("/userbot/logout")
async def api_userbot_logout() -> dict:
    try:
        return await userbot_runtime.logout()
    except RuntimeError as exc:
        raise _userbot_http_error(exc) from exc


def _default_alert_chat_id():
    for value in DEVICE_CHAT_MAP.values():
        if value is not None:
            try:
                return int(value)
            except (TypeError, ValueError):
                continue
    return None


def _watch_payload() -> dict:
    data = get_watch_settings()
    if data.get("alert_chat_id") is None:
        data["alert_chat_id"] = _default_alert_chat_id()
    return data


@api.get("/userbot/watch")
async def api_userbot_watch() -> dict:
    return _watch_payload()


@api.post("/userbot/watch")
async def api_userbot_watch_save(request: Request) -> dict:
    body = await request.json()
    if not isinstance(body, dict):
        raise HTTPException(status_code=400, detail="Некорректные настройки")
    save_watch_settings(body)
    return _watch_payload()


@api.get("/card-flags")
async def api_card_flags() -> dict:
    return {"flags": get_card_flag_defs()}


@api.post("/card-flags")
async def api_card_flags_save(request: Request) -> dict:
    _require_admin(request)
    body = await request.json()
    items = body.get("flags") if isinstance(body, dict) else None
    if not isinstance(items, list):
        raise HTTPException(status_code=400, detail="Некорректный список флагов")
    try:
        saved = save_card_flag_defs(items)
    except ValueError as exc:
        raise HTTPException(status_code=400, detail=str(exc)) from exc
    return {"flags": saved}


@api.get("/userbot/dialogs")
async def api_userbot_dialogs() -> dict:
    try:
        dialogs = await userbot_runtime.list_dialogs()
    except RuntimeError as exc:
        raise _userbot_http_error(exc) from exc
    return {"dialogs": dialogs}


@api.get("/devices")
@api.get("/")
async def api_list_devices() -> dict:
    try:
        rows = list_devices()
        hosts = {item.get("device"): item for item in list_hosts()}
    except Exception as exc:
        raise HTTPException(status_code=500, detail=f"Ошибка БД: {exc}") from exc
    return {"devices": [serialize_device(row, hosts.get(row.get("device")) or {}) for row in rows], "servers": _server_summaries(rows)}


@api.post("/devices")
@api.post("/")
async def api_create_device(request: Request) -> dict:
    body = await request.json()
    device_id = _parse_device_id(body.get("id") or body.get("device") or body.get("name"))
    host_id = str(body.get("host_id") or body.get("host_device") or "").strip()
    host_row = None
    host_ssh = None
    if host_id:
        _require_admin(request)
        host_row = get_device(host_id)
        if not host_row:
            raise HTTPException(status_code=404, detail=f"Устройство '{host_id}' не найдено")
        host_ssh = get_host(host_id) or {}
        ip = _parse_ip(host_row.get("ip"))
        ready = False
        server_login = (
            str(host_ssh.get("username") or "").strip()
            or str(body.get("server_login") or body.get("ssh_user") or "").strip()
            or SSH_DEFAULT_USER
        )
        server_port = _parse_tcp_port(
            host_ssh.get("port") if host_ssh.get("port") not in (None, "") else body.get("server_port", body.get("ssh_port")),
            label="Порт сервера",
            required=False,
            default=SSH_DEFAULT_PORT,
        )
        password = body.get("password")
        if password is None:
            password = body.get("ssh_password")
        if password in (None, ""):
            password = host_ssh.get("password") or ""
        if not password:
            raise HTTPException(status_code=400, detail="У сервера нет пароля SSH")
        password = str(password)
    else:
        ip = _parse_ip(body.get("ip"))
        ready = _parse_ready_flag(body.get("ready"), default=True)
        password = body.get("password")
        if password is None:
            password = body.get("ssh_password")
        if password is None or str(password) == "":
            raise HTTPException(status_code=400, detail="Укажите пароль сервера")
        password = str(password)
        server_login = (str(body.get("server_login") or body.get("ssh_user") or "")).strip() or SSH_DEFAULT_USER
        server_port = _parse_tcp_port(
            body.get("server_port", body.get("ssh_port")),
            label="Порт сервера",
            required=False,
            default=SSH_DEFAULT_PORT,
        )

    has_siblings = bool(list_devices_by_ip(ip))
    slot = next_device_slot(ip, min_slot=2 if (host_row or has_siblings) else 1)
    port_raw = body.get("port")
    if port_raw in (None, ""):
        port = _suggested_adb_port(ip, slot)
    else:
        port = _parse_tcp_port(port_raw, label="Порт девайса")

    use_pool = False
    proxy = None
    if not ready:
        _require_admin(request)
        use_pool = _wants_pool_proxy(body)
        if use_pool:
            if proxy_counts().get("available", 0) <= 0:
                raise HTTPException(status_code=409, detail="Нет свободных прокси в пуле")
        else:
            proxy = _parse_proxy_credentials(body, port_fallback=False)

    duplicate = find_device_by_ip_port(ip, port)
    if duplicate:
        raise HTTPException(
            status_code=409,
            detail=f"IP {ip} уже занят девайсом '{duplicate.get('device')}'",
        )

    if ready:
        try:
            await asyncio.to_thread(probe_adb, ip, port)
        except Exception as exc:
            raise HTTPException(status_code=400, detail=_without_device_port(exc)) from exc

    try:
        create_device(device_id, ip=ip, port=port, slot=slot)
    except sqlite3.IntegrityError:
        raise HTTPException(status_code=409, detail=f"Девайс '{device_id}' уже есть")
    except Exception as exc:
        raise HTTPException(status_code=500, detail=f"Не удалось создать девайс: {exc}") from exc

    try:
        upsert_host(device_id, ip, password, username=server_login, port=server_port)
    except Exception as exc:
        delete_device(device_id)
        raise HTTPException(status_code=500, detail=f"Не удалось сохранить SSH: {exc}") from exc

    if ready:
        return {"device": serialize_device(_require_device(device_id)), "servers": _server_summaries(list_devices())}

    if use_pool:
        try:
            proxy = _take_pool_proxy(device_id)
        except HTTPException:
            delete_device(device_id)
            delete_host(device_id)
            raise

    update_device(device_id, {"needs_setup": 1, "slot": slot})
    _spawn_device_action(
        device_id,
        provision_new_device,
        device_id,
        proxy["login"],
        proxy["password"],
        proxy["ip"],
        proxy["proxy_port"],
        kind="provision",
        notify_setup=True,
        created=True,
    )
    return {
        "device": serialize_device(_require_device(device_id)),
        "ok": True,
        "started": True,
        "slot": slot,
    }


def _pick_imported_device_id(wanted: str, ip: str, port: int, public_id: str) -> str:
    by_public = find_device_by_public_id(public_id)
    if by_public:
        return str(by_public.get("device") or wanted)
    duplicate = find_device_by_ip_port(ip, port)
    if duplicate:
        return str(duplicate.get("device") or wanted)
    if not get_device(wanted):
        return wanted
    base = f"pub-{wanted}"
    candidate = base[:64]
    if not get_device(candidate):
        return candidate
    n = 2
    while n < 1000:
        suffix = f"pub{n}-{wanted}"
        candidate = suffix[:64]
        if not get_device(candidate):
            return candidate
        n += 1
    raise HTTPException(status_code=409, detail=f"Не удалось подобрать имя для '{wanted}'")


def _import_snapshot_fields(body: dict, *, ip: str, port: int, slot: int, public_id: str) -> dict:
    lk_password = body.get("lk_password")
    if lk_password is None:
        lk_password = body.get("password")
    cards = body.get("cards")
    if isinstance(cards, (list, dict)):
        cards = json.dumps(cards, ensure_ascii=False)
    flags = body.get("card_flags")
    if isinstance(flags, (list, dict)):
        flags = json.dumps(flags, ensure_ascii=False)
    blocked = body.get("blocked")
    return {
        "ip": ip,
        "port": port,
        "slot": slot,
        "number": str(body.get("number") or ""),
        "password": "" if lk_password is None else str(lk_password),
        "name": str(body.get("name") or "").strip(),
        "balance": _to_number(body.get("balance")),
        "income": _to_number(body.get("income")),
        "outcome": _to_number(body.get("outcome")),
        "cards": "" if cards is None else str(cards),
        "card_flags": None if flags is None else str(flags),
        "blocked": 1 if blocked in (True, 1, "1", "true", "True") else 0,
        "needs_setup": 0,
        "from_public": 1,
        "public_id": public_id,
    }


@api.post("/devices/import-public")
async def api_import_public_device(request: Request) -> dict:
    if not _link_token_ok(request):
        raise HTTPException(status_code=401, detail="Нет доступа")
    try:
        body = await request.json()
    except Exception:
        body = {}
    if not isinstance(body, dict):
        body = {}

    public_id = _parse_device_id(body.get("public_id") or body.get("id") or body.get("device") or body.get("name"))
    wanted_id = _parse_device_id(body.get("id") or body.get("device") or public_id)
    ip = _parse_ip(body.get("ip"))
    port = _parse_tcp_port(body.get("port"), label="Порт девайса")
    try:
        slot = int(body.get("slot") or 1)
    except (TypeError, ValueError):
        slot = 1
    if slot < 1:
        slot = 1

    device_id = _pick_imported_device_id(wanted_id, ip, port, public_id)
    existing = get_device(device_id)
    for row in list_devices_by_ip(ip):
        other_id = str(row.get("device") or "")
        if other_id == device_id:
            continue
        if device_slot_value(row) != slot:
            continue
        try:
            other_port = int(row.get("port"))
        except (TypeError, ValueError):
            other_port = None
        if other_port != port:
            raise HTTPException(
                status_code=409,
                detail=f"Слот {slot} на {ip} уже занят девайсом '{other_id}'",
            )

    snapshot = _import_snapshot_fields(body, ip=ip, port=port, slot=slot, public_id=public_id)
    created = existing is None
    if created:
        try:
            create_device(
                device_id,
                ip=ip,
                port=port,
                number=snapshot["number"],
                password=snapshot["password"],
                name=snapshot["name"],
                balance=snapshot["balance"],
                income=snapshot["income"],
                outcome=snapshot["outcome"],
                cards=snapshot["cards"],
                slot=slot,
            )
        except sqlite3.IntegrityError:
            raise HTTPException(status_code=409, detail=f"Девайс '{device_id}' уже есть")
        except Exception as exc:
            raise HTTPException(status_code=500, detail=f"Не удалось создать девайс: {exc}") from exc
    if not update_device(device_id, snapshot):
        if created:
            delete_device(device_id)
            raise HTTPException(status_code=500, detail="Не удалось сохранить слот")
        if not get_device(device_id):
            raise HTTPException(status_code=500, detail="Не удалось сохранить слот")

    ssh_password = body.get("ssh_password")
    if ssh_password is None:
        ssh_password = body.get("server_password")
    ssh_password = "" if ssh_password is None else str(ssh_password)
    server_login = (str(body.get("server_login") or body.get("ssh_user") or "")).strip() or SSH_DEFAULT_USER
    server_port = _parse_tcp_port(
        body.get("server_port", body.get("ssh_port")),
        label="Порт сервера",
        required=False,
        default=SSH_DEFAULT_PORT,
    )
    host = get_host(device_id)
    if ssh_password or host:
        try:
            upsert_host(
                device_id,
                ip,
                ssh_password if ssh_password else None,
                username=server_login,
                port=server_port,
            )
        except Exception as exc:
            if created:
                delete_device(device_id)
            raise HTTPException(status_code=500, detail=f"Не удалось сохранить SSH: {exc}") from exc

    print(f"import-public: {public_id} -> {device_id} created={created} ip={ip} slot={slot}")
    return {
        "ok": True,
        "created": created,
        "id": device_id,
        "public_id": public_id,
        "device": serialize_device(_require_device(device_id)),
    }


@api.post("/devices/rebuild")
async def api_rebuild_device(request: Request) -> dict:
    _require_admin(request)
    body = await request.json()
    device_id = _parse_device_id(body.get("id") or body.get("device") or body.get("name"))
    _require_device(device_id)
    if device_id in _checking:
        raise HTTPException(status_code=409, detail="Девайс уже проверяется")

    if _wants_pool_proxy(body):
        proxy = _take_pool_proxy(device_id)
    else:
        proxy = _parse_proxy_credentials(body, port_fallback=True)
    login = proxy["login"]
    password = proxy["password"]
    ip = proxy["ip"]
    proxy_port = proxy["proxy_port"]

    _spawn_device_action(
        device_id,
        rebuild_device,
        device_id,
        login,
        password,
        ip,
        proxy_port,
        kind="rebuild",
        notify_setup=True,
        created=False,
    )
    return {
        "device": serialize_device(_require_device(device_id)),
        "ok": True,
        "started": True,
    }


@api.post("/devices/{device_id}/install-apps")
async def api_install_device_apps(device_id: str, request: Request) -> dict:
    _require_admin(request)
    _require_device(device_id)
    if device_id in _checking:
        raise HTTPException(status_code=409, detail="Девайс уже проверяется")
    _spawn_device_action(device_id, install_device_apps, device_id, kind="install")
    return {
        "device": serialize_device(_require_device(device_id)),
        "ok": True,
        "started": True,
    }


@api.get("/{device_id}")
async def api_get_device(device_id: str) -> dict:
    return {"device": serialize_device(_require_device(device_id))}


def _command_vars_payload(device_id: str) -> dict:
    row = _require_device(device_id)
    host = get_host(device_id) or {}
    ip = str(row.get("ip") or "").strip()
    name = str(row.get("name") or "").strip()
    port_raw = row.get("port")
    try:
        port_text = str(int(port_raw)) if port_raw not in (None, "") else ""
    except (TypeError, ValueError):
        port_text = str(port_raw or "").strip()
    ssh_user = str(host.get("username") or "").strip() or SSH_DEFAULT_USER
    ssh_port_raw = host.get("port")
    try:
        ssh_port = int(ssh_port_raw) if ssh_port_raw not in (None, "") else SSH_DEFAULT_PORT
    except (TypeError, ValueError):
        ssh_port = SSH_DEFAULT_PORT
    endpoint = f"{ip}:{port_text}" if ip and port_text else ip
    slot = device_slot_value(row)
    layout = redroid_layout(slot)
    vars_map = {
        "device": device_id,
        "name": name or device_id,
        "device_ip": ip,
        "device_port": port_text,
        "ip": ip,
        "port": port_text,
        "device_endpoint": endpoint,
        "slot": str(slot),
        "redroid_dir": layout["dir"],
        "redroid_container": layout["container"],
        "ssh_user": ssh_user,
        "ssh_login": ssh_user,
        "server_login": ssh_user,
        "ssh_port": str(ssh_port),
        "server_port": str(ssh_port),
        "proxy_login": "",
        "proxy_user": "",
        "proxy_password": "",
        "proxy_ip": "",
        "proxy_port": "",
    }
    return {"id": device_id, "name": name, "vars": vars_map}


@api.get("/{device_id}/command-vars")
async def api_device_command_vars(device_id: str, request: Request) -> dict:
    _require_admin(request)
    return _command_vars_payload(device_id)


@api.post("/{device_id}/update")
async def api_update_device(device_id: str, request: Request) -> dict:
    row = _require_device(device_id)
    if device_id in _checking:
        raise HTTPException(status_code=409, detail="Дождитесь окончания проверки")

    body = await request.json()
    new_id = _parse_device_id(body.get("id") or body.get("device") or device_id)
    ip = _parse_ip(body.get("ip") if "ip" in body else row.get("ip"))
    port_raw = body.get("port") if "port" in body else None
    if port_raw in (None, ""):
        port = _parse_tcp_port(row.get("port"), label="Порт девайса")
    else:
        port = _parse_tcp_port(port_raw, label="Порт девайса")

    password = body.get("password")
    if password is None:
        password = body.get("ssh_password")
    password = "" if password is None else str(password)

    host = get_host(device_id)
    server_login = (str(body.get("server_login") or body.get("ssh_user") or "")).strip()
    if not server_login:
        server_login = str((host or {}).get("username") or "").strip() or SSH_DEFAULT_USER
    server_port = _parse_tcp_port(
        body.get("server_port", body.get("ssh_port")),
        label="Порт сервера",
        required=False,
        default=int((host or {}).get("port") or SSH_DEFAULT_PORT),
    )

    if new_id != device_id and get_device(new_id):
        raise HTTPException(status_code=409, detail=f"Девайс '{new_id}' уже есть")

    duplicate = find_device_by_ip_port(ip, port)
    if duplicate and duplicate.get("device") not in {device_id, new_id}:
        raise HTTPException(
            status_code=409,
            detail=f"IP {ip} уже занят девайсом '{duplicate.get('device')}'",
        )

    old_ip = str(row.get("ip") or "")
    try:
        old_port = int(row.get("port"))
    except (TypeError, ValueError):
        old_port = None
    endpoint_changed = old_ip != ip or old_port != port
    if endpoint_changed:
        try:
            await asyncio.to_thread(probe_adb, ip, port)
        except Exception as exc:
            raise HTTPException(status_code=400, detail=_without_device_port(exc)) from exc

    current_id = device_id
    if new_id != device_id:
        if not rename_device(device_id, new_id):
            raise HTTPException(status_code=500, detail="Не удалось переименовать девайс")
        try:
            rename_host(device_id, new_id)
        except ValueError as exc:
            rename_device(new_id, device_id)
            raise HTTPException(status_code=409, detail=str(exc)) from exc
        except Exception as exc:
            rename_device(new_id, device_id)
            raise HTTPException(status_code=500, detail=f"Не удалось переименовать SSH: {exc}") from exc
        _retarget_runtime(device_id, new_id)
        current_id = new_id
        host = get_host(current_id) or host

    fields = {"ip": ip, "port": port}
    if old_ip != ip:
        fields["slot"] = next_device_slot(ip, exclude=current_id, min_slot=2 if list_devices_by_ip(ip, exclude=current_id) else 1)
    update_device(current_id, fields)
    if not get_device(current_id):
        raise HTTPException(status_code=500, detail="Не удалось сохранить девайс")

    if host or password:
        try:
            upsert_host(
                current_id,
                ip,
                password if password else None,
                username=server_login,
                port=server_port,
            )
        except Exception as exc:
            raise HTTPException(status_code=500, detail=f"Не удалось сохранить SSH: {exc}") from exc

    return {"device": serialize_device(_require_device(current_id))}


@api.delete("/{device_id}")
async def api_delete_device(device_id: str) -> dict:
    _require_device(device_id)
    if device_id in _checking:
        raise HTTPException(status_code=409, detail="Дождитесь окончания проверки")
    if not delete_device(device_id):
        raise HTTPException(status_code=500, detail="Не удалось удалить девайс")
    try:
        delete_host(device_id)
    except Exception as exc:  # noqa: BLE001 — отсутствие SSH-записи не должно ломать удаление
        print(f"delete_host({device_id}) failed: {exc}")
    return {"ok": True, "id": device_id}


@api.post("/{device_id}/ready")
async def api_device_ready(device_id: str) -> dict:
    row = _require_device(device_id)
    if device_id in _checking:
        raise HTTPException(status_code=409, detail="Дождитесь окончания проверки")
    if not row.get("needs_setup"):
        raise HTTPException(status_code=400, detail="Девайс уже готов к работе")
    if not update_device(device_id, {"needs_setup": 0}):
        raise HTTPException(status_code=500, detail="Не удалось сохранить статус")
    return {"device": serialize_device(_require_device(device_id)), "ok": True}


@api.post("/{device_id}/logout")
async def api_logout_device(device_id: str) -> dict:
    row = _require_device(device_id)
    if device_id in _checking:
        raise HTTPException(status_code=409, detail="Девайс уже проверяется")
    if not row.get("number"):
        raise HTTPException(status_code=400, detail="ЛК не привязан")

    _begin_action(device_id)
    lock = _lock_for(device_id)
    cancelled = False
    try:
        async with lock:
            await asyncio.to_thread(_call_device_action, device_id, logout_lk, device_id)
    except ActionCancelled:
        cancelled = True
    except HTTPException:
        raise
    except Exception as exc:
        raise HTTPException(status_code=500, detail=str(exc)) from exc
    finally:
        _end_action(device_id)

    payload = {"device": serialize_device(_require_device(device_id))}
    if cancelled:
        payload["cancelled"] = True
    return payload


@api.post("/{device_id}/force-logout")
async def api_force_logout_device(device_id: str, request: Request) -> dict:
    _require_admin(request)
    row = _require_device(device_id)
    if not row.get("number"):
        raise HTTPException(status_code=400, detail="ЛК не привязан")
    clear_lk_session(device_id)
    return {"device": serialize_device(_require_device(device_id)), "ok": True}


@api.post("/{device_id}/cards/add")
async def api_add_card(device_id: str) -> dict:
    _require_device(device_id)
    if device_id in _checking:
        raise HTTPException(status_code=409, detail="Девайс уже проверяется")

    _begin_action(device_id)
    lock = _lock_for(device_id)
    result = None
    cancelled = False
    try:
        async with lock:
            result = await asyncio.to_thread(_call_device_action, device_id, add_card, device_id)
    except ActionCancelled:
        cancelled = True
    except HTTPException:
        raise
    except Exception as exc:
        raise HTTPException(status_code=500, detail=str(exc)) from exc
    finally:
        _end_action(device_id)

    if cancelled:
        return {"device": serialize_device(_require_device(device_id)), "cancelled": True}

    if not result or not result.get("success"):
        raise HTTPException(status_code=500, detail=(result or {}).get("message") or "Не удалось выпустить карту")

    return {"device": serialize_device(_require_device(device_id)), "message": result.get("message")}


@api.post("/{device_id}/color")
async def api_set_device_color(device_id: str, request: Request) -> dict:
    row = _require_device(device_id)
    if not str(row.get("number") or "").strip():
        raise HTTPException(status_code=400, detail="Цвет можно задать только у добавленного ЛК")
    try:
        body = await request.json()
    except Exception:
        body = {}
    if not isinstance(body, dict):
        body = {}
    raw = body.get("color")
    if raw in (None, "", "default"):
        color = ""
    else:
        color = _lk_color(raw)
        if not color:
            raise HTTPException(status_code=400, detail="Укажите цвет в формате #rrggbb")
    update_device(device_id, {"color": color})
    if not get_device(device_id):
        raise HTTPException(status_code=500, detail="Не удалось сохранить цвет")
    return {"device": serialize_device(_require_device(device_id))}


@api.post("/{device_id}/cards/flags")
async def api_update_card_flags(device_id: str, request: Request) -> dict:
    _require_device(device_id)
    body = await request.json()
    number = re.sub(r"\D", "", str(body.get("number") or ""))
    if not number:
        raise HTTPException(status_code=400, detail="Укажите номер карты")

    row = get_device(device_id)
    flags = parse_card_flags(row.get("card_flags") if row else None)
    keys = card_flag_keys()
    current = flags.get(number) or {key: False for key in keys}
    for key in keys:
        if key in body:
            current[key] = bool(body.get(key))
    if any(current.get(key) for key in keys):
        flags[number] = {key: bool(current.get(key)) for key in keys}
    else:
        flags.pop(number, None)

    if not update_card_flags(device_id, json.dumps(flags, ensure_ascii=False)):
        raise HTTPException(status_code=500, detail="Не удалось сохранить флаги карты")
    return {"device": serialize_device(_require_device(device_id))}


@api.post("/{device_id}/check/{kind}")
async def api_check_device(device_id: str, kind: str) -> dict:
    if kind not in CHECKERS:
        raise HTTPException(status_code=400, detail="Неизвестный тип проверки")

    _require_device(device_id)
    if device_id in _checking:
        raise HTTPException(status_code=409, detail="Девайс уже проверяется")

    _begin_action(device_id)
    lock = _lock_for(device_id)
    cancelled = False
    try:
        async with lock:
            await asyncio.to_thread(_call_device_action, device_id, CHECKERS[kind], device_id)
    except ActionCancelled:
        cancelled = True
    except HTTPException:
        raise
    except Exception as exc:
        raise HTTPException(status_code=500, detail=str(exc)) from exc
    finally:
        _end_action(device_id)

    payload = {"device": serialize_device(_require_device(device_id))}
    if cancelled:
        payload["cancelled"] = True
    return payload


@api.post("/{device_id}/login")
async def api_login_start(device_id: str, request: Request) -> dict:
    row = _require_device(device_id)
    if row.get("needs_setup"):
        raise HTTPException(status_code=409, detail="Сначала нажмите «Готов к работе»")

    body = await request.json()
    number = re.sub(r"\D", "", str(body.get("number") or ""))
    password = re.sub(r"\D", "", str(body.get("password") or ""))
    if not number:
        raise HTTPException(status_code=400, detail="Укажите номер телефона")
    if not password:
        raise HTTPException(status_code=400, detail="Укажите пароль (код-пароль)")

    existing = _login_sessions.get(device_id)
    if existing and existing.status in _ACTIVE_LOGIN_STATES:
        raise HTTPException(status_code=409, detail="Вход уже выполняется")
    if device_id in _checking:
        raise HTTPException(status_code=409, detail="Девайс занят проверкой")

    session = LoginSession(device_id, number, password)
    _login_sessions[device_id] = session
    _checking.add(device_id)
    session.thread = threading.Thread(target=_run_login, args=(session,), daemon=True)
    session.thread.start()
    return {"status": session.status}


@api.get("/{device_id}/login")
async def api_login_status(device_id: str) -> dict:
    session = _login_sessions.get(device_id)
    if not session:
        return {"status": "idle"}
    payload = {
        "status": session.status,
        "method": session.method,
        "target": session.target,
        "resend_available": session.resend_available,
    }
    if session.status == "error":
        payload["error"] = session.error
    if session.status == "done" and session.device:
        payload["device"] = session.device
    return payload


@api.post("/{device_id}/login/code")
async def api_login_code(device_id: str, request: Request) -> dict:
    session = _login_sessions.get(device_id)
    if not session or session.status != "awaiting_code":
        raise HTTPException(status_code=409, detail="Сейчас код не ожидается")

    body = await request.json()
    code = re.sub(r"\D", "", str(body.get("code") or ""))
    if len(code) != 6:
        raise HTTPException(status_code=400, detail="Код должен состоять из 6 цифр")

    session.submit_code(code)
    return {"status": "verifying"}


@api.post("/{device_id}/login/resend")
async def api_login_resend(device_id: str) -> dict:
    session = _login_sessions.get(device_id)
    if not session or session.status != "awaiting_code":
        raise HTTPException(status_code=409, detail="Сейчас код не ожидается")
    if not session.resend_available:
        raise HTTPException(status_code=409, detail="Кнопка ещё не активна")
    session.request_resend()
    return {"status": "awaiting_code"}


@api.post("/{device_id}/login/cancel")
async def api_login_cancel(device_id: str) -> dict:
    session = _login_sessions.get(device_id)
    if not session or session.status not in _ACTIVE_LOGIN_STATES:
        return {"status": "idle"}
    session.request_cancel()
    return {"status": session.status}


def _normalize_lk_number(raw) -> str:
    number = re.sub(r"\D", "", str(raw or ""))
    if len(number) == 11 and number[0] in ("7", "8"):
        number = number[1:]
    return number


@api.post("/{device_id}/login/refresh")
async def api_login_refresh(device_id: str) -> dict:
    _require_device(device_id)
    if device_id in _checking:
        raise HTTPException(status_code=409, detail="Девайс занят")
    try:
        state = await asyncio.to_thread(check_login_state, device_id)
    except Exception as exc:
        raise HTTPException(status_code=500, detail=str(exc)) from exc
    return state


@api.post("/{device_id}/login/claim")
async def api_login_claim(device_id: str, request: Request) -> dict:
    """Сохранить номер и код-пароль, если вход в ЛК на устройстве уже есть."""
    row = _require_device(device_id)
    if row.get("needs_setup"):
        raise HTTPException(status_code=409, detail="Сначала нажмите «Готов к работе»")
    if row.get("number"):
        raise HTTPException(status_code=409, detail="ЛК уже привязан")
    if device_id in _checking:
        raise HTTPException(status_code=409, detail="Девайс занят")

    body = await request.json()
    number = _normalize_lk_number(body.get("number"))
    password = re.sub(r"\D", "", str(body.get("password") or ""))
    if len(number) != 10:
        raise HTTPException(status_code=400, detail="Укажите номер телефона")
    if len(password) < 4:
        raise HTTPException(status_code=400, detail="Укажите пароль (код-пароль)")

    if not update_device(device_id, {"number": number, "password": password}):
        raise HTTPException(status_code=500, detail="Не удалось сохранить ЛК")
    return {"device": serialize_device(_require_device(device_id))}


@api.get("/{device_id}/screen")
async def api_device_screen(device_id: str) -> Response:
    """Снять PNG-скриншот текущего экрана устройства.

    Во время активного входа переиспользуем живое ADB-соединение из сессии
    (если оно уже получено на шаге ввода кода), иначе подключаемся заново.
    Эндпоинт не блокируется на `_checking`, чтобы экран можно было смотреть
    прямо во время входа и проверок.
    """
    _require_device(device_id)
    session = _login_sessions.get(device_id)
    try:
        if session is not None and session._device is not None and session.status in _ACTIVE_LOGIN_STATES:
            png = await asyncio.to_thread(screencap_png, session._device)
        else:
            png = await asyncio.to_thread(capture_screen, device_id)
    except Exception as exc:
        raise HTTPException(status_code=502, detail=f"Не удалось снять экран: {exc}") from exc
    return Response(
        content=png,
        media_type="image/png",
        headers={"Cache-Control": "no-store, max-age=0"},
    )


@api.post("/{device_id}/back")
async def api_device_back(device_id: str) -> dict:
    _require_device(device_id)
    if device_id in _checking:
        raise HTTPException(status_code=409, detail="Недоступно во время действия")
    try:
        await asyncio.to_thread(press_device_back, device_id)
    except Exception as exc:
        raise HTTPException(status_code=502, detail=str(exc)) from exc
    return {"ok": True}


@api.post("/{device_id}/cancel")
async def api_cancel_action(device_id: str) -> dict:
    _require_device(device_id)
    session = _login_sessions.get(device_id)
    if session is not None and session.status in _ACTIVE_LOGIN_STATES:
        session.request_cancel()
        return {"ok": True, "kind": "login"}
    if device_id not in _checking:
        raise HTTPException(status_code=409, detail="Нет активного действия")
    flag = _cancel_flags.get(device_id)
    if flag is not None:
        flag.set()
    try:
        await asyncio.to_thread(press_device_back, device_id)
    except Exception as exc:
        print(f"cancel: press_back failed: {exc}")
    return {"ok": True, "kind": "action"}


def _parse_ssh_commands(body: dict) -> list[str]:
    raw = body.get("commands")
    if raw is None:
        raw = body.get("command")
    if isinstance(raw, str):
        items = raw.splitlines()
    elif isinstance(raw, list):
        items = raw
    else:
        items = []
    return [str(item).strip() for item in items if str(item).strip()]


@api.post("/{device_id}/ssh")
async def api_ssh_command(device_id: str, request: Request) -> dict:
    _require_device(device_id)
    body = await request.json()
    commands = _parse_ssh_commands(body)
    if not commands:
        raise HTTPException(status_code=400, detail="Укажите команду")

    timeout = min(300.0, SSH_COMMAND_TIMEOUT * max(1, len(commands)))
    try:
        results = await asyncio.to_thread(
            run_ssh_commands,
            device_id,
            commands,
            timeout,
            False,
        )
    except ValueError as exc:
        raise HTTPException(status_code=400, detail=str(exc)) from exc
    except RuntimeError as exc:
        raise HTTPException(status_code=400, detail=str(exc)) from exc
    except Exception as exc:
        raise HTTPException(status_code=502, detail=str(exc)) from exc

    payload = {"results": results}
    if len(results) == 1:
        payload["stdout"] = results[0].get("stdout") or ""
        payload["stderr"] = results[0].get("stderr") or ""
        payload["exit_code"] = results[0].get("exit_code")
    return payload


@api.post("/{device_id}/reboot")
async def api_reboot_device(device_id: str) -> dict:
    _require_device(device_id)
    if device_id in _checking:
        raise HTTPException(status_code=409, detail="Девайс уже проверяется")

    _begin_action(device_id)
    lock = _lock_for(device_id)
    cancelled = False
    try:
        async with lock:
            await asyncio.to_thread(_call_device_action, device_id, reboot_device, device_id)
    except ActionCancelled:
        cancelled = True
    except HTTPException:
        raise
    except Exception as exc:
        raise HTTPException(status_code=500, detail=str(exc)) from exc
    finally:
        _end_action(device_id)

    payload = {"device": serialize_device(_require_device(device_id)), "ok": True}
    if cancelled:
        payload["cancelled"] = True
    return payload


@api.post("/{device_id}/password")
async def api_update_password(device_id: str, request: Request) -> dict:
    _require_device(device_id)
    body = await request.json()
    password = re.sub(r"\D", "", str(body.get("password") or ""))
    if len(password) < 4:
        raise HTTPException(status_code=400, detail="Код-пароль: минимум 4 цифры")
    if not update_password(device_id, password):
        raise HTTPException(status_code=500, detail="Не удалось сохранить пароль")
    return {"ok": True, "id": device_id}


app.include_router(api)


@app.get("/")
async def index() -> FileResponse:
    return FileResponse(
        WEBAPP_DIR / "index.html",
        headers={"Cache-Control": "no-store, max-age=0"},
    )


app.mount("/", StaticFiles(directory=WEBAPP_DIR), name="webapp")


if __name__ == "__main__":
    import uvicorn

    kwargs = {
        "host": "0.0.0.0",
        "port": PORT,
        "reload": False,
        "proxy_headers": True,
        "forwarded_allow_ips": "127.0.0.1",
    }

    cert = Path(SSL_CERTFILE) if SSL_CERTFILE else None
    key = Path(SSL_KEYFILE) if SSL_KEYFILE else None
    if cert and key and cert.exists() and key.exists():
        kwargs["ssl_certfile"] = str(cert)
        kwargs["ssl_keyfile"] = str(key)

    uvicorn.run("panel_server:app", **kwargs)
