#!/usr/bin/env python3
"""
Отправитель отчётов MOS Activity Monitor.

Читает сгенерированные JSON-отчёты, складывает их в локальную очередь,
gzip-сжимает тело запроса и отправляет его в ingest HTTP endpoint.
Для HTTPS можно указать CA-файл через tls_ca_file, чтобы доверять
самоподписанному сертификату сервера.
"""

import gzip
import json
import re
import ssl
import sys
import time
import uuid
from dataclasses import dataclass
from pathlib import Path
from typing import Dict
from urllib import error, request

DEFAULT_CONFIG = "/etc/mos-activity-monitor/monitor.conf"
DEFAULT_OUTPUT_DIR = "/var/lib/mos-activity-monitor"
DEFAULT_SPOOL_DIR = "/var/spool/mos-activity-monitor"
DEFAULT_SERVER_ROOT = (
    "https://uds-test.mos.ru/uds-telemetry-metrics-ingest-api-service"
)
DEFAULT_ENDPOINT = "/api/telemetry/device-usage/v1/report"
DEFAULT_MAX_PAYLOAD = 1048576
DEFAULT_TLS_CA_FILE = ""


def read_config(path: str) -> Dict[str, str]:
    """
    Читает простой key=value конфиг агента.

    path: путь к monitor.conf.
    return: словарь параметров конфигурации без значений по умолчанию.
    """
    cfg = {}
    try:
        with open(path, "r", encoding="utf-8") as f:
            for raw in f:
                line = raw.strip()
                if not line or line.startswith("#") or "=" not in line:
                    continue
                key, value = line.split("=", 1)
                cfg[key.strip()] = value.strip()
    except FileNotFoundError:
        pass
    return cfg


def cfg_get(cfg: Dict[str, str], key: str, default: str) -> str:
    """
    Возвращает значение конфига или default, если ключ пустой/отсутствует.

    cfg: словарь конфигурации.
    key: имя параметра в конфиге.
    default: значение по умолчанию.
    return: итоговое строковое значение.
    """
    value = cfg.get(key, "")
    return value if value else default


def cfg_bool(cfg: Dict[str, str], key: str, default: bool) -> bool:
    """
    Разбирает bool-значения в формате, удобном для monitor.conf.

    cfg: словарь конфигурации.
    key: имя boolean-параметра.
    default: значение по умолчанию.
    return: True для 1/true/yes/on/да, иначе False.
    """
    value = cfg_get(cfg, key, "true" if default else "false").strip().lower()
    return value in ("1", "true", "yes", "on", "да")


def cfg_int(cfg: Dict[str, str], key: str, default: int) -> int:
    """
    Безопасно читает целое число из конфига.

    cfg: словарь конфигурации.
    key: имя числового параметра.
    default: значение по умолчанию при отсутствии ключа или ошибке разбора.
    return: целое число.
    """
    try:
        return int(cfg_get(cfg, key, str(default)))
    except ValueError:
        return default


def json_string_value(payload: bytes, key: str) -> str:
    """
    Достаёт строковое поле из JSON; regex оставлен как fallback для
    битых файлов.

    payload: содержимое JSON-файла в байтах.
    key: имя поля, которое нужно извлечь.
    return: строковое значение поля или пустая строка.
    """
    try:
        data = json.loads(payload.decode("utf-8"))
        value = data.get(key)
        return str(value) if value is not None else ""
    except Exception:
        pass

    try:
        text = payload.decode("utf-8", errors="replace")
        match = re.search(r'"' + re.escape(key) + r'"\s*:\s*"([^"]*)"', text)
        return match.group(1) if match else ""
    except Exception:
        return ""


def sanitize_filename(value: str) -> str:
    """
    Преобразует messageId/имя файла в безопасное имя для spool-файлов.

    value: исходная строка идентификатора.
    return: строка только из ASCII-букв, цифр, дефиса, подчёркивания и точки.
    """
    result = []
    for char in value:
        if char.isascii() and (char.isalnum() or char in "-_."):
            result.append(char)
        else:
            result.append("_")
    return "".join(result) or "unknown"


def read_uuid() -> str:
    """
    Генерирует X-Request-Id, предпочитая kernel uuid как в C++ версии.

    return: UUID-строка для заголовка X-Request-Id.
    """
    try:
        value = Path("/proc/sys/kernel/random/uuid").read_text(
            encoding="utf-8"
        ).strip()
        if value:
            return value
    except Exception:
        pass
    return str(uuid.uuid4())


def list_files(path: Path):
    """
    Возвращает отсортированный список обычных файлов в каталоге.

    path: каталог для просмотра.
    return: список Path-объектов; при отсутствии каталога
    возвращается пустой список.
    """
    try:
        return sorted(p for p in path.iterdir() if p.is_file())
    except FileNotFoundError:
        return []
    except NotADirectoryError:
        return []


def ensure_writable_dir(path: Path, purpose: str) -> bool:
    """
    Проверяет, что каталог доступен для записи текущему пользователю.

    path: каталог, который будет использоваться для spool-состояния.
    purpose: человекочитаемое описание назначения каталога.
    return: True, если удалось создать и удалить временный файл.
    """
    probe = path / f".write-test-{uuid.uuid4().hex}"
    try:
        probe.write_text("ok\n", encoding="utf-8")
        probe.unlink()
        return True
    except OSError as exc:
        print(
            f"Spool directory is not writable for {purpose}: {path}: {exc}",
            file=sys.stderr,
        )
        try:
            if probe.exists():
                probe.unlink()
        except OSError:
            pass
        return False


@dataclass
class SendResult:
    """
    Результат одной HTTP-отправки.

    http_code: код HTTP-ответа, если запрос дошёл до сервера.
    error_text: текст сетевой/локальной ошибки, если запрос не был выполнен.
    body: тело ответа сервера в текстовом виде.
    """

    http_code: int = 0
    error_text: str = ""
    body: str = ""


def gzip_compress(payload: bytes, level: int) -> bytes:
    """
    Сжимает тело отчёта в gzip-формат для HTTP-запроса.

    payload: исходный JSON в байтах.
    level: уровень gzip-сжатия от 1 до 9.
    return: gzip-сжатое тело запроса.
    """
    return gzip.compress(payload, compresslevel=level)


def build_ssl_context(tls_ca_file: str) -> ssl.SSLContext:
    """
    Создаёт TLS-контекст для HTTPS-отправки.

    tls_ca_file: путь к CA/самоподписанному сертификату сервера.
    return: SSLContext с проверкой цепочки и имени хоста.
    """
    if tls_ca_file:
        return ssl.create_default_context(cafile=tls_ca_file)
    return ssl.create_default_context()


def send_report(
    url: str,
    payload: bytes,
    device_id: str,
    gzip_level: int,
    timeout_seconds: int,
    tls_ca_file: str,
) -> SendResult:
    """
    Отправляет один отчёт в ingest endpoint.

    url: полный URL endpoint'а.
    payload: исходный JSON-отчёт в байтах.
    device_id: идентификатор устройства для заголовка X-Device-Id.
    gzip_level: уровень gzip-сжатия.
    timeout_seconds: общий timeout HTTP-запроса.
    tls_ca_file: путь к CA/самоподписанному сертификату сервера.
    return: SendResult с HTTP-кодом, телом ответа или текстом ошибки.
    """
    compressed = gzip_compress(payload, gzip_level)

    headers = {
        "Content-Type": "application/json;charset=utf-8",
        "Content-Encoding": "gzip",
        "Accept": "application/json",
        "X-Request-Id": read_uuid(),
        "X-Device-Id": device_id,
    }
    req = request.Request(url, data=compressed, headers=headers, method="POST")

    try:
        ssl_context = build_ssl_context(tls_ca_file)
        with request.urlopen(
            req, timeout=timeout_seconds, context=ssl_context
        ) as response:
            body = response.read().decode("utf-8", errors="replace")
            return SendResult(http_code=response.getcode(), body=body)
    except error.HTTPError as exc:
        body = exc.read().decode("utf-8", errors="replace")
        return SendResult(http_code=exc.code, body=body)
    except error.URLError as exc:
        return SendResult(error_text=str(exc.reason))
    except Exception as exc:
        return SendResult(error_text=str(exc))


def queue_new_reports(
    reports_dir: Path,
    queued_dir: Path,
    sent_dir: Path,
    max_uncompressed_bytes: int,
) -> bool:
    """
    Копирует новые отчёты в очередь, пропуская уже отправленные и уже queued.

    reports_dir: каталог с generated report JSON.
    queued_dir: каталог локальной очереди отправки.
    sent_dir: каталог marker-файлов для уже отправленных отчётов.
    max_uncompressed_bytes: максимальный размер несжатого JSON.
    return: True, если все подходящие отчёты были обработаны без
    локальных ошибок.
    """
    ok = True
    for path in list_files(reports_dir):
        base = path.name
        if not base.startswith("mos_activity_report_") or not base.endswith(
            ".json"
        ):
            continue
        if base == "latest_report.json":
            continue

        payload = path.read_bytes()
        if not payload:
            continue

        message_id = json_string_value(payload, "messageId")
        key = sanitize_filename(message_id or base)
        sent_marker = sent_dir / f"{key}.sent"
        queued_path = queued_dir / f"{key}.json"

        # sent marker и queued file вместе защищают от повторной
        # постановки в очередь.
        if sent_marker.exists() or queued_path.exists():
            continue

        if len(payload) > max_uncompressed_bytes:
            print(
                "Report is larger than local payload limit, not queued: "
                f"{path} "
                f"size={len(payload)} limit={max_uncompressed_bytes}",
                file=sys.stderr,
            )
            ok = False
            continue

        try:
            queued_path.write_bytes(payload)
            print(f"Queued report: {queued_path}")
        except OSError:
            print(
                f"Failed to queue report: {path} -> {queued_path}",
                file=sys.stderr,
            )
            ok = False
    return ok


def build_url(cfg: Dict[str, str]) -> str:
    """
    Собирает полный URL из server_root и report_endpoint.

    cfg: словарь конфигурации.
    return: полный URL для HTTP POST.
    """
    server_root = cfg_get(cfg, "server_root", DEFAULT_SERVER_ROOT).rstrip("/")
    endpoint = cfg_get(cfg, "report_endpoint", DEFAULT_ENDPOINT)
    if endpoint and not endpoint.startswith("/"):
        endpoint = "/" + endpoint
    return server_root + endpoint


def write_sent_marker(
    marker: Path,
    message_id: str,
    queued_path: Path,
    http_code: int,
) -> None:
    """
    Пишет короткий marker, по которому отчёт больше не попадёт в очередь.

    marker: путь marker-файла в sent_dir.
    message_id: messageId отправленного отчёта.
    queued_path: путь queued-файла, из которого была отправка.
    http_code: успешный HTTP-код ответа сервера.
    """
    marker.write_text(
        f"messageId={message_id}\n"
        f"queuedFile={queued_path}\n"
        f"httpCode={http_code}\n"
        f"sentAt={int(time.time())}\n",
        encoding="utf-8",
    )


def main(argv) -> int:
    """
    Основной oneshot-сценарий: подготовить очередь, отправить всё доступное.

    argv: аргументы командной строки; argv[1] может содержать путь к
    monitor.conf.
    return: process exit code, 0 при успешной обработке всей очереди.
    """
    config_path = argv[1] if len(argv) > 1 else DEFAULT_CONFIG
    cfg = read_config(config_path)

    if not cfg_bool(cfg, "send_enabled", True):
        print("Report sender is disabled: send_enabled=false")
        return 0

    output_dir = Path(cfg_get(cfg, "output_dir", DEFAULT_OUTPUT_DIR))
    reports_dir = Path(
        cfg_get(cfg, "reports_dir", str(output_dir / "reports"))
    )
    spool_dir = Path(cfg_get(cfg, "spool_dir", DEFAULT_SPOOL_DIR))
    queued_dir = spool_dir / "queued"
    sent_dir = spool_dir / "sent"

    gzip_level = cfg_int(cfg, "gzip_level", 6)
    if gzip_level < 1 or gzip_level > 9:
        gzip_level = 6

    timeout_seconds = cfg_int(cfg, "send_timeout_seconds", 30)
    if timeout_seconds <= 0:
        timeout_seconds = 30

    tls_ca_file = cfg_get(cfg, "tls_ca_file", DEFAULT_TLS_CA_FILE)

    max_payload = cfg_int(cfg, "send_max_payload_bytes", DEFAULT_MAX_PAYLOAD)
    if max_payload <= 0:
        max_payload = DEFAULT_MAX_PAYLOAD

    try:
        spool_dir.mkdir(parents=True, exist_ok=True)
        queued_dir.mkdir(parents=True, exist_ok=True)
        sent_dir.mkdir(parents=True, exist_ok=True)
    except OSError as exc:
        print(
            f"Failed to prepare spool directories under {spool_dir}: {exc}",
            file=sys.stderr,
        )
        return 1

    if not ensure_writable_dir(queued_dir, "queued reports"):
        return 1
    if not ensure_writable_dir(sent_dir, "sent markers"):
        return 1

    # Снача добираем новые report-файлы в spool, затем
    # отправляем всю очередь.
    queue_new_reports(reports_dir, queued_dir, sent_dir, max_payload)

    url = build_url(cfg)
    all_ok = True

    for queued_path in list_files(queued_dir):
        if queued_path.suffix != ".json":
            continue

        payload = queued_path.read_bytes()
        if not payload:
            continue
        if len(payload) > max_payload:
            print(
                "Queued report exceeds local payload limit, keeping for "
                f"manual handling: {queued_path}",
                file=sys.stderr,
            )
            all_ok = False
            continue

        message_id = json_string_value(payload, "messageId")
        device_id = json_string_value(payload, "deviceId") or cfg_get(
            cfg, "device_id", ""
        )
        if not device_id:
            print(
                "Report has empty deviceId and config device_id is empty: "
                f"{queued_path}",
                file=sys.stderr,
            )

        print(
            f"Sending report: {queued_path} "
            f"messageId={message_id or '<empty>'} "
            f"deviceId={device_id or '<empty>'} "
            f"url={url}"
        )

        result = send_report(
            url,
            payload,
            device_id,
            gzip_level,
            timeout_seconds,
            tls_ca_file,
        )
        if result.error_text:
            print(
                f"Send failed: error={result.error_text} body={result.body}",
                file=sys.stderr,
            )
            all_ok = False
            continue

        print(f"HTTP {result.http_code} response: {result.body}")
        if 200 <= result.http_code < 300:
            # Удаляем queued-файл только после успешного ответа сервера.
            key = sanitize_filename(message_id or queued_path.name)
            try:
                write_sent_marker(
                    sent_dir / f"{key}.sent",
                    message_id,
                    queued_path,
                    result.http_code,
                )
            except OSError as exc:
                print(
                    "Failed to write sent marker after successful send, "
                    "keeping queued file: "
                    f"{sent_dir / f'{key}.sent'}: {exc}",
                    file=sys.stderr,
                )
                all_ok = False
                continue
            try:
                queued_path.unlink()
            except OSError as exc:
                print(
                    "Failed to remove queued file after successful send: "
                    f"{queued_path}: {exc}",
                    file=sys.stderr,
                )
                all_ok = False
        else:
            all_ok = False

    return 0 if all_ok else 1


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