"""rows_parsed обязан считать ЗАПИСАННОЕ, а не положенное в буфер.

Зачем этот файл отдельно от test_buffered_write.py: там проверяется, что строки
доезжают до базы, а здесь — что ЧИСЛО, которым прогон о себе отчитывается, не
врёт. Это разные вещи, и вторая опаснее: потерю строк заметят по данным, а
завышенный счётчик выглядит здоровым прогоном (получено = разобрано), пока
строки лежат в parser_error.

Разошлось это при переходе на пакетную запись: раньше вызов записи и был
записью, поэтому счёт сразу после вызова был верен. С буфером вызов больше не
пишет, и прежний счёт стал считать буферизацию.
"""
from __future__ import annotations

import sys
from pathlib import Path
from typing import Any
from unittest.mock import MagicMock

import pymysql.err
import pytest

sys.path.insert(0, str(Path(__file__).resolve().parent.parent))

from collectors.base import (BATCH_SIZE, MAX_STATEMENT_ROWS,  # noqa: E402
                             BaseCollector)


class _Run:
    """Минимальный дублёр прогона: нас интересуют только счётчики."""

    def __init__(self) -> None:
        self.rows_received = 0
        self.rows_parsed = 0
        self.rows_skipped = 0
        self.request_count = 0
        self.rate_limit_wait_sec = 0.0
        self.retry_count = 0

    def record_http(self, *_a, **_kw) -> None:  # pragma: no cover
        pass


class _CursorCtx:
    def __init__(self, cur):
        self._cur = cur

    def __enter__(self):
        return self._cur

    def __exit__(self, *_exc):
        return False


class _Collector(BaseCollector):
    """Коллектор с пакетной записью: одна входная строка → N DB-строк."""

    endpoint_code = "test.rows_parsed"
    mp_id = 99
    collector_class = "_Collector"
    # Объявляем явно: без этого признака коллектор уходит на построчный путь.
    supports_batch_write = True

    def __init__(self, rows_per_input: int = 1) -> None:  # noqa: D107
        super().__init__()
        self.run = _Run()
        self._rows_per_input = rows_per_input

    def fetch(self):  # pragma: no cover - строки подаются напрямую
        raise NotImplementedError
        yield

    def parse_and_store(self, payload_id: int, raw_payload: Any) -> None:
        raise NotImplementedError("пакетный путь не должен сюда падать")

    def _ensure_table(self, cur) -> None:
        pass

    def _insert_sql(self) -> str:
        return "INSERT INTO t (v) VALUES (%s)"

    def _row_to_params_list(self, payload_id: int, row: Any):
        return [(payload_id, (f"{row}-{i}",), None)
                for i in range(self._rows_per_input)]


@pytest.fixture()
def cursor():
    return MagicMock()


def _patch_cursor(monkeypatch, cur):
    import collectors.base as base_mod

    monkeypatch.setattr(base_mod, "get_cursor", lambda: _CursorCtx(cur))
    monkeypatch.setattr(base_mod, "log_parser_error", MagicMock())


def test_буферизация_не_засчитывает_строку(monkeypatch, cursor):
    """Пока строка в буфере — она НЕ разобрана. Это суть правки."""
    _patch_cursor(monkeypatch, cursor)
    collector = _Collector()

    parsed_now = collector._buffer_row(1, "a")

    assert parsed_now == 0, "буферизация не должна давать прибавку счётчику"
    assert collector.run.rows_parsed == 0
    assert len(collector._buffer) == 1


def test_счёт_появляется_только_после_записи(monkeypatch, cursor):
    _patch_cursor(monkeypatch, cursor)
    collector = _Collector()

    for i in range(3):
        collector._buffer_row(i, f"row{i}")
    assert collector.run.rows_parsed == 0, "до сброса засчитывать нечего"

    collector._flush_buffer()
    assert collector.run.rows_parsed == 3


def test_входная_строка_из_многих_db_строк_считается_один_раз(monkeypatch, cursor):
    """lamoda: один заказ → N строк. Счёт идёт по заказам, а не по строкам."""
    _patch_cursor(monkeypatch, cursor)
    collector = _Collector(rows_per_input=4)

    collector._buffer_row(1, "заказ")
    collector._flush_buffer()

    assert len(collector._buffer) == 0
    assert collector.run.rows_parsed == 1, "четыре DB-строки — это одна входная"


def test_упавшая_при_сбросе_строка_не_засчитывается(monkeypatch, cursor):
    """Главный случай: executemany упал, откат построчный, часть строк не легла.

    Раньше такая строка всё равно числилась разобранной."""
    _patch_cursor(monkeypatch, cursor)
    collector = _Collector()

    calls = {"n": 0}

    def execute(_sql, params=None):
        # Вторая строка при построчном откате падает.
        calls["n"] += 1
        if params and params[0] == "b-0":
            raise pymysql.err.IntegrityError(1062, "duplicate")

    cursor.executemany.side_effect = pymysql.err.IntegrityError(1062, "duplicate")
    cursor.execute.side_effect = execute

    for value in ("a", "b", "c"):
        collector._buffer_row(1, value)
    collector._flush_buffer()

    assert collector.run.rows_parsed == 2, "упавшая строка не разобрана"
    assert collector.run.rows_skipped == 1


def test_группа_не_разрывается_между_пачками(monkeypatch, cursor):
    """Входная строка не должна оказаться записанной наполовину.

    Иначе ответить «разобрана она или нет» будет нечем."""
    _patch_cursor(monkeypatch, cursor)
    collector = _Collector(rows_per_input=3)

    # Набиваем буфер почти до предела, затем добавляем группу, которая не влезает.
    total_inputs = BATCH_SIZE // 3 + 2
    for i in range(total_inputs):
        collector._buffer_row(i, f"row{i}")
    collector._flush_buffer()

    assert len(collector._buffer) == 0
    assert collector.run.rows_parsed == total_inputs
    # Ни одна пачка не должна была разрезать группу: число DB-строк в каждом
    # вызове executemany кратно трём.
    for call in cursor.executemany.call_args_list:
        assert len(call.args[1]) % 3 == 0, "группа разорвана между пачками"


def test_старый_синхронный_путь_считается_сразу(monkeypatch, cursor):
    """Коллекторы без пакетной записи пишут синхронно — счёт там же."""
    _patch_cursor(monkeypatch, cursor)

    class _Legacy(_Collector):
        # Признак снят ЯВНО: путь выбирается по нему, а не по наличию метода.
        supports_batch_write = False

        def parse_and_store(self, payload_id, raw_payload):
            pass  # запись состоялась

    collector = _Legacy()
    assert collector._buffer_row(1, "a") == 1, "синхронная запись даёт прибавку"
    assert collector.run.rows_parsed == 0, "прибавку делает вызывающий, не метод"


def test_огромная_группа_режется_на_куски_но_считается_раз(monkeypatch, cursor):
    """Группа больше предела оператора уходит несколькими executemany.

    Зачем: группа одной входной строки не разрывается между пачками, поэтому
    пачка может оказаться сильно больше BATCH_SIZE. Одним оператором такое
    упрётся в max_allowed_packet, а ошибка выглядит как «сервер закрыл
    соединение» и уводит расследование в сторону.

    При этом куски идут в ОДНОЙ транзакции, значит входная строка по-прежнему
    засчитывается один раз и целиком."""
    _patch_cursor(monkeypatch, cursor)
    huge = MAX_STATEMENT_ROWS * 2 + 7
    collector = _Collector(rows_per_input=huge)

    collector._buffer_row(1, "огромная")
    collector._flush_buffer()

    assert cursor.executemany.call_count == 3, "должно быть три куска"
    written = sum(len(c.args[1]) for c in cursor.executemany.call_args_list)
    assert written == huge, "ни одна строка не потеряна при нарезке"
    assert collector.run.rows_parsed == 1, "входная строка засчитана один раз"


def test_ошибка_внутри_подготовки_строк_не_прячется(monkeypatch, cursor):
    """NotImplementedError из глубины больше НЕ уводит на построчный путь.

    Раньше признаком возможности был перехват этого исключения, и настоящая
    ошибка внутри подготовки строк молча превращалась в откат к parse_and_store —
    у коллектора с пакетной записью там заглушка, и падение случалось далеко от
    причины."""
    _patch_cursor(monkeypatch, cursor)

    class _Broken(_Collector):
        def _row_to_params_list(self, payload_id, row):
            raise NotImplementedError("что-то внутри не готово")

    collector = _Broken()
    with pytest.raises(NotImplementedError, match="что-то внутри не готово"):
        collector._buffer_row(1, "a")


def test_ненормально_большая_группа_идёт_отдельной_пачкой(monkeypatch, cursor, caplog):
    """Группа сверх порога не копится рядом с чужими строками.

    Подтверждено проверяющим прогоном: группа не разрывается между пачками
    намеренно, и платой служит отсутствие потолка памяти — 50 000 строк одной
    группы держатся целиком. Настоящее лечение — генератор вместо списка, то
    есть переработка договора коллекторов. Пока её нет, такая группа хотя бы
    видна в журнале и не удлиняет общую транзакцию.
    """
    import logging

    from collectors.base import OVERSIZE_GROUP_ROWS

    _patch_cursor(monkeypatch, cursor)
    collector = _Collector(rows_per_input=1)

    # Сначала обычная строка — она НЕ должна уехать вместе с большой.
    collector._buffer_row(1, "обычная")
    assert len(collector._buffer) == 1

    collector._rows_per_input = OVERSIZE_GROUP_ROWS + 1
    with caplog.at_level(logging.WARNING):
        collector._buffer_row(2, "огромная")

    assert collector._buffer == [], "буфер обязан быть пуст после сброса"
    assert collector.run.rows_parsed == 2, "засчитаны обе входные строки"
    assert any("одна входная строка дала" in r.message for r in caplog.records), (
        "ненормальная группа обязана быть видна в журнале, а не проходить молча"
    )
