"""GROUNDING-V2 §2-3.5 — 분류·계획 스텝.

계획 §1.8 의 **★A 경계**다:

    13.5  entity_merge          short_id 중복 제거 (DB canon 아님)
      ↓
    13.6  ★grounding_plan       검색 0 · short_id 기준 · Sol classifier · route 확정
      ↓
    13.7  entity_filter         보호 목록에 있으면 안 지운다  ← union 을 여기서 받는다

★모드별로 하는 일이 다르다 (계약 §12):

    legacy       **안 돈다** — manifest 의 `if_grounding_v2` 가 걸러 낸다
    shadow_plan  **안 돈다** — legacy 와 구별이 안 된다
    v2           체크포인트를 만든다 — `entity_filter` 가 읽는다

★``shadow_plan`` 을 파이프라인에 들여놓지 않는 이유: 이 스텝들은 production
step id 를 쓰고 DAG 에 하류가 걸려 있어, `force` 로 한 번 돌리면 `_execute`
전에 그 하류가 통째로 무효화된다. shadow 관측은 **저장 CP 만 읽는 offline
재생기**(§2-3a)가 한다.

★``classify_samples`` 를 **실제로 부르는 자리**다.
"""
from __future__ import annotations

import logging
from typing import Any, Dict, List

from app.core.grounding_mode import (
    GROUNDING_MODE_LEGACY, GROUNDING_MODE_V2, resolve_grounding_mode,
)
from app.core.step_runner import StepRunner
from app.modules.pipeline import grounding_carry as _carry
from app.modules.pipeline.grounding_subject import SUBJECT_ID_CONTRACT_VERSION
from app.core.steps.entity_steps import _EntityStepMixin

logger = logging.getLogger(__name__)


#: 표본 수. ★판정기를 결정적으로 만들 수 없어(``_NO_TEMPERATURE_ALIASES``)
#: 한 표본으로 route 를 못 정한다.
DEFAULT_SAMPLES = 3


def check_provenance(runs: List[Any], *, expected_judges=None,
                     expected_samples: int = 0) -> List[str]:
    """표본들이 **같은 것을 잰 것인지** 본다. 어긋나면 사유 목록.

    ★기록만 하면 「같은 것을 N회 쟀다」가 성립하는지 아무도 안 본다.
    측정 도구에만 있던 gate 를 여기로 올린다 (계획 §2-3.5 완료 조건).

    ★판정자가 여럿(패널)이면 **무너졌는지**까지 본다. production 은
    ``strict_single_attempt=False`` 라 tier fallback 으로 gpt 슬롯이 gemini 로
    넘어갈 수 있다. 그러면 실제 alias 는 둘 다 gemini 이고, 「판정자마다
    한 값」만 보는 gate 는 **한 모델 패널을 통과시킨다**(Codex).
    그래서 ①부탁한 집합 ②판정자별 표본 수 ③실제 물리 모델의 **다름**을 본다.
    무너지면 fail-closed — 「미확정」이지 통과가 아니다.
    """
    if not runs:
        return ["표본 지문이 하나도 없다"]
    bad: List[str] = []
    # ★판정자가 **여럿**일 수 있다(패널). 그때 alias/physical 이 다른 것은
    #  정상이다 — 대신 **판정자마다** 나머지가 한 값이어야 한다.
    for col in ("requested_judge_alias", "judge_model_alias",
                "judge_physical_model"):
        vals = {(r or {}).get(col) for r in runs}
        if not vals or None in vals or "" in vals:
            bad.append(f"{col} 이 비었다")

    if expected_judges:
        want = sorted(str(a) for a in expected_judges)
        got = sorted({str((r or {}).get("requested_judge_alias")) for r in runs})
        if got != want:
            bad.append(f"부탁한 판정자 집합이 다르다 — 기대 {want}, 실제 {got}")
        for alias in want:
            grp = [r for r in runs
                   if str((r or {}).get("requested_judge_alias")) == alias]
            if expected_samples and len(grp) != expected_samples:
                bad.append(
                    f"{alias}: 표본이 {len(grp)}개 — {expected_samples}개여야 한다")
        # ★부탁한 판정자마다 **실제로 다른 물리 모델**이 답했어야 한다.
        actual = {alias: {str((r or {}).get("judge_physical_model"))
                          for r in runs
                          if str((r or {}).get("requested_judge_alias")) == alias}
                  for alias in want}
        flat = [m for ms in actual.values() for m in ms]
        if len(want) > 1 and len(set(flat)) < len(want):
            bad.append(
                f"패널이 한 모델로 무너졌다 — 부탁 {want}, 실제 {sorted(set(flat))}")

    by_judge: Dict[str, List[Any]] = {}
    for r in runs:
        by_judge.setdefault(str((r or {}).get("requested_judge_alias")), []).append(r)
    for alias, group in sorted(by_judge.items()):
        for col in ("payload_hash", "prompt_version", "judge_physical_model"):
            vals = {(r or {}).get(col) for r in group}
            if not vals or None in vals or "" in vals:
                bad.append(f"{alias}: {col} 이 비었다")
            elif len(vals) > 1:
                bad.append(
                    f"{alias}: {col} 이 표본마다 다르다 {sorted(map(str, vals))}")
    return bad


#: ★결속은 `grounding_carry` 한 자리에 있다 — 측정 도구와 **같은 함수**를 쓴다.
#: 여기 다시 적으면 둘이 갈리고, 갈리면 잰 것이 뜻을 잃는다.
build_carry_index = _carry.build_carry_index
match_candidate = _carry.match_candidate


def _pack_fingerprint(module: str, version: str) -> Dict[str, str]:
    """팩의 **실효 바이트**를 지문에 넣을 수 있는 모양으로. 못 읽으면 그대로 표시.

    ★버전 문자열만 접으면 **같은 버전 디렉토리의 내용을 고쳤을 때** stale 을
    못 잡는다. 실제로 이 저장소에서 「팩에 스템만 추가했는데 지문이 안
    움직였다」를 두 번 냈다(`test_published_pack_ratchet` 참조).
    """
    from app.core.errors import AppError
    from app.modules.prompt_loader import resolve_effective

    stems = _PACK_FILES.get(module)
    if not stems:
        raise AppError(
            code="grounding.pack_stems_unknown",
            message=f"{module} 의 지문 대상 파일을 안 적어 뒀다",
            status_code=500,
        )
    out: Dict[str, str] = {"module": module, "version": version}
    for stem, kind in sorted(stems):
        try:
            eff = resolve_effective(module, stem, kind=kind, version=version,
                                    db=None)
        except Exception as exc:
            # ★조용히 넘어가면 지문이 **상수**가 되고, 팩을 고쳐도 resume 이
            #  옛 체크포인트를 그대로 건너뛴다 — 고치려던 결함 그 자체다.
            raise AppError(
                code="grounding.pack_unreadable",
                message=(f"{module}/{stem}({kind}) 팩을 못 읽어 지문을 못 만든다: "
                         f"{type(exc).__name__}: {exc}"),
                status_code=500,
            ) from exc
        out[stem] = eff["raw_content_hash"]
    return out


#: 지문에 접을 **프롬프트 파일** (stem, kind) — 팩마다 무엇을 읽는지 못박는다.
#: ★뜻을 가르는 낱말 목록이 아니다. 파일 이름이라 닫혀 있다.
#: ★스키마는 `kind="schema"` 다. `"prompt"` 로 읽으면 못 찾고, 그것을 삼키면
#:  지문이 상수가 된다 (실제로 그럴 뻔했다).
_PACK_FILES = {
    "grounding_a0": (("system", "prompt"), ("a0_schema", "schema")),
    "grounding_classify": (("system", "prompt"), ("classify_schema", "schema")),
    # ★조사 팩도 **바이트**로 접는다 — 버전 문자열만 접으면 같은 버전
    #  디렉토리의 내용을 고쳤을 때 stale 을 못 잡는다.
    "grounding_claims_search": (("system", "prompt"),
                                ("claims_schema", "schema")),
    # ★C(c) — 구간 판독과 합치기 팩 **둘 다**. 하나만 접으면 나머지를 고쳐도
    #  resume 이 옛 체크포인트를 건너뛴다.
    "grounding_chunk": (("system", "prompt"), ("chunk_schema", "schema"),
                        ("merge_system", "prompt"),
                        ("merge_schema", "schema")),
}


def _hash_payload(payload: Dict[str, Any]) -> str:
    import hashlib
    import json as _json

    return hashlib.sha256(
        _json.dumps(payload, sort_keys=True, ensure_ascii=False).encode("utf-8")
    ).hexdigest()[:16]


class GroundingA0Step(_EntityStepMixin, StepRunner):
    """★**거르기 전에** 원문에서 고증 후보를 건진다. 검색을 하지 않는다.

    뒤 단계의 추출 규칙들이(한 번만 나옴 · 착용물 · 고정 설비 · 텍스트만 다른
    종이류) **고증 대상을 정확히 같이 걸러 낸다.** 그래서 그 앞에 선다.

    ★`legacy` 에서는 manifest 의 `if_grounding_v2` 가 not_applicable 로
    걸러 **이 스텝 자체가 안 돈다.**
    """

    def _execute(self, mode: str = "resume") -> Dict[str, Any]:
        from app.modules.pipeline.grounding_a0 import collect_candidates

        g_mode = resolve_grounding_mode(self.project_config)
        if g_mode == GROUNDING_MODE_LEGACY:
            logger.info("grounding_a0: mode=legacy — 아무것도 안 한다")
            return {"completed_count": 0, "applicable_count": 0, "failed_count": 0,
                    "data": {"mode": g_mode, "candidates": [], "skipped": True,
                             "search_calls": 0}}

        fulltext = self._load_cleaned_text()
        rules_cp = self._load_prev_checkpoint("visual_world_rules")
        rules = (rules_cp or {}).get("data") or {}
        out = collect_candidates(
            fulltext,
            project_id=self.project_id, episode_id=self.episode_id,
            era=str(rules.get("era") or "").strip(),
            region=str(rules.get("region") or "").strip(),
            project_config=self.project_config,
            opik_metadata=self.build_opik_metadata(),
        )
        cands = out.get("candidates") or []
        owners: Dict[str, int] = {}  # noqa: E501
        for c in cands:
            owners[c.get("owner_type", "?")] = owners.get(c.get("owner_type", "?"), 0) + 1
        logger.info("grounding_a0: 후보 %d개 %s · 지어낸 인용 %d건",
                    len(cands), owners, len(out.get("hallucinated") or []))
        return {
            "completed_count": len(cands),
            "applicable_count": len(cands),
            "failed_count": len(out.get("hallucinated") or []),
            "config_hash": self._config_hash(),
            "data": {"mode": g_mode, **out, "owner_counts": owners},
        }

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

        안 접으면 팩을 고쳐도 resume 이 **옛 체크포인트를 그대로 건너뛴다** —
        고친 것이 아무 데도 안 닿는다.
        """
        from app.core.step_runner import compute_config_hash
        from app.modules.pipeline import grounding_a0 as _a0

        return _hash_payload({
            "base": compute_config_hash(self.project_config),
            "a0_contract": _a0.A0_CONTRACT_VERSION,
            "subject_contract": SUBJECT_ID_CONTRACT_VERSION,
            "pack": _pack_fingerprint("grounding_a0", _a0.PROMPT_PACK_VERSION),
        })


class GroundingPlanStep(_EntityStepMixin, StepRunner):
    """고증 후보 분류·계획. ★**검색을 하지 않는다.**"""

    def _execute(self, mode: str = "resume") -> Dict[str, Any]:
        from app.modules.pipeline.grounding_classifier import (
            DEFAULT_JUDGES, classify_samples,
        )
        from app.modules.pipeline.grounding_planner import decide_route_from_samples
        from app.modules.pipeline.grounding_overlay import completeness_report
        from app.modules.pipeline.grounding_subject import build_subject

        g_mode = resolve_grounding_mode(self.project_config)
        if g_mode == GROUNDING_MODE_LEGACY:
            # ★여기까지 오면 안 된다 — manifest 의 `if_grounding_v2` 가
            #  legacy 를 not_applicable 로 걸러 스텝 자체가 안 돈다. 그래도
            #  단독 호출로 들어올 수 있으니 **아무것도 안 하고** 선다.
            logger.info("grounding_plan: mode=legacy — 아무것도 안 한다")
            return self._wrap(g_mode, decided=[], subject_count=0, skipped=True,
                              config_hash=self._config_hash())

        merge_cp = self._load_prev_checkpoint("entity_merge")
        entities = (merge_cp or {}).get("data") or {}
        rules_cp = self._load_prev_checkpoint("visual_world_rules")
        rules = (rules_cp or {}).get("data") or {}
        era = str(rules.get("era") or "").strip()
        region = str(rules.get("region") or "").strip()
        if not era or not region:
            from app.core.errors import AppError

            # ★빈 시대로 판정하면 분류기가 상상 묘사만 보고 답한다 — 실측으로 확인된 결함.
            raise AppError(
                code="grounding_plan.missing_world_context",
                message=(f"visual_world_rules 에 시대/지역이 없다 "
                         f"(era={era!r}, region={region!r})"),
                status_code=400,
            )

        # ★A0 후보를 **물려받는다.** 여기서 subject 를 새로 발급하면 같은 대상에
        #  다른 id 가 붙고, 분류기가 보는 근거도 **원문 인용이 아니라 LLM 이
        #  상상해 쓴 description** 이 된다 — 판정의 전제가 무너진다.
        #  ★결속 자체는 `grounding_carry` 가 한다 — 측정 도구와 같은 함수다.
        a0_cp = self._load_prev_checkpoint("grounding_a0")
        a0_cands = ((a0_cp or {}).get("data") or {}).get("candidates") or []
        built = _carry.build_subjects(
            entities, project_id=self.project_id, episode_id=self.episode_id,
            source_step="entity_merge", a0_candidates=a0_cands)
        subjects = built["subjects"]
        unbound = built["unbound"]
        carry_reasons = built["carry_reasons"]
        carried = built["carried"]
        # ★★**장부를 체크포인트에 남긴다.** 문서가 「SOT 는 disposition 산출과
        #  체크포인트」라고 적어 놓고 체크포인트에 없으면 그 문장이 거짓이다
        #  (Codex).
        #  ★정정: 지금 이 칸을 **읽는 코드는 없다.** offline 재생기
        #  (`grounding_shadow`)는 `entity_merge`+A0 에서 `build_subjects` 를
        #  다시 불러 **같은 함수로 재구축**한다. 이 칸은 「그때 무엇이 어디로
        #  갔나」를 나중에 사람이 되짚기 위한 기록이다.
        candidate_ledger = built.get("candidate_ledger") or {}

        logger.info("grounding_plan: A0 후보 %d개 중 %d개를 물려받았다 %s",
                    len(a0_cands), carried, carry_reasons)
        _unbound = (carry_reasons.get("ambiguous", 0)
                    + carry_reasons.get("duplicate_surface", 0))
        if _unbound:
            # ★못 정한 것을 「없었다」로 읽지 않는다 — 기록에 남긴다.
            logger.warning("grounding_plan: 표면형이 모호해 결속 안 한 것 %d건",
                           _unbound)

        # ★승격된 후보는 **엔티티 목록에 없는 것이 정상**이다 — 없어서 승격한
        #  것이라 「엔티티에 있나」로 물으면 당연히 없다. 그것을 `missing` 으로
        #  세면 스텝이 `partial` 이 되어 하류가 막힌다.
        _promoted = [r["research_subject_id"]
                     for r in (candidate_ledger.get("rows") or ())
                     if r.get("disposition") == "promoted"]
        report = completeness_report(a0_cands, {
            "character": entities.get("characters") or [],
            "location": entities.get("locations") or [],
            "prop": entities.get("props") or [],
        }, promoted_ids=_promoted)
        if not report["complete"]:
            # ★여기서 예외를 던지지는 않는다 — 세우면 **얼마나 사라졌는지도
            #  못 본다.** 대신 `failed_count` 로 올려 스텝을 `partial` 로 만들고,
            #  manifest 의 `allow_partial_downstream=False` 가 `entity_filter`
            #  를 막는다. 기록은 남기고 진행은 멈춘다.
            logger.warning("grounding_plan: A0 후보 %d개가 산출에 안 남았다 %s",
                           report["missing_count"], sorted(report["missing"]))
        # ★★**승격으로 살렸다고 조용히 넘어가지 않는다.** 엔티티 행이 사라진
        #  것은 그대로 사실이고(추출이 지웠다), 그게 안 보이면 「승격이 있으니
        #  괜찮다」로 추출 결함이 영원히 안 드러난다. 이 수를 아무도 안 읽으면
        #  칸만 있고 뜻이 없다.
        _saved = int(report.get("entity_missing_count") or 0) \
            - int(report.get("missing_count") or 0)
        if _saved > 0:
            logger.warning(
                "grounding_plan: 엔티티 행이 없어진 A0 후보 %d개를 **승격으로 "
                "살렸다** — 추출이 지운 것이다. 조사는 계속되지만 그 사실은 "
                "남는다 (entity_missing=%d · missing=%d · promoted=%d)",
                _saved, report.get("entity_missing_count"),
                report.get("missing_count"), report.get("promoted_count"))

        if not subjects:
            # ★못 붙인 것만 있으면 **분류기를 안 부른다** — 살 이유가 없다.
            #  그래도 그 행들은 판정을 낸 행이라 집계에 들어간다.
            only: Dict[str, int] = {}
            for d in unbound:
                only[d["route"]] = only.get(d["route"], 0) + 1
            # ★★**여기서도 장부를 넘긴다.** 이 갈래는 「전부 얽혀서 아무것도
            #  분류 못 한 판」이라 장부가 **제일 필요한 자리**인데, 조기 반환이
            #  칸을 안 넘겨 그 판에서만 장부가 통째로 비었다 (Codex).
            #  ★반환이 둘이면 한쪽만 고쳐진다 — 인자를 더할 때 **두 자리를
            #  다 본다**.
            return self._wrap(g_mode, decided=list(unbound), subject_count=0,
                              counts=only,
                              completeness=report, a0_carried=carried,
                              carry_reasons=carry_reasons,
                              candidate_ledger=candidate_ledger,
                              config_hash=self._config_hash())

        smp = classify_samples(
            subjects, samples=DEFAULT_SAMPLES, era=era, region=region,
            project_config=self.project_config,
            opik_metadata=self.build_opik_metadata(),
        )
        # ★기대 집합은 **계약**(`DEFAULT_JUDGES`)에서 받는다. 호출이 돌려준
        #  `smp["judges"]` 에서 받으면 **자기보고**다 — 한 판정자만 돌고 그
        #  목록도 한 판정자로 줄면 그대로 통과한다(Codex).
        runs = smp.get("runs") or []
        reported = list(smp.get("judges") or [])
        bad = check_provenance(
            runs,
            expected_judges=list(DEFAULT_JUDGES),
            expected_samples=DEFAULT_SAMPLES)
        if reported != list(DEFAULT_JUDGES):
            bad.append(
                f"분류기가 돌린 판정자가 계약과 다르다 — "
                f"계약 {list(DEFAULT_JUDGES)}, 실제 {reported}")
        if bad:
            from app.core.errors import AppError

            raise AppError(
                code="grounding_plan.provenance_mismatch",
                message="같은 것을 N회 잰 것이 아니다: " + "; ".join(bad),
                status_code=502,
            )

        decided: List[Dict[str, Any]] = []
        for s in subjects:
            recs = smp["by_subject"].get(s["research_subject_id"], [])
            d = decide_route_from_samples(recs)
            decided.append({
                **d,
                "research_subject_id": s["research_subject_id"],
                "_surface_form": s["surface_form"],
                "_owner_type": s["owner_type"],
                "_short_id": (s.get("provenance") or {}).get("short_id"),
            })

        # ★집계 모집단은 **최종 decided** 다. 분류기가 낸 것만 세면
        #  `unresolved`(못 붙여 분류 안 한 것)가 빠져 `sum(counts) != len(decided)`
        #  가 되고, 전부 못 붙인 판에서는 `counts` 가 통째로 빈다.
        final = decided + unbound
        counts: Dict[str, int] = {}
        for d in final:
            counts[d["route"]] = counts.get(d["route"], 0) + 1
        unstable = sum(1 for d in final if d.get("unstable"))
        logger.info("grounding_plan: mode=%s · %d후보 · %s · 흔들림 %d",
                    g_mode, len(decided), counts, unstable)

        return self._wrap(
            g_mode, decided=final, subject_count=len(subjects),
            counts=counts, unstable_count=unstable,
            logical_calls=smp.get("logical_calls", 0),
            runs=smp.get("runs") or [],
            completeness=report, a0_carried=carried,
            carry_reasons=carry_reasons, candidate_ledger=candidate_ledger,
            config_hash=self._config_hash(),
        )

    @staticmethod
    def _wrap(g_mode: str, *, decided: List[Dict[str, Any]], subject_count: int,
              counts: Dict[str, int] | None = None, unstable_count: int = 0,
              logical_calls: int = 0, runs: List[Any] | None = None,
              skipped: bool = False, completeness: Dict[str, Any] | None = None,
              a0_carried: int = 0, config_hash: str = "",
              carry_reasons: Dict[str, int] | None = None,
              candidate_ledger: Dict[str, Any] | None = None
              ) -> Dict[str, Any]:
        """★체크포인트 모양을 **다른 스텝과 같게** 만든다.

        `run()` 은 ``save_checkpoint({"status": ..., **result})`` 로 **평평하게**
        펼친다. 그래서 하류가 읽는 자리는 ``manifest["data"]`` 다 —
        평평하게 돌려주면 `entity_filter` 의 ``data.decided`` 가 **늘 비어**
        등록을 해도 보호가 하나도 안 걸린다. 실제로 그랬다.
        """
        reasons = carry_reasons or {}
        # ★못 정한 결속도 **실패로 센다** — 세지 않으면 A0 근거를 잃은 채
        #  completed 로 지나간다. `failed>0` 이면 `run()` 이 `partial` 로
        #  확정하고 하류가 막힌다.
        unresolved_carry = sum(reasons.get(r, 0)
                               for r in _carry.UNBOUND_REASONS)
        missing = (int((completeness or {}).get("missing_count") or 0)
                   + unresolved_carry)
        # ★못 붙인 것도 **판정을 낸 행**이다(unresolved). 완료 수에 세지 않으면
        #  `completed==0` 이 되어 status 가 `partial` 이 아니라 `failed` 로 간다 —
        #  스텝이 깨진 것이 아니라 「못 붙였다」를 제대로 낸 것이다.
        done = subject_count + len(decided) - len([
            d for d in decided if d.get("research_subject_id")])

        payload: Dict[str, Any] = {
            "mode": g_mode,
            "subject_count": subject_count,
            "counts": counts or {},
            "unstable_count": unstable_count,
            "logical_calls": logical_calls,
            "search_calls": 0,
            "runs": runs or [],
            "decided": [],
            "a0_carried": a0_carried,
            "carry_reasons": reasons,
            "unresolved_carry_count": unresolved_carry,
            "completeness": completeness or {},
            # ★★후보 **하나하나**의 행선지. 문서가 「SOT 는 체크포인트」라고
            #  적어 놓고 여기 없으면 그 문장이 거짓이다 (Codex).
            "candidate_ledger": candidate_ledger or {},
        }
        if skipped:
            payload["skipped"] = True
        # ★`v2` 에서만 쓴다. `shadow_plan` 은 이 스텝이 **안 돈다** —
        #  manifest 의 `if_grounding_v2` 가 not_applicable 로 거른다.
        #  shadow 관측은 저장 CP 만 읽는 offline 재생기(§2-3a)가 한다.
        if g_mode == GROUNDING_MODE_V2:
            payload["decided"] = decided
        return {
            "completed_count": done,
            "applicable_count": done,
            # ★사라진 후보 수 = 실패 수. `failed>0 and completed>0` 이면
            #  `run()` 이 `partial` 로 확정하고 하류가 막힌다.
            "failed_count": missing,
            "config_hash": config_hash,
            "data": payload,
        }

    def _config_hash(self) -> str:
        """★팩·계약·**표본 수**·**판정자**·대상 모델을 접는다.

        표본 수를 바꾸면 판정이 달라지는데 지문이 안 움직이면 옛 판을 그대로
        쓴다. ★**판정자 목록**도 마찬가지다 — `gpt+gemini-pro` 를
        `gpt+grok` 으로 바꿔도 지문이 안 움직이면 resume 이 **옛 판정을 그대로
        건너뛴다**. 누가 답했는지가 바뀌었는데 다시 안 묻는 것이다 (Codex).
        """
        from app.core.step_runner import compute_config_hash
        from app.modules.pipeline import grounding_classifier as _gc
        from app.modules.pipeline.grounding_overlay import OVERLAY_CONTRACT_VERSION
        from app.modules.pipeline.grounding_planner import PLANNER_CONTRACT_VERSION

        return _hash_payload({
            "base": compute_config_hash(self.project_config),
            "planner_contract": PLANNER_CONTRACT_VERSION,
            "carry_contract": _carry.CARRY_CONTRACT_VERSION,
            "overlay_contract": OVERLAY_CONTRACT_VERSION,
            "subject_contract": SUBJECT_ID_CONTRACT_VERSION,
            "samples": DEFAULT_SAMPLES,
            # ★**집합**으로 접는다 — 순서는 신원이 아니다.
            #  `classify_samples` 는 판정자를 순서대로 다 부르고 결과를 합칠
            #  뿐이라 첫 판정자에게 우선권이 없고, `decide_route_from_samples`
            #  도 순서에 안 흔들린다(실측). 그런데 지문만 순서를 세면
            #  **같은 판정이 나오는데 지문이 달라 다시 산다.**
            "judges": sorted(str(j) for j in _gc.DEFAULT_JUDGES),
            "classifier_contract": _gc.CLASSIFIER_CONTRACT_VERSION,
            # ★2026-08-30: 이미지 모델 좌표를 **뺐다.** 분류가 그 좌표를 안
            #  보는데 지문에 남기면, 이미지 backend 만 바꿔도 같은 분류를
            #  전부 다시 산다.
            "pack": _pack_fingerprint("grounding_classify",
                                      _gc.PROMPT_PACK_VERSION),
        })

class GroundingResearchStep(_EntityStepMixin, StepRunner):
    """★**출처 붙은 조사를 실제로 산다.** GROUNDING-V2 §2-4b.

    `grounding_plan` 이 `route="research"` 로 표시한 것만 조사한다 — `skip` 은
    안 산다. 그게 계획의 뜻이다.

    ## 시간 통제 세 겹 안에서 돈다 (§8.5)

    `research_run_scope(cap=, deadline_seconds=)` 로 열고, 그 안에서
    `search_claims` 가 호출마다 마감을 건다. 셋 중 어디에 걸리든 판정은 같다 —
    `time_capped` · `unresolved` · retryable. 무엇에 걸렸는지는 `limit_kind` 다.

    ## ★**main thread 만** 결과를 남긴다

    `call_with_deadline` 은 **취소가 아니라 포기**라, 버려진 worker 가 나중에
    값을 들고 온다. 그 값은 `search_claims` 가 이미 버렸고 여기서는 **돌아온
    행만** 저장한다 — worker 는 바깥 호출과 감사 기록만 한다.

    ## batch 크기는 **1** 이다

    §4b 실험이 네 크기를 다 돌고 「후보 없음 · 개별 호출 유지」로 끝났다.
    어긋남은 어느 크기에서도 0이었지만 **크기를 키울수록 조사량이 말랐다**
    (계약 통과 claim 19 → 16 → 10 → 9). 그래서 기본값을 안 바꾼다.
    """

    #: ★한 주행의 **물리 전송** 상한. 대상 수와 다르다 — 슬롯 failover 가
    #:  논리 하나를 물리 둘로 만든다.
    DEFAULT_TRANSMISSION_CAP = 60
    #: 주행 전체 벽시계. ★넘으면 남은 대상은 **`unresolved`** 다 — `skip` 이
    #:  아니다(계약 §3: 「모른다를 아니다로 닫지 않는다」).
    DEFAULT_RUN_DEADLINE_SECONDS = 1800.0
    #: ★§4b 실험이 정한 값. 바꾸려면 그 실험을 다시 해야 한다.
    BATCH_SIZE = 1

    def _execute(self, mode: str = "resume") -> Dict[str, Any]:
        from app.core.openai_keys import openai_client
        from app.core.research_call_budget import research_run_scope
        from app.modules.pipeline import grounding_claims_search as gcs

        g_mode = resolve_grounding_mode(self.project_config)
        if g_mode == GROUNDING_MODE_LEGACY:
            logger.info("grounding_research: mode=legacy — 아무것도 안 한다")
            return self._wrap_research(g_mode, rows=[], skipped=True)

        subjects, era, region, no_source, missing = self._research_targets()
        if no_source or missing:
            # ★유료 호출 **전에** 말한다 — 사고 나서 알면 늦다.
            logger.warning("grounding_research: 인용 없는 대상 %d · 못 세운 것 %d",
                           len(no_source), len(missing))
        if not subjects:
            # ★인용이 없어서 못 사는 것은 「대상이 없다」가 **아니다**.
            logger.info("grounding_research: 살 대상이 없다 (인용 없음 %d)",
                        len(no_source))
            return self._wrap_research(
                g_mode, rows=[], era=era, region=region,
                no_source=no_source, missing=missing,
                config_hash=self._config_hash())

        # ★★★**이미 끝난 것은 다시 안 산다** (§2-4a ②). 배선을 빼면 재개할
        #  때마다 같은 입력을 통째로 다시 산다 — 그게 이 프로젝트에서 제일 비싼
        #  결함 부류다.
        #  ★상한을 넘은 것은 **버리지 않고** `time_capped`(admission_limit)로
        #   남는다 — 「조사할 필요가 없었다」가 아니다.
        admission = self._plan_admission(subjects, era=era, region=region)
        # ★★이번 판이 기대하는 **입력 신원**. 정책 스텝이 이것과 맞는 revision
        #  만 본다 — 안 그러면 입력이 바뀐 옛 결과가 섞인다 (Codex).
        expected = dict(admission.get("expected_input_hashes") or {})
        keep = set(admission["admitted"])
        deferred = list(admission["capped"])
        already = list(admission["already_done"])
        subjects = [s for s in subjects if s["research_subject_id"] in keep]
        if already or deferred:
            logger.info("grounding_research: 이미 끝난 것 %d · 상한에 밀린 것 %d",
                        len(already), len(deferred))
        if not subjects:
            logger.info("grounding_research: 새로 살 것이 없다")
            return self._wrap_research(
                g_mode, rows=[], era=era, region=region,
                no_source=no_source, missing=missing,
                already_done=already, deferred=deferred,
                expected=expected, config_hash=self._config_hash())

        cap = int(self.project_config.get("research_transmission_cap")
                  or self.DEFAULT_TRANSMISSION_CAP)
        deadline = float(self.project_config.get("research_run_deadline_seconds")
                         or self.DEFAULT_RUN_DEADLINE_SECONDS)
        logger.info("grounding_research: 대상 %d개 · 전송 상한 %d · 마감 %.0f초",
                    len(subjects), cap, deadline)

        # ★★상한과 마감을 **먼저** 연다. 열기 전에 부르면 그 호출은 아무 문도
        #  안 지난다 — 그게 이 스텝이 있는 이유의 절반이다.
        saved = 0

        def _persist(row):
            """★★**산 것을 한 호출마다 바로 남긴다.** 9개 중 7개째에 끊기면
            앞의 7개를 통째로 잃고 다시 사야 한다 — 실제로 그 부류의 결함을
            이 프로젝트에서 냈다.

            ★이 콜백은 `search_claims` 가 **main thread 에서** 부른다. 버려진
            worker 는 여기 못 온다.
            """
            nonlocal saved
            saved += self._save_revisions([row], subjects, era=era,
                                          region=region)

        with research_run_scope(cap=cap, deadline_seconds=deadline) as budget:
            out = gcs.search_claims(
                openai_client(max_retries=gcs.SDK_RETRIES), subjects,
                model=self._search_model(), era=era, region=region,
                batch_size=self.BATCH_SIZE, db=self.db,
                opik_metadata=self.build_opik_metadata(),
                on_batch=_persist)
        rows = out.get("batches") or []
        spent = budget.snapshot()
        capped = [r for r in rows if r.get("limit_kind")]
        logger.info("grounding_research: 논리 %d회 · 물리 %s · 상한에 걸린 것 %d "
                    "· 저장 %d행",
                    out.get("logical_calls", 0), spent, len(capped), saved)
        return self._wrap_research(
            g_mode, rows=rows, era=era, region=region,
            budget=spent, logical_calls=out.get("logical_calls", 0),
            bought_calls=out.get("bought_calls", 0),
            no_source=no_source, missing=missing, saved=saved,
            already_done=already, deferred=deferred,
            expected=expected, config_hash=self._config_hash())

    def _plan_admission(self, subjects, *, era: str, region: str):
        """무엇을 새로 살지. ★**바깥 호출 전에** 정한다.

        ★이미 끝난 행(`completed`)은 다시 안 산다. 상한을 넘은 것은 버리지 않고
        `capped` 로 남아 다음 주행이 **앞에서부터 이어간다** — 상한만 두고
        이어가는 계약이 없으면 영구 partial 이 된다.
        """
        from app.modules.pipeline.grounding_claims import (
            plan_admission, subject_payload_hash)

        ids = [s["research_subject_id"] for s in subjects]
        return plan_admission(
            ids,
            subject_payload_hashes={s["research_subject_id"]:
                                    subject_payload_hash(s) for s in subjects},
            research_inputs=self._research_inputs(era=era, region=region),
            done_rows=self._done_rows(),
            cap=int(self.project_config.get("research_subject_cap") or 0) or None,
            batch_size=self.BATCH_SIZE)

    def _pack(self):
        """★★★**활성 DB 팩**을 한 번 읽어 세 자리가 같이 쓴다 (Codex).

        `search_claims(db=self.db)` 는 DB 에 올라간 prompt/schema 를 쓰는데,
        신원과 지문을 `db=None` 으로 읽으면 **셋이 갈린다**:
        ①admission 이 낸 hash 와 저장된 hash 가 달라 **재개마다 다시 사고**
        ②config 지문이 DB 변경을 못 봐 **완료 CP 를 건너뛴다**.
        """
        from app.modules.pipeline import grounding_claims_search as gcs

        cached = getattr(self, "_pack_cache", None)
        if cached is None:
            cached = gcs.load_pack(db=self.db)
            self._pack_cache = cached
        return cached

    def _research_inputs(self, *, era: str, region: str, prov=None):
        """입력 신원 중 **대상과 무관한** 좌표. ★한 곳에서 만든다.

        ★★`prov` 를 주면 그 호출이 **실제로 쓴** 좌표를 쓴다. 입력 신원은
        「무엇을 요청했나」이고 provenance 는 「무엇이 답했나」라, 계약이 그
        일곱 쌍의 **일치**를 검사한다 — 저장할 때 살 때의 값을 그대로 쓰면
        팩이 바뀐 판에서 둘이 갈린다.
        ★프로덕션에서는 같은 주행이 같은 팩으로 사고 저장하므로 값이 같다.
        """
        from app.modules.pipeline import grounding_claims_search as gcs
        from app.modules.pipeline.grounding_planner import (
            PLANNER_CONTRACT_VERSION)

        pv = dict(prov or {})
        pack = self._pack() if not pv else None
        return {
            "era": era, "region": region,
            "claims_pack_version": (pv.get("prompt_pack_version")
                                    if pv else pack["version"]),
            "prompt_raw_hash": (pv.get("prompt_raw_hash") if pv else
                                pack["stems"][gcs.SYSTEM_STEM][
                                    "raw_content_hash"]),
            "schema_raw_hash": (pv.get("schema_raw_hash") if pv else
                                pack["stems"][gcs.SCHEMA_STEM][
                                    "raw_content_hash"]),
            "policy_version": str(PLANNER_CONTRACT_VERSION),
            "policy_contract_hash": _hash_payload(
                {"planner": PLANNER_CONTRACT_VERSION}),
            "provider": pv.get("provider") or "openai",
            "model": pv.get("model") or self._search_model(),
            "model_version": pv.get("model") or self._search_model(),
        }

    def _done_rows(self):
        """이 에피소드의 **저장된 조사**. ★없으면 전부 새로 사게 된다."""
        from app.models.project import GroundingResearchRevision

        if self.db is None:
            return []
        rows = (self.db.query(GroundingResearchRevision)
                .filter(GroundingResearchRevision.project_id == self.project_id,
                        GroundingResearchRevision.episode_id == self.episode_id)
                .all())
        return [{"research_subject_id": r.research_subject_id,
                 "research_input_hash": r.research_input_hash,
                 "status": r.status} for r in rows]

    def _save_revisions(self, rows, subjects, *, era: str, region: str) -> int:
        """산 것을 **정본 테이블**에 남긴다. ★판정은 `build_revision_row` 가 한다.

        ★호출부가 `decide_delta` 를 따로 부르고 결과만 넘기면 저장된 판정과
        저장된 claims 가 갈릴 수 있다 — 그래서 한 곳에서 한다.
        ★같은 시도를 두 번 안 넣는다(유일 제약). 재개하면 `attempt_id` 가
        달라져 새 행이 되고, 앞 행은 그대로 남는다.
        """
        import uuid
        from datetime import datetime, timezone

        from app.models.project import GroundingResearchRevision
        from app.modules.pipeline.grounding_claims import (
            build_revision_row, subject_payload_hash)
        by_id = {s["research_subject_id"]: s for s in subjects}
        now = datetime.now(timezone.utc).isoformat()
        n = 0
        for row in rows:
            prov = dict(row.get("provenance") or {})
            if not prov:
                # ★★안 보낸 행에도 provenance 가 **있다** — `search_claims` 가
                #  팩 좌표와 그 호출의 신원을 채워 준다. 여기까지 비어 오면
                #  그건 다른 결함이니 **조용히 지어내지 않고** 건너뛴다.
                logger.warning("grounding_research: provenance 가 없는 행 — "
                               "저장하지 않는다 %s", row.get("requested"))
                continue
            parsed = (row.get("parsed") or {}).get("results") or []
            for sid in (row.get("requested") or []):
                subj = by_id.get(sid)
                if not subj:
                    continue
                mine = [r for r in parsed if isinstance(r, dict)
                        and str(r.get("research_subject_id") or "") == sid]
                claims = (mine[0].get("claims") or []) if mine else []
                gaps = (mine[0].get("gaps") or []) if mine else []
                # ★`research_subject_id` 와 `subject_payload_hash` 는 **따로**
                #  넘긴다 — `canonical_research_input` 이 그 둘을 직접 받으므로
                #  여기 또 넣으면 같은 인자를 두 번 주는 것이 된다.
                sph = subject_payload_hash(subj)
                # ★★admission 이 쓴 것과 **같은 함수**에서 받는다. 여기서 따로
                #  조립하면 「살지 말지 정할 때 본 신원」과 「저장된 신원」이
                #  갈리고, 갈리면 다음 주행이 같은 것을 다시 산다.
                ri = self._research_inputs(era=era, region=region,
                                           prov=prov)
                built = build_revision_row(
                    project_id=self.project_id, episode_id=self.episode_id,
                    research_subject_id=sid,
                    subject_payload_hash=sph,
                    research_inputs=ri, claims=claims, gaps=gaps,
                    # ★`research_input_hash` 는 `build_revision_row` 가
                    #  `research_inputs` 에서 낸다 — provenance 에 또 넣으면
                    #  「모르는 칸」으로 거부된다(계약이 그렇게 막는다).
                    # ★provenance 는 **입력 신원에 결속**된다 — 겹치는 좌표를
                    #  손으로 맞추면 갈리고, 갈리면 그 행이 정말 그 입력의
                    #  것인지 못 밝힌다. 계약이 그 결속을 검사한다.
                    provenance={**prov,
                                "policy_version": ri["policy_version"],
                                "model_version": ri["model_version"]},
                    attempt_id=str(uuid.uuid4()),
                    # ★상한에 걸린 것은 **시간에 잘린 것**이다 — 「다 보고 못
                    #  정했다」와 갈라야 다음 판에 다시 산다.
                    time_capped=bool(row.get("limit_kind")),
                    created_at=now)
                if self.db is not None:
                    self.db.add(GroundingResearchRevision(**built))
                n += 1
        if self.db is not None and n:
            self.db.flush()
        return n

    # ── 재료 ──────────────────────────────────────────────────────────
    def _research_targets(self):
        """계획이 `research` 로 표시한 대상 + 그 **원문 인용**.

        ★subject 는 `grounding_plan` 과 **같은 함수**로 다시 세운다
        (`_carry.build_subjects`). 여기서 새로 발급하면 같은 대상에 다른 id 가
        붙어 계획과 조사가 갈린다.
        """
        plan = (self._load_prev_checkpoint("grounding_plan") or {}
                ).get("data") or {}
        want = {d.get("research_subject_id") for d in (plan.get("decided") or [])
                if d.get("route") == "research"}
        rules = (self._load_prev_checkpoint("visual_world_rules") or {}
                 ).get("data") or {}
        era = str(rules.get("era") or "").strip()
        region = str(rules.get("region") or "").strip()
        if not era or not region:
            from app.core.errors import AppError

            # ★빈 시대로 조사하면 무엇과 비교할지가 없다.
            raise AppError(
                code="grounding_research.missing_world_context",
                message=f"시대/지역이 없다 (era={era!r}, region={region!r})",
                status_code=400)
        if not want:
            # ★★모양을 **다섯으로 맞춘다.** 셋만 돌려주면 호출부가 풀다가
            #  죽는다 — 「조사할 대상이 하나도 없다」는 **정상 판**인데
            #  그 판이 통째로 예외가 됐다 (Codex).
            return [], era, region, [], []

        entities = (self._load_prev_checkpoint("entity_merge") or {}
                    ).get("data") or {}
        a0 = ((self._load_prev_checkpoint("grounding_a0") or {}).get("data")
              or {}).get("candidates") or []
        built = _carry.build_subjects(
            entities, project_id=self.project_id, episode_id=self.episode_id,
            source_step="entity_merge", a0_candidates=a0)
        picked = [s for s in built["subjects"]
                  if s["research_subject_id"] in want]
        # ★★원문 인용이 없는 대상은 **안 산다** — 상상 묘사로 조사하면 그
        #  결과가 무엇의 것인지 알 수 없다(§2-4a 계약).
        #  ★★★그렇다고 **조용히 버리지 않는다** (Codex). 그건 「조사할 대상이
        #   없다」가 아니라 `no_source`·`unresolved` 다. 버리면 유료 호출 0회로
        #   `completed` 가 되어 하류가 「조사할 게 없었다」로 읽는다.
        quoted = [s for s in picked
                  if str(s.get("source_quote") or "").strip()]
        no_source = [s["research_subject_id"] for s in picked
                     if not str(s.get("source_quote") or "").strip()]
        # ★계획이 research 로 표시했는데 subject 를 아예 못 세운 것도 남긴다 —
        #  「전수」가 아니면 어디로 갔는지 못 찾는다.
        missing = sorted(want - {s["research_subject_id"] for s in picked})
        return quoted, era, region, no_source, missing

    def _search_model(self) -> str:
        from app.core.config import settings

        return str(self.project_config.get("grounding_search_model")
                   or settings.openai_model)

    def _config_hash(self) -> str:
        """★팩 바이트·계약 버전에 더해 **요청 계약**까지 접는다.

        안 접으면 `strict` 나 `include` 를 고쳐도 resume 이 **옛 체크포인트를
        그대로 건너뛴다** — 고친 것이 아무 데도 안 닿는다.
        """
        import json as _json

        from app.core.step_runner import compute_config_hash
        from app.modules.pipeline import grounding_claims_search as gcs

        # ★★지문도 **DB 실효 팩**에서 낸다. `db=None` 으로 읽으면 DB 에 올라간
        #  프롬프트를 고쳐도 지문이 안 움직여 **완료 CP 를 그대로 건너뛴다**.
        pack = self._pack()
        return _hash_payload({
            "base": compute_config_hash(self.project_config),
            "pack": {
                "module": pack["module"], "version": pack["version"],
                "manifest": pack["pack_manifest_hash"],
                # ★스템 **바이트**까지 — 같은 버전 안에서 내용만 바뀐 판을 잡는다
                "stems": {k: v.get("raw_content_hash", "")
                          for k, v in sorted(pack["stems"].items())},
            },
            "request": _json.dumps(gcs.request_contract(), sort_keys=True,
                                   ensure_ascii=False),
            "batch_size": self.BATCH_SIZE,
        })

    @staticmethod
    def _wrap_research(g_mode: str, *, rows, era: str = "", region: str = "",
                       budget=None, logical_calls: int = 0,
                       bought_calls: int = 0, skipped: bool = False,
                       no_source=None, missing=None, saved: int = 0,
                       already_done=None, deferred=None, expected=None,
                       config_hash: str = "") -> Dict[str, Any]:
        """★체크포인트 모양을 다른 스텝과 같게. 하류가 읽는 자리는 `data` 다.

        ★**상한에 걸린 것을 실패로 센다.** 세지 않으면 「시간이 없어 못 산 것」이
        `completed` 로 지나가고, 하류가 그걸 「조사할 게 없었다」로 읽는다.
        """
        capped = [r for r in (rows or []) if r.get("limit_kind")]
        errored = [r for r in (rows or [])
                   if r.get("error") and not r.get("limit_kind")]
        ns = list(no_source or ())
        ms = list(missing or ())
        return {
            "completed_count": len(rows or []) - len(capped) - len(errored),
            # ★못 산 것도 **모집단에 넣는다** — 빼면 「전수」가 아니다.
            "applicable_count": len(rows or []) + len(ns) + len(ms),
            # ★인용 없음·못 세움도 **실패로 센다**. 안 세면 유료 호출 0회로
            #  `completed` 가 되어 하류가 「조사할 게 없었다」로 읽는다.
            "failed_count": (len(capped) + len(errored) + len(ns) + len(ms)
                             + len(list(deferred or ()))),
            "config_hash": config_hash,
            "data": {
                "mode": g_mode, "skipped": skipped,
                "era": era, "region": region,
                "batch_size": GroundingResearchStep.BATCH_SIZE,
                "rows": rows or [],
                "budget": budget or {},
                "logical_calls": logical_calls,
                "bought_calls": bought_calls,
                # ★무엇에 걸렸는지 갈래별로 남긴다 — 고칠 곳이 다르다
                "limit_kinds": sorted({str(r.get("limit_kind"))
                                       for r in capped}),
                "capped_count": len(capped),
                # ★어디로 갔는지 **전수**로 남긴다
                "no_source": ns,
                "unbuilt": ms,
                "saved_revisions": saved,
                # ★이미 끝난 것과 상한에 밀린 것 — **다시 안 산 이유**를 남긴다
                "already_done": list(already_done or ()),
                "deferred": list(deferred or ()),
                # ★정책 스텝이 **이 신원과 맞는 revision 만** 본다
                "expected_input_hashes": dict(expected or {}),
            },
        }
