"""앞쪽 **판별** 스텝 (13.68) — 「무엇이 참조 사진을 사야 하는가」.

    13.5  entity_merge        ← 결속할 엔티티가 여기서 정해진다
      ↓
    13.65 grounding_plan      route 확정 (기존 축)
      ↓
    13.68 ★grounding_screen   **판별만** — `era_research.assess_plan_cached`
      ↓
    13.7  entity_filter       판별이 연 것을 **안 지운다** + 없는 행은 만든다
      ↓
    19.x  중앙 획득           ★사는 곳은 **여기 하나**

★이 스텝은 **참조를 사지 않는다.** 검색·이미지 검색·VLM 호출 0.
사는 것은 중앙 19.x 한 곳이다 — 두 곳에서 사면 같은 대상을 두 번 산다.

★`generation_difficulty` 를 **근거로 쓰지 않는다.** 같은 유료 판단을 두 번
사지 않기로 확정됐다 — 판별 SOT 는 `era_research.assess_subjects` 다.

★결속(`binding`)은 `grounding_carry` 가 SOT 다. 이 스텝은 그 장부를 읽을 뿐
이름을 다시 짝짓지 않는다.
"""
from __future__ import annotations

import logging
from typing import Any, Dict

from app.core.grounding_mode import (GROUNDING_MODE_LEGACY,
                                     resolve_grounding_mode)
from app.core.step_runner import StepRunner
from app.core.steps.entity_steps import _EntityStepMixin
from app.modules.pipeline import grounding_carry as _carry
from app.modules.pipeline import grounding_entity_contract as _ec
from app.modules.pipeline import grounding_screen as _screen
from app.modules.pipeline import era_research as _era
from app.modules.pipeline.grounding_subject import SUBJECT_ID_CONTRACT_VERSION

logger = logging.getLogger(__name__)

#: ★판별 캐시가 앉는 칸. CP 안에 두어 **resume 이 다시 안 사게** 한다.
CACHE_KEY = "assess_cache"
#: ★★한 줄이 날 때마다 **여기에 먼저** 쌓인다. 스텝이 중간에 끊겨도 이
#:  파일이 남아 다음 판이 앞의 유료 결과를 **다시 안 산다**.
PARTIAL_NAME = "screen_partial.json"


class GroundingScreenStep(_EntityStepMixin, StepRunner):
    """참조 획득 **의무**를 정한다. 사지는 않는다."""

    def _execute(self, mode: str = "resume") -> Dict[str, Any]:
        g_mode = resolve_grounding_mode(self.project_config)
        if g_mode == GROUNDING_MODE_LEGACY:
            # ★manifest 의 `if_grounding_v2` 가 이미 거른다. 단독 호출로 들어와도
            #  **아무것도 안 하고** 선다.
            logger.info("grounding_screen: mode=legacy — 아무것도 안 한다")
            return self._wrap(g_mode, screened={}, skipped=True)

        # ★★★C(c) 판이면 **새 producer 의 산출**을 읽는다 (Codex BLOCK · 09-01).
        #  앞 판은 어느 판이든 `entity_merge`+`grounding_a0` 를 읽었는데,
        #  C(c) 에서는 그 둘이 **안 돌므로** 판별이 빈손이거나 **옛 CP 를
        #  되쓴다**. 새 producer 는 같은 판독에서 **엔티티 줄과 후보를 함께**
        #  내므로 그것을 그대로 쓴다 — 같은 것을 다시 판단하지 않는다.
        from app.core.grounding_mode import uses_chunk_producer

        if uses_chunk_producer(g_mode):
            chunk_cp = self._load_prev_checkpoint("grounding_chunk")
            data = (chunk_cp or {}).get("data") or {}
            if not data:
                from app.core.errors import AppError
                raise AppError(
                    code="step.no_input",
                    message=("grounding_chunk 산출 없음 — C(c) 판에서 판별은 "
                             "그 producer 의 줄로만 선다"),
                    status_code=400)
            entities = data
            a0_cands = list(data.get("grounding_candidates") or [])
        else:
            merge_cp = self._load_prev_checkpoint("entity_merge")
            entities = (merge_cp or {}).get("data") or {}
            a0_cp = self._load_prev_checkpoint("grounding_a0")
            a0_cands = ((a0_cp or {}).get("data") or {}).get(
                "candidates") or []

        # ★★결속은 `grounding_carry` **한 함수**가 한다 — `grounding_plan` 과
        #  offline 재생기가 쓰는 바로 그 함수다. 여기서 따로 짝지으면 세 곳이
        #  갈리고, 갈린 것을 아무도 못 본다.
        built = _carry.build_subjects(
            entities, project_id=self.project_id, episode_id=self.episode_id,
            source_step="entity_merge", a0_candidates=a0_cands)
        ledger = built.get("candidate_ledger") or {}
        # ★★★**두 장부를 합친다** (Codex BLOCK · 09-01).
        #  후보 장부만 보면 「A0 후보가 없는 기존 엔티티 행」을 표현할 수 없어
        #  `binding_of` 가 그것을 `unbound` 로 떨어뜨린다 — 실측: 등록 LP +
        #  A0=[] → `unbound`. subject 쪽 처분은 `build_subjects` 가 낸다.
        dispositions = dict(built.get("subject_dispositions") or {})
        dispositions.update({
            str(r.get("research_subject_id") or ""): str(r.get("disposition") or "")
            for r in (ledger.get("rows") or [])
        })
        # ★★**다섯 갈래를 전부** 모집단에 넣는다. `build_subjects` 의
        #  `subjects` 는 `DISP_PROMOTED` 만 승격하므로 facet 두 갈래가 통째로
        #  빠진다 — 그것만 보내 놓고 「다섯 owner 판별」이라고 쓰면 거짓이다.
        subjects = _screen.build_population(
            built["subjects"], a0_cands, dispositions)

        from app.core.world_context import build_world_facts_block

        rules_cp = self._load_prev_checkpoint("visual_world_rules")
        world = build_world_facts_block(rules_cp)
        if not world.strip():
            from app.core.errors import AppError

            # ★세계관 없이 판별하면 「그 시대에 특별한가」를 물을 기준이 없다.
            #  빈 블록으로 사면 그 돈이 그냥 버려진다.
            raise AppError(
                code="grounding_screen.missing_world_context",
                message="visual_world_rules 가 비었다 — 시대 기준 없이 판별할 수 없다",
                status_code=400,
            )

        # ★이전 판의 판별 캐시를 이어받는다 — resume 이 같은 것을 **다시 사면**
        #  안 된다. 빈 목록(비대상)도 캐시라 그것까지 이어받아야 뜻이 있다.
        prev = ((self._load_prev_checkpoint(self.step_id) or {}).get("data")
                or {}).get(CACHE_KEY) or {}
        cache: Dict[str, Any] = dict(prev) if isinstance(prev, dict) else {}
        # ★★**끊긴 판이 남긴 것도 이어받는다.** 체크포인트는 스텝이 끝나야
        #  써지므로, N번째에서 끊기면 앞의 N-1 유료 결과가 CP 에는 없다.
        cache.update(self._load_partial())
        failed_memo: set = set()

        def _flush_row(_row, _plan, _cache=cache) -> None:
            # ★캐시 전체를 다시 쓴다. 줄 단위로 덧붙이면 파일이 반쪽인 채로
            #  끊길 수 있고, 반쪽 JSON 은 다음 판이 통째로 못 읽는다.
            self._save_partial(_cache)

        # ★★★C(c) 판이면 **결정적 투영**이다 — 같은 뜻을 다시 안 판단한다.
        #  producer 가 이미 두 축을 냈으므로 여기서 또 사면 ①두 번 사고
        #  ②두 판정이 갈릴 때 정본이 없다 (Codex BLOCK · 09-01).
        if uses_chunk_producer(g_mode):
            projected = _screen.project_from_producer(
                subjects, dispositions=dispositions)
            finals = {str(r.get("research_subject_id") or ""):
                      str(((r.get("producer_link") or {}).get("final_id"))
                          or "")
                      for r in (ledger.get("rows") or [])
                      if (r.get("producer_link") or {}).get("final_id")}
            return self._wrap(g_mode, screened=projected, final_ids=finals)

        out = _screen.screen_subjects(
            subjects,
            dispositions=dispositions,
            world_facts_block=world,
            cache_get=lambda k: (cache.get(k)
                                 if isinstance(cache.get(k), dict) else None),
            cache_put=lambda k, v: cache.__setitem__(k, v),
            project_config=self.project_config,
            opik_metadata=self.build_opik_metadata(),
            failed_memo=failed_memo,
            max_assess=(self.project_config or {}).get(
                "grounding_screen_max_assess"),
            on_row=_flush_row,
        )
        logger.info(
            "grounding_screen: 대상 %d개 · 판별 호출 %d · %s · 갈래별 %s",
            len(subjects), out["assess_calls"], out["counts"], out["by_owner"])

        facet_debt = _screen.facet_obligations(out)
        if facet_debt:
            # ★**사야 하는데 붙을 데가 없는 것들.** §2-6.5a producer 몫이다.
            #  이 목록이 남은 채로 「다섯 갈래 완료」를 말하면 거짓이다.
            logger.warning(
                "grounding_screen: producer 가 없어 미룬 의무 %d건 — §2-6.5a",
                len(facet_debt))
        # ★★**지우지 않는다** (Codex BLOCK-2). manifest 저장은
        #  `StepRunner.run → _execute_and_finalize → save_checkpoint` 에서
        #  `_execute` **가 돌아온 뒤**에 일어난다. 여기서 지우면 반환 직후
        #  crash 나 CP 저장 실패에 산 것이 **아무 데도 안 남는다**.
        #  ★남겨 두는 비용은 파일 하나뿐이다 — 캐시 키에 팩·모델·정책이
        #   접혀 있어 계약이 바뀌면 옛 칸은 그냥 안 맞는다.
        # ★검증된 링크에서만 정본 ID 를 낸다 — 이름으로 다시 짝짓지 않는다
        finals = {str(r.get("research_subject_id") or ""):
                  str(((r.get("producer_link") or {}).get("final_id")) or "")
                  for r in (ledger.get("rows") or [])
                  if (r.get("producer_link") or {}).get("final_id")}
        return self._wrap(g_mode, screened=out, cache=cache, final_ids=finals)

    def _partial_path(self):
        return self._cp_dir / PARTIAL_NAME

    def _load_partial(self) -> Dict[str, Any]:
        import json

        p = self._partial_path()
        if not p.exists():
            return {}
        try:
            got = json.loads(p.read_text(encoding="utf-8"))
        except Exception:  # noqa: BLE001 — 반쪽 파일이 판별을 막으면 안 된다
            logger.warning("grounding_screen: 중간 저장이 깨졌다 — 무시하고 간다")
            return {}
        return got if isinstance(got, dict) else {}

    def _save_partial(self, cache: Dict[str, Any]) -> None:
        import json
        import os

        p = self._partial_path()
        p.parent.mkdir(parents=True, exist_ok=True)
        tmp = p.with_suffix(".tmp")
        tmp.write_text(json.dumps(cache, ensure_ascii=False), encoding="utf-8")
        # ★**원자적으로** 바꾼다. 그냥 덮어쓰면 쓰는 도중에 끊겼을 때
        #  다음 판이 반쪽 파일을 읽는다.
        os.replace(tmp, p)



    def _wrap(self, g_mode: str, *, screened: Dict[str, Any],
              cache: Dict[str, Any] | None = None,
              skipped: bool = False,
              final_ids: Dict[str, str] | None = None) -> Dict[str, Any]:
        rows = screened.get("rows") or []
        counts = screened.get("counts") or {}
        # ★못 선 것과 상한에 걸린 것은 **실패로 센다** — `partial` 이 되어
        #  하류가 막힌다. 「미확정」을 「비대상」으로 읽으면 참조 없이 내려간다.
        failed = (counts.get(_screen.SCREEN_UNRESOLVED, 0)
                  + counts.get(_screen.SCREEN_CAPPED, 0))
        done = len(rows) - failed
        payload: Dict[str, Any] = {
            "mode": g_mode,
            "contract_version": _screen.SCREEN_CONTRACT_VERSION,
            "subject_count": len(rows),
            "search_calls": 0,
            "image_calls": 0,
            **{k: v for k, v in screened.items() if k != "contract_version"},
            "facet_debt": _screen.facet_obligations(screened),
            # ★★subject → **검증된 정본 ID**. 중앙 조사기가 이것을 읽는다.
            #  ★앞 판은 그것이 `grounding_plan` CP 에만 있었는데 C(c) 에서는
            #   그 스텝이 **안 돈다** — 빈 표거나 **옛 CP 의 stale 좌표**가
            #   된다 (Codex 재현 · 09-01).
            "final_id_by_subject": dict(final_ids or {}),
        }
        if cache is not None:
            payload[CACHE_KEY] = cache
        if skipped:
            payload["skipped"] = True
        return {
            "completed_count": done,
            "applicable_count": len(rows),
            "failed_count": failed,
            "config_hash": self._config_hash(),
            "data": payload,
        }

    def _config_hash(self) -> str:
        """★판별 계약·팩 바이트를 지문에 접는다.

        안 접으면 팩이나 모델을 바꿔도 resume 이 **옛 판별을 그대로 건너뛴다**.
        """
        import hashlib
        import json

        from app.core.step_runner import compute_config_hash

        raw = json.dumps({
            "base": compute_config_hash(self.project_config),
            "screen_contract": _screen.SCREEN_CONTRACT_VERSION,
            "carry_contract": _carry.CARRY_CONTRACT_VERSION,
            # ★★★producer 계약도 **접는다** (Codex NON-BLOCK · 09-01).
            #  이 스텝의 정본 행이 `grounding_producer_payload` 를 싣기
            #  시작했으므로, 안 접으면 **payload 가 없던 옛 CP** 를 resume 이
            #  그대로 재사용한다 — 그러면 하류가 판정 칸 없이 돈다.
            "producer_contract": _ec.PRODUCER_CONTRACT_VERSION,
            "subject_contract": SUBJECT_ID_CONTRACT_VERSION,
            "era_policy": _era.ERA_RESEARCH_POLICY_VERSION,
            "era_pack": _era.resolve_era_pack(),
            "era_pack_hash": _era.era_pack_content_hash(),
            "assess_model": _era.ASSESS_MODEL,
            # ★★alias 는 그대로인데 **물리 모델만 바뀌면** outer resume 이
            #  옛 CP 를 건너뛴다 — 다른 모델이 판별한 것을 그대로 쓰게 된다
            #  (Codex). `era_research` 가 캐시 키에 접는 것과 같은 값이다.
            "assess_model_physical": _era.resolve_model_physical(
                _era.ASSESS_MODEL),
        }, sort_keys=True, ensure_ascii=False)
        return hashlib.sha256(raw.encode("utf-8")).hexdigest()[:16]
