"""Legacy combined ProductSpec collector with parser/intake writes.

The independent raw process never imports this module.  It exists only for the
active combined collector during migration and preserves its historical public
``SpecCollector`` API.
"""
from __future__ import annotations

from datetime import date
from typing import Any

from .next_mapper import Mapper
from .mapping_helpers import explode_values, parent_val
from .runner import SpecAcquisitionCollector


class SpecCollector(SpecAcquisitionCollector):
    """Combined acquisition + intake parser retained until the old cutover."""

    collector_class = "SpecCollector"
    # The legacy process is intentionally still allowed to consult its typed
    # intake state for historical deduplication.  The raw acquisition runtime
    # must never inherit that behaviour.
    is_raw_acquisition = False

    def __init__(
        self,
        spec: dict[str, Any],
        max_items: int | None = None,
        **kwargs: Any,
    ) -> None:
        super().__init__(spec, max_items=max_items, **kwargs)
        self.intake_table = spec["intake"]["table"]
        self.required_intake_columns = tuple(spec["intake"]["fields"].keys())
        self._mapper = Mapper(spec["intake"])
        self._mappers_by_kind: dict[str, list[Mapper]] | None = None
        self._kind_field: str | None = None
        if spec.get("intakeByKind"):
            by_kind = spec["intakeByKind"]
            self._kind_field = by_kind["field"]
            self._mappers_by_kind = {
                kind: [Mapper(target) for target in targets]
                for kind, targets in by_kind["targets"].items()
            }

    def parse_and_store(self, payload_id: int, raw_payload: Any) -> None:
        # Recognised evidence marker for the currentMonthNotFoundOk 404 path:
        # no typed output, no mapper lookup and no parser error.
        if (
            isinstance(raw_payload, dict)
            and raw_payload.get(self._kind_field) == "expected_empty_404"
        ):
            return
        context = {
            "payload_id": payload_id,
            "TODAY": self.today,
            "mp_id": self.mp_id,
            "collected_at": self._collected_at,
        }
        if self._iter and isinstance(self._iter.get("date"), date):
            context["iter_date"] = self._iter["date"]
        if self._mappers_by_kind is not None:
            kind = str((raw_payload or {}).get(self._kind_field) or "")
            if kind not in self._mappers_by_kind and self._spec["intakeByKind"].get(
                "defaultKind"
            ):
                kind = self._spec["intakeByKind"]["defaultKind"]
            mappers = self._mappers_by_kind.get(kind)
            if mappers is None:
                raise RuntimeError(
                    f"{self.endpoint_code}: неизвестный kind «{kind}» (intakeByKind)"
                )
            targets = self._spec["intakeByKind"]["targets"][kind]
            for mapper, target_config in zip(mappers, targets):
                self._write_target(
                    mapper,
                    target_config.get("explode"),
                    raw_payload,
                    context,
                )
            return
        explode = self._spec["intake"].get("explode")
        self._write_target(self._mapper, explode, raw_payload, context)

    def _reused_payload_is_complete(self, data: Any) -> bool:
        """Accept cached data only when the mapper restores every source field.

        Поля, помеченные reuse_required: false, не участвуют в проверке:
        их отсутствие в сохранённом ответе не делает payload «неполным».
        Это позволяет добавлять извлечение новых полей из ответа без
        превращения ранее собранных ответов в недействительные.
        """
        try:
            items = self._extract_items(data)
        except Exception:  # noqa: BLE001
            return False
        if not items:
            return False
        expected = set()
        for name, field in (self._spec.get("intake") or {}).get("fields", {}).items():
            if not isinstance(field, dict):
                source = getattr(field, "json_path", None)
                require_reuse = True
            else:
                source = field.get("json_path")
                # По умолчанию True: поле, бывшее в спеке до появления флага,
                # обязано восстанавливаться при переиспользовании ответа.
                require_reuse = field.get("reuse_required", True)
            if source and require_reuse:
                expected.add(name)
        if not expected:
            return True
        context = {
            "payload_id": 0,
            "TODAY": self.today,
            "mp_id": self.mp_id,
            "collected_at": self._collected_at,
        }
        restored: set[str] = set()
        for item in items:
            try:
                row = self._mapper.render_intake(item, context=context)
            except Exception:  # noqa: BLE001
                return False
            if not row:
                continue
            restored |= {name for name in expected if row.get(name) is not None}
            if restored == expected:
                return True
        missing = sorted(expected - restored)
        if missing:
            self.logger.debug(
                "%s: cached payload rejected; attributes not restored: %s",
                self.endpoint_code,
                ", ".join(missing),
            )
        return not missing

    @staticmethod
    def _child_mapper(child_target: dict) -> Mapper:
        """Фабрика дочерних Mapper для nested explode — переопределяется в тестах.

        Выделена отдельным методом, чтобы тесты могли перехватить создание
        Mapper'а (через monkeypatch или замену атрибута экземпляра) и
        подставить _CaptureMapper без правки самой _write_target.
        """
        return Mapper(dict(child_target))

    @staticmethod
    def _write_target(
        mapper: Mapper,
        explode: dict | None,
        raw_payload: Any,
        context: dict,
    ) -> None:
        if not explode:
            mapper.write_intake(raw_payload, context=context)
            return
        parent = raw_payload or {}
        unique_by = explode.get("uniqueBy")
        children = explode.get("children")
        seen_rows: set[tuple] = set()
        index_field = explode.get("indexField")
        for i, element in enumerate(explode_values(parent, explode)):
            if not isinstance(element, dict):
                continue
            merged = dict(element)
            for destination, source in explode["parentFields"].items():
                merged[destination] = parent_val(parent, source)
            if index_field:
                merged[index_field] = i
            row = mapper.render_intake(merged, context=context)
            if row is None:
                continue
            if unique_by:
                key = tuple(row.get(column) for column in unique_by)
                if key in seen_rows:
                    continue
                seen_rows.add(key)
            mapper.write_intake_row(row)
            # nested explode — каждый дочерний таргет даёт строки в своей таблице.
            # Передаём merged с наложенным row: element-поля (apps, date, …)
            # нужны для listField в child explode, а row-поля (stat_date, …)
            # нужны для parentFields дочернего уровня.
            if children:
                merged_with_row = dict(merged)
                merged_with_row.update(row)
                for child_target in children:
                    child_mapper = SpecCollector._child_mapper(child_target)
                    SpecCollector._write_target(
                        child_mapper,
                        child_target.get("explode"),
                        merged_with_row,
                        context,
                    )
