#!/usr/bin/env python3
"""Проверяет, что ProductSpec не требует несуществующих intake-колонок.

Проверка появилась после аварии F-89: код уже писал новое поле Ozon, а
подготовленная миграция ещё не была применена.  Поэтому источник требований —
ровно ``outputs.intake.fields`` ProductSpec, а режим работы только читает
``information_schema`` либо переданный JSON-снимок.  Скрипт ничего не меняет
в базе и не запускает миграции сам.
"""
from __future__ import annotations

import argparse
import json
import os
import re
from collections.abc import Iterable, Mapping, Sequence
from pathlib import Path
from typing import Any

import yaml


INTAKE_SCHEMA = "gwptd_intake"
COLLECTOR_ROOT = Path(__file__).resolve().parents[1]
REPO_ROOT = COLLECTOR_ROOT.parent
PRODUCTS_DIR = REPO_ROOT / "specs/products"
MIGRATIONS_DIR = COLLECTOR_ROOT / "kernel/schemas/migrations"

Schema = dict[str, set[str]]


class IntakeSchemaContractError(RuntimeError):
    """Контракт нельзя проверить без надёжного источника требований или схемы."""


def collect_required_columns(products_dir: Path = PRODUCTS_DIR) -> Schema:
    """Возвращает объединение target-колонок всех plain ``outputs.intake``."""
    if not products_dir.is_dir():
        raise IntakeSchemaContractError(
            f"каталог ProductSpec не найден: {products_dir}"
        )

    paths = sorted(
        path for path in products_dir.glob("*.yaml")
        if not path.name.startswith("._")
    )
    if not paths:
        raise IntakeSchemaContractError(
            f"в каталоге ProductSpec нет YAML-файлов: {products_dir}"
        )

    required: Schema = {}
    for path in paths:
        try:
            document = yaml.safe_load(path.read_text(encoding="utf-8")) or {}
        except (OSError, yaml.YAMLError) as exc:
            raise IntakeSchemaContractError(
                f"не удалось прочитать ProductSpec {path}: {exc}"
            ) from exc
        if not isinstance(document, Mapping):
            raise IntakeSchemaContractError(
                f"ProductSpec {path} должен быть YAML-объектом"
            )

        outputs = document.get("outputs")
        if outputs is None:
            continue
        if not isinstance(outputs, Mapping):
            raise IntakeSchemaContractError(
                f"ProductSpec {path}: outputs должен быть объектом"
            )
        intake = outputs.get("intake")
        if intake is None:
            continue
        if not isinstance(intake, Mapping):
            raise IntakeSchemaContractError(
                f"ProductSpec {path}: outputs.intake должен быть объектом"
            )

        table = intake.get("table")
        fields = intake.get("fields")
        if not isinstance(table, str) or not table:
            raise IntakeSchemaContractError(
                f"ProductSpec {path}: outputs.intake.table должен быть непустой строкой"
            )
        if not isinstance(fields, Mapping):
            raise IntakeSchemaContractError(
                f"ProductSpec {path}: outputs.intake.fields должен быть объектом"
            )
        if any(not isinstance(column, str) or not column for column in fields):
            raise IntakeSchemaContractError(
                f"ProductSpec {path}: имя колонки в outputs.intake.fields некорректно"
            )
        required.setdefault(table, set()).update(fields)

    return {table: required[table] for table in sorted(required)}


def _snapshot_columns(value: Any, table: str) -> set[str]:
    """Пустой список означает существующую таблицу без колонок, не её отсутствие."""
    if isinstance(value, Mapping):
        if set(value) != {"columns"}:
            raise IntakeSchemaContractError(
                f"снимок таблицы {table} должен содержать только ключ columns"
            )
        value = value["columns"]
    if not isinstance(value, list) or any(
        not isinstance(column, str) or not column for column in value
    ):
        raise IntakeSchemaContractError(
            f"снимок таблицы {table} должен быть списком непустых имён колонок"
        )
    return set(value)


def load_schema_snapshot(path: Path) -> Schema:
    """Читает JSON вида {"tables": {"table_name": ["column_name"]}}."""
    try:
        document = json.loads(path.read_text(encoding="utf-8"))
    except (OSError, json.JSONDecodeError) as exc:
        raise IntakeSchemaContractError(
            f"не удалось прочитать снимок схемы {path}: {exc}"
        ) from exc
    if not isinstance(document, Mapping):
        raise IntakeSchemaContractError("снимок схемы должен быть JSON-объектом")

    tables = document.get("tables")
    if not isinstance(tables, Mapping):
        raise IntakeSchemaContractError(
            "снимок схемы должен содержать объект tables"
        )
    if any(not isinstance(table, str) or not table for table in tables):
        raise IntakeSchemaContractError("снимок схемы содержит некорректное имя таблицы")
    return {
        table: _snapshot_columns(value, table)
        for table, value in sorted(tables.items())
    }


def _connection_settings() -> dict[str, Any]:
    names = (
        "V3_MYSQL_HOST",
        "V3_MYSQL_PORT",
        "V3_MYSQL_USER",
        "V3_MYSQL_PASSWORD",
    )
    missing = [name for name in names if not os.getenv(name)]
    if missing:
        raise IntakeSchemaContractError(
            "не заданы переменные окружения для --database: " + ", ".join(missing)
        )
    try:
        port = int(os.environ["V3_MYSQL_PORT"])
    except ValueError as exc:
        raise IntakeSchemaContractError(
            "V3_MYSQL_PORT должен быть целым числом"
        ) from exc
    if not 1 <= port <= 65535:
        raise IntakeSchemaContractError(
            "V3_MYSQL_PORT должен быть в диапазоне 1..65535"
        )
    return {
        "host": os.environ["V3_MYSQL_HOST"],
        "port": port,
        "user": os.environ["V3_MYSQL_USER"],
        "password": os.environ["V3_MYSQL_PASSWORD"],
        "charset": "utf8mb4",
        "autocommit": True,
    }


def _field(row: Mapping[str, Any], name: str) -> Any:
    """Поддерживает варианты DictCursor, возвращающие metadata в верхнем регистре."""
    for candidate in (name, name.upper(), name.lower()):
        if candidate in row:
            return row[candidate]
    raise IntakeSchemaContractError(f"ответ information_schema не содержит поле {name}")


def _select_rows(
    connection: Any, sql: str, params: Sequence[Any]
) -> list[dict[str, Any]]:
    if not sql.lstrip().upper().startswith("SELECT"):
        raise IntakeSchemaContractError("проверка схемы разрешает только SELECT")
    with connection.cursor() as cursor:
        cursor.execute(sql, params)
        return list(cursor.fetchall())


def read_database_schema(connection: Any) -> Schema:
    """Читает существующие base tables и колонки только из information_schema."""
    table_rows = _select_rows(
        connection,
        """
        SELECT table_name
        FROM information_schema.tables
        WHERE table_schema = %s AND table_type = 'BASE TABLE'
        ORDER BY table_name
        """,
        (INTAKE_SCHEMA,),
    )
    column_rows = _select_rows(
        connection,
        """
        SELECT table_name, column_name
        FROM information_schema.columns
        WHERE table_schema = %s
        ORDER BY table_name, ordinal_position
        """,
        (INTAKE_SCHEMA,),
    )
    actual: Schema = {
        str(_field(row, "table_name")): set() for row in table_rows
    }
    for row in column_rows:
        table = str(_field(row, "table_name"))
        actual.setdefault(table, set()).add(str(_field(row, "column_name")))
    return actual


def read_schema_from_database() -> Schema:
    """Подключается только после проверки явных переменных окружения V3 MySQL."""
    try:
        import pymysql
        import pymysql.cursors
    except ImportError as exc:
        raise IntakeSchemaContractError(
            "для режима --database требуется установленный пакет pymysql"
        ) from exc

    settings = _connection_settings()
    settings["cursorclass"] = pymysql.cursors.DictCursor
    connection = None
    try:
        connection = pymysql.connect(**settings)
        return read_database_schema(connection)
    finally:
        if connection is not None:
            connection.close()


def compare_schema(required: Mapping[str, Iterable[str]], actual: Mapping[str, Iterable[str]]) -> dict[str, list[Any]]:
    """Разделяет блокирующие отсутствия и информационные лишние колонки."""
    required_sets = {table: set(columns) for table, columns in required.items()}
    actual_sets = {table: set(columns) for table, columns in actual.items()}
    missing_tables = sorted(set(required_sets) - set(actual_sets))
    missing_columns = sorted(
        (table, column)
        for table in sorted(set(required_sets) & set(actual_sets))
        for column in required_sets[table] - actual_sets[table]
    )
    extra_columns = sorted(
        (table, column)
        for table in sorted(set(required_sets) & set(actual_sets))
        for column in actual_sets[table] - required_sets[table]
    )
    return {
        "missing_tables": missing_tables,
        "missing_columns": missing_columns,
        "extra_columns": extra_columns,
    }


def find_migrations(
    table: str,
    column: str,
    migrations_dir: Path = MIGRATIONS_DIR,
) -> list[Path]:
    """Ищет миграцию с конкретным ALTER TABLE ... ADD COLUMN, не по словам поодиночке."""
    schema = rf"`?{re.escape(INTAKE_SCHEMA)}`?"
    table_name = rf"`?{re.escape(table)}`?"
    column_name = rf"`?{re.escape(column)}`?"
    pattern = re.compile(
        rf"ALTER\s+TABLE\s+(?:{schema}\s*\.\s*)?{table_name}"
        rf"\s+ADD\s+COLUMN\s+(?:IF\s+NOT\s+EXISTS\s+)?{column_name}(?![A-Za-z0-9_])",
        re.IGNORECASE,
    )
    if not migrations_dir.is_dir():
        raise IntakeSchemaContractError(
            f"каталог миграций не найден: {migrations_dir}"
        )
    return [
        path
        for path in sorted(migrations_dir.glob("*.sql"))
        if not path.name.startswith("._")
        if pattern.search(path.read_text(encoding="utf-8"))
    ]


def _print_migration_hint(table: str, column: str, migrations_dir: Path) -> None:
    candidates = find_migrations(table, column, migrations_dir)
    if len(candidates) == 1:
        print(f"  Примените миграцию: {candidates[0].name}")
    elif not candidates:
        print("  Подходящая миграция с ADD COLUMN для этой колонки не найдена.")
    else:
        names = ", ".join(path.name for path in candidates)
        print(f"  Однозначно выбрать миграцию нельзя: {names}.")


def print_report(
    result: Mapping[str, list[Any]],
    migrations_dir: Path = MIGRATIONS_DIR,
    *,
    show_extra: bool = False,
) -> None:
    """Печатает ошибки до заметок, чтобы сломанный контракт был виден первым."""
    for table in result["missing_tables"]:
        print(f"ОШИБКА: в базе нет требуемой таблицы {INTAKE_SCHEMA}.{table}.")
        print("  Миграция для создания таблицы не найдена: поиск покрывает ADD COLUMN.")
    for table, column in result["missing_columns"]:
        print(
            "ОШИБКА: в базе нет требуемой колонки "
            f"{INTAKE_SCHEMA}.{table}.{column}."
        )
        _print_migration_hint(table, column, migrations_dir)
    extra = result["extra_columns"]
    if extra and not show_extra:
        # На живой схеме таких колонок 131: служебные id, collected_at,
        # synced_at и прочее, что заводит миграция, а не спека.  Печатать их
        # построчно — значит утопить в них две настоящие ошибки, а проверку
        # после каждой выкладки читают глазом.  Поэтому сводка, а подробности
        # по требованию.
        tables = sorted({table for table, _ in extra})
        print(
            f"ЗАМЕТКА: {len(extra)} колонок в {len(tables)} таблицах есть в базе, "
            "но ProductSpec их не требует — это не ошибка "
            "(подробности: --show-extra)."
        )
    elif extra:
        for table, column in extra:
            print(
                "ЗАМЕТКА: колонка "
                f"{INTAKE_SCHEMA}.{table}.{column} есть в базе, но ProductSpec её не требует."
            )
    if result["missing_tables"] or result["missing_columns"]:
        print("ИТОГ: контракт схемы приёма нарушен.")
    else:
        print("OK: все обязательные таблицы и колонки схемы приёма совпадают с ProductSpec.")


def _parser() -> argparse.ArgumentParser:
    parser = argparse.ArgumentParser(
        description="Проверка контракта схемы приёма ProductSpec только на чтение",
        allow_abbrev=False,
    )
    mode = parser.add_mutually_exclusive_group(required=True)
    mode.add_argument(
        "--database",
        action="store_true",
        help="прочитать gwptd_intake из information_schema",
    )
    mode.add_argument(
        "--schema-snapshot",
        type=Path,
        metavar="JSON",
        help='прочитать JSON-снимок вида {"tables": {"table": ["column"]}}',
    )
    parser.add_argument(
        "--show-extra",
        action="store_true",
        help="перечислить поимённо колонки, которых спеки не требуют (их 131)",
    )
    parser.add_argument(
        "--products-dir",
        type=Path,
        default=PRODUCTS_DIR,
        help="каталог specs/products (нужен для изолированных проверок)",
    )
    return parser


def main(argv: Sequence[str] | None = None) -> int:
    args = _parser().parse_args(argv)
    try:
        required = collect_required_columns(args.products_dir)
        actual = (
            read_schema_from_database()
            if args.database
            else load_schema_snapshot(args.schema_snapshot)
        )
        result = compare_schema(required, actual)
        print_report(result, show_extra=args.show_extra)
        return 1 if result["missing_tables"] or result["missing_columns"] else 0
    except Exception as exc:
        print(f"ОШИБКА: проверка контракта схемы приёма не выполнена: {exc}")
        return 2


if __name__ == "__main__":
    raise SystemExit(main())
