"""요소 Phase StepRunner — entity_extract, entity_review, entity_detail, entity_t2i."""
import json
import logging
import time
from concurrent.futures import ThreadPoolExecutor, as_completed
from pathlib import Path
from typing import Any, Dict, List, Optional

from app.core.prompt_assembly import assemble_user_prompt
from app.core.step_runner import StepRunner
from app.modules.llm.llm_client import call_structured
from app.modules.prompt_loader import load_prompt

logger = logging.getLogger(__name__)


def build_distinctive_trait_block(
    distinctive_visual_traits: Optional[List[Dict[str, Any]]],
) -> tuple[str, List[Dict[str, Any]]]:
    """② P0a generic fix (2026-06-18) — evidence-gated distinctive 외형 블록.

    entity_detail 의 ``distinctive_visual_traits`` 중 **렌더용으로 채택할 항목만**
    골라 entity_t2i context 로 넘길 별도 블록 문자열을 만든다. 채택 조건은 두 개의
    **구조적 신호**다 (trait 텍스트를 글자/regex/substring 으로 의미 판정하지 않음
    — feedback-no-literal-substring-meaning. literal/metaphor 판단은 LLM 이
    ``interpretation_kind`` 로 내려준 구조화 결과를 코드가 그대로 신뢰한다):

      1. ``source_quote`` 가 비어 있지 않다 (본문 직접 근거 존재).
      2. ``interpretation_kind == "literal_visual"`` (신체 부위의 색/형태/표식을
         평이하게 사실로 진술한 경우. metaphorical_impression / ambiguous 는 제외).

    이로써 (a) 본문 근거 없는 invention(빈 quote), (b) 근거는 있으나 시적·인상적
    비유(예: '금빛 물결이 일렁이는 눈동자')는 둘 다 t2i_prompt(ref/scene SOT)로
    승격되지 않는다. 반대로 평이한 색 진술(예: '붉은 눈동자')은 보존된다.

    Returns:
        (block, excluded): ``block`` 은 t2i context 에 append 할 문자열(채택 항목이
        없으면 빈 문자열), ``excluded`` 는 채택되지 않은 distinctive 항목 dict 목록
        (각 trait/interpretation_kind/interpretation_reason 포함 — '근거 있는데 왜
        빠졌나' 추적용 diagnostic).
    """
    dvts = distinctive_visual_traits or []
    woven: List[Dict[str, Any]] = []
    excluded: List[Dict[str, Any]] = []
    for d in dvts:
        if not isinstance(d, dict):
            continue
        has_quote = bool((d.get("source_quote") or "").strip())
        is_literal = d.get("interpretation_kind") == "literal_visual"
        (woven if (has_quote and is_literal) else excluded).append(d)

    if not woven:
        return "", excluded

    lines = "\n".join(
        f"- {d.get('trait', '')} (원문 인용: {d.get('source_quote', '')}"
        + (
            f"; 연결근거: {d['attribution_note']}"
            if (d.get("attribution_note") or "").strip()
            else ""
        )
        + ")"
        for d in woven
    )
    return f"\n\n[원문 인용으로 확인된 문자적 특이 외형 — 반드시 반영]\n{lines}", excluded


class _EntityStepMixin:
    """요소 단계 공통 — fulltext 로드, 이전 단계 결과 로드."""

    def _load_fulltext(self) -> str:
        from app.models.project import Episode
        from sqlalchemy.orm import undefer

        ep = (
            self.db.query(Episode)
            .options(undefer(Episode.fulltext))
            .filter(Episode.id == self.episode_id)
            .first()
        )
        if not ep or not ep.fulltext:
            from app.core.errors import AppError
            raise AppError(
                code="step.no_fulltext",
                message="시나리오 텍스트가 없습니다.",
                status_code=400,
            )
        return ep.fulltext

    def _load_cleaned_text(self) -> str:
        cp = self._load_prev_checkpoint("text_cleanup")
        if cp and cp.get("data", {}).get("cleaned_text"):
            return cp["data"]["cleaned_text"]
        return self._load_fulltext()

    def _load_prev_checkpoint(self, step_id: str) -> Optional[Dict]:
        from app.core.config import settings

        cp = (
            Path(settings.projects_dir)
            / self.project_id
            / "checkpoints"
            / "episodes"
            / self.episode_id
            / step_id
            / "manifest.json"
        )
        if cp.exists():
            return json.loads(cp.read_text(encoding="utf-8"))
        return None


# ── entity_all: 리스팅 단계 (이름 + 출현 씬 + short_id 부여) ──


def _assign_short_ids(entities: List[Dict], prefix: str) -> List[Dict]:
    """엔티티 리스트에 short_id를 부여한다. prefix: C/L/P."""
    for i, e in enumerate(entities, 1):
        e["short_id"] = f"{prefix}{i:02d}"
    return entities


class EntityAllCharacterStep(_EntityStepMixin, StepRunner):
    """인물 리스팅 — shot 기반 1회 호출 (fallback: 씬 체이닝)."""

    def _execute(self, mode="resume") -> Dict[str, Any]:
        rules_cp = self._load_prev_checkpoint("visual_world_rules")
        visual_rules = json.dumps(rules_cp.get("data", {}), ensure_ascii=False) if rules_cp else ""

        # shot 기반 추출
        shot_cp = self._load_prev_checkpoint("shot_validator")
        if shot_cp:
            from app.core.steps.shot_validator_step import assert_no_failed_scenes
            assert_no_failed_scenes(shot_cp, self.project_config, consumer_step="entity_all_character")
        shots_scenes = shot_cp.get("data", {}).get("scenes", []) if shot_cp else []
        total_shots = sum(len(s.get("shots", [])) for s in shots_scenes)

        if total_shots > 0:
            from app.modules.pipeline.entity_lister import list_characters_from_shots
            characters = list_characters_from_shots(
                shots_scenes, visual_rules,
                self.project_config, self.build_opik_metadata(),
            )
        else:
            logger.warning("entity_all_character: no shots found, falling back to scene chaining")
            save_cp = self._load_prev_checkpoint("scene_save")
            scenes = save_cp.get("data", {}).get("segments", []) if save_cp else []
            from app.modules.pipeline.entity_lister import list_entities_by_type
            characters = list_entities_by_type("", "character", visual_rules, self.project_config, self.build_opik_metadata(), scenes=scenes)

        _assign_short_ids(characters, "C")
        return {"completed_count": len(characters), "applicable_count": len(characters), "failed_count": 0, "data": {"characters": characters}}


class EntityAllLocationStep(_EntityStepMixin, StepRunner):
    """배경 리스팅 — shot 기반 1회 호출."""

    def _execute(self, mode="resume") -> Dict[str, Any]:
        rules_cp = self._load_prev_checkpoint("visual_world_rules")
        visual_rules = json.dumps(rules_cp.get("data", {}), ensure_ascii=False) if rules_cp else ""

        shot_cp = self._load_prev_checkpoint("shot_validator")
        shots_scenes = shot_cp.get("data", {}).get("scenes", []) if shot_cp else []
        total_shots = sum(len(s.get("shots", [])) for s in shots_scenes)

        if total_shots > 0:
            from app.modules.pipeline.entity_lister import list_entities_from_shots
            locations = list_entities_from_shots(shots_scenes, "location", visual_rules, self.project_config, self.build_opik_metadata())
        else:
            logger.warning("entity_all_location: no shots, falling back to scene chaining")
            save_cp = self._load_prev_checkpoint("scene_save")
            scenes = save_cp.get("data", {}).get("segments", []) if save_cp else []
            from app.modules.pipeline.entity_lister import list_entities_by_type
            locations = list_entities_by_type("", "location", visual_rules, self.project_config, self.build_opik_metadata(), scenes=scenes)

        _assign_short_ids(locations, "L")
        return {"completed_count": len(locations), "applicable_count": len(locations), "failed_count": 0, "data": {"locations": locations}}


class EntityAllPropStep(_EntityStepMixin, StepRunner):
    """소품 리스팅 — shot 기반 1회 호출."""

    def _execute(self, mode="resume") -> Dict[str, Any]:
        rules_cp = self._load_prev_checkpoint("visual_world_rules")
        visual_rules = json.dumps(rules_cp.get("data", {}), ensure_ascii=False) if rules_cp else ""

        shot_cp = self._load_prev_checkpoint("shot_validator")
        shots_scenes = shot_cp.get("data", {}).get("scenes", []) if shot_cp else []
        total_shots = sum(len(s.get("shots", [])) for s in shots_scenes)

        if total_shots > 0:
            from app.modules.pipeline.entity_lister import list_entities_from_shots
            props = list_entities_from_shots(shots_scenes, "prop", visual_rules, self.project_config, self.build_opik_metadata())
        else:
            logger.warning("entity_all_prop: no shots, falling back to scene chaining")
            save_cp = self._load_prev_checkpoint("scene_save")
            scenes = save_cp.get("data", {}).get("segments", []) if save_cp else []
            from app.modules.pipeline.entity_lister import list_entities_by_type
            props = list_entities_by_type("", "prop", visual_rules, self.project_config, self.build_opik_metadata(), scenes=scenes)

        _assign_short_ids(props, "P")
        return {"completed_count": len(props), "applicable_count": len(props), "failed_count": 0, "data": {"props": props}}


class EntityExtractStep(_EntityStepMixin, StepRunner):
    """Step 8 (legacy): 요소 추출 (3턴: 인물 -> 배경 -> 소품)."""

    def _execute(self, mode="resume") -> Dict[str, Any]:
        fulltext = self._load_cleaned_text()

        rules_cp = self._load_prev_checkpoint("visual_world_rules")
        visual_rules = (
            json.dumps(rules_cp.get("data", {}), ensure_ascii=False)
            if rules_cp
            else ""
        )

        save_cp = self._load_prev_checkpoint("scene_save")
        segments_json = (
            json.dumps(
                save_cp.get("data", {}).get("segments", []), ensure_ascii=False,
            )
            if save_cp
            else ""
        )

        from app.modules.pipeline.entity_extractor_v4 import extract_all_entities

        result = extract_all_entities(
            fulltext,
            visual_rules,
            segments_json,
            self.project_config,
            self.build_opik_metadata(),
        )

        total = (
            len(result.get("characters", []))
            + len(result.get("locations", []))
            + len(result.get("props", []))
        )
        return {
            "completed_count": total,
            "applicable_count": total,
            "failed_count": 0,
            "data": result,
        }


def _inject_short_ids(entities: List[Dict], sid_map: Dict[str, str]) -> List[Dict]:
    """entity_all의 name→short_id 매핑을 extract 결과에 주입."""
    for e in entities:
        if not e.get("short_id"):
            e["short_id"] = sid_map.get(e.get("name", ""), "")
    return entities


class EntityExtractCharacterStep(_EntityStepMixin, StepRunner):
    """Step 8: 인물 추출 — entity_all_character 리스트 기반으로 상세 설명 추가."""

    def _execute(self, mode="resume") -> Dict[str, Any]:
        fulltext = self._load_cleaned_text()
        rules_cp = self._load_prev_checkpoint("visual_world_rules")
        visual_rules = json.dumps(rules_cp.get("data", {}), ensure_ascii=False) if rules_cp else ""

        all_cp = self._load_prev_checkpoint("entity_all_character")
        if all_cp and all_cp.get("data", {}).get("characters"):
            name_list = all_cp["data"]["characters"]
            sid_map = {e["name"]: e.get("short_id", "") for e in name_list}
            from app.modules.pipeline.entity_extractor_v4 import extract_entities_by_type_with_list
            characters = extract_entities_by_type_with_list(fulltext, "character", name_list, visual_rules, self.project_config, self.build_opik_metadata())
            _inject_short_ids(characters, sid_map)
        else:
            from app.modules.pipeline.entity_extractor_v4 import extract_entities_by_type
            characters = extract_entities_by_type(fulltext, "character", visual_rules, "", self.project_config, self.build_opik_metadata())
        return {"completed_count": len(characters), "applicable_count": len(characters), "failed_count": 0, "data": {"characters": characters}}


class EntityExtractLocationStep(_EntityStepMixin, StepRunner):
    """Step 9: 배경 추출 — entity_all_location 리스트 기반으로 상세 설명 추가."""

    def _execute(self, mode="resume") -> Dict[str, Any]:
        fulltext = self._load_cleaned_text()
        rules_cp = self._load_prev_checkpoint("visual_world_rules")
        visual_rules = json.dumps(rules_cp.get("data", {}), ensure_ascii=False) if rules_cp else ""

        all_cp = self._load_prev_checkpoint("entity_all_location")
        if all_cp and all_cp.get("data", {}).get("locations"):
            name_list = all_cp["data"]["locations"]
            sid_map = {e["name"]: e.get("short_id", "") for e in name_list}
            from app.modules.pipeline.entity_extractor_v4 import extract_entities_by_type_with_list
            locations = extract_entities_by_type_with_list(fulltext, "location", name_list, visual_rules, self.project_config, self.build_opik_metadata())
            _inject_short_ids(locations, sid_map)
        else:
            from app.modules.pipeline.entity_extractor_v4 import extract_entities_by_type
            locations = extract_entities_by_type(fulltext, "location", visual_rules, "", self.project_config, self.build_opik_metadata())
        return {"completed_count": len(locations), "applicable_count": len(locations), "failed_count": 0, "data": {"locations": locations}}


class EntityExtractPropStep(_EntityStepMixin, StepRunner):
    """Step 10: 소품 추출 — entity_all_prop 리스트 기반으로 상세 설명 추가."""

    def _execute(self, mode="resume") -> Dict[str, Any]:
        fulltext = self._load_cleaned_text()
        rules_cp = self._load_prev_checkpoint("visual_world_rules")
        visual_rules = json.dumps(rules_cp.get("data", {}), ensure_ascii=False) if rules_cp else ""

        all_cp = self._load_prev_checkpoint("entity_all_prop")
        if all_cp and all_cp.get("data", {}).get("props"):
            name_list = all_cp["data"]["props"]
            sid_map = {e["name"]: e.get("short_id", "") for e in name_list}
            from app.modules.pipeline.entity_extractor_v4 import extract_entities_by_type_with_list
            props = extract_entities_by_type_with_list(fulltext, "prop", name_list, visual_rules, self.project_config, self.build_opik_metadata())
            _inject_short_ids(props, sid_map)
        else:
            from app.modules.pipeline.entity_extractor_v4 import extract_entities_by_type
            props = extract_entities_by_type(fulltext, "prop", visual_rules, "", self.project_config, self.build_opik_metadata())
        return {"completed_count": len(props), "applicable_count": len(props), "failed_count": 0, "data": {"props": props}}


# ── Step 11: entity_filter ──


class EntityFilterStep(_EntityStepMixin, StepRunner):
    """Step 11: 저빈도 요소 필터링 (3씬 이하 -> LLM 판단)."""

    def _execute(self, mode="resume") -> Dict[str, Any]:
        fulltext = self._load_cleaned_text()

        # entity_merge 결과 로드 (병합 후 필터링)
        merge_cp = self._load_prev_checkpoint("entity_merge")
        if merge_cp and merge_cp.get("data"):
            entities = {
                "characters": merge_cp["data"].get("characters", []),
                "locations": merge_cp["data"].get("locations", []),
                "props": merge_cp["data"].get("props", []),
            }
        else:
            # fallback: 개별 추출 결과
            char_cp = self._load_prev_checkpoint("entity_extract_character")
            loc_cp = self._load_prev_checkpoint("entity_extract_location")
            prop_cp = self._load_prev_checkpoint("entity_extract_prop")
            entities = {
                "characters": (char_cp or {}).get("data", {}).get("characters", []),
                "locations": (loc_cp or {}).get("data", {}).get("locations", []),
                "props": (prop_cp or {}).get("data", {}).get("props", []),
            }

        # Load segments for scene context
        save_cp = self._load_prev_checkpoint("scene_save")
        segments = save_cp.get("data", {}).get("segments", []) if save_cp else []

        # entity_all에서 shot_count 복원 (extract/merge에서 유실됨)
        for etype, prefix in [("characters", "entity_all_character"), ("locations", "entity_all_location"), ("props", "entity_all_prop")]:
            all_cp = self._load_prev_checkpoint(prefix)
            if all_cp and all_cp.get("data"):
                key = etype
                count_map = {e["name"]: e.get("shot_count") or e.get("scene_count", 0)
                             for e in all_cp["data"].get(key, [])}
                for e in entities.get(etype, []):
                    if "shot_count" not in e:
                        e["shot_count"] = count_map.get(e["name"], 0)

        # entity_relation에서 변형 관계(visual_similarity=true)가 있는 요소는 보호
        protected_sids: set = set()
        rel_cp = self._load_prev_checkpoint("entity_relation")
        if rel_cp and rel_cp.get("data", {}).get("relations"):
            for rel in rel_cp["data"]["relations"]:
                if rel.get("visual_similarity"):
                    protected_sids.add(rel.get("base_short_id", ""))
                    protected_sids.add(rel.get("variant_short_id", ""))
            protected_sids.discard("")

        from app.modules.pipeline.entity_filter import filter_low_frequency_entities

        result = filter_low_frequency_entities(
            entities=entities,
            segments=segments,
            fulltext=fulltext,
            max_scenes=2,
            protected_short_ids=protected_sids if protected_sids else None,
            project_config=self.project_config,
            opik_metadata=self.build_opik_metadata(),
        )

        return {
            "completed_count": 1,
            "applicable_count": 1,
            "failed_count": 0,
            "data": result,
        }


# ── Step 12: entity_review ──


class EntityReviewStep(_EntityStepMixin, StepRunner):
    """Step 12: 요소 교차 검증 — 다른 모델로 추출 결과 교차 확인."""

    def _execute(self, mode="resume") -> Dict[str, Any]:
        fulltext = self._load_cleaned_text()

        # Try entity_filter checkpoint first (has filtered results)
        filter_cp = self._load_prev_checkpoint("entity_filter")
        if filter_cp and filter_cp.get("data", {}).get("filtered_entities"):
            extract_data = filter_cp["data"]["filtered_entities"]
        else:
            # Fallback: Load from 3 separate checkpoints (v3 split) or legacy single checkpoint
            char_cp = self._load_prev_checkpoint("entity_extract_character")
            loc_cp = self._load_prev_checkpoint("entity_extract_location")
            prop_cp = self._load_prev_checkpoint("entity_extract_prop")

            if char_cp or loc_cp or prop_cp:
                extract_data = {
                    "characters": (char_cp or {}).get("data", {}).get("characters", []),
                    "locations": (loc_cp or {}).get("data", {}).get("locations", []),
                    "props": (prop_cp or {}).get("data", {}).get("props", []),
                }
            else:
                # Legacy fallback: single entity_extract checkpoint
                prev = self._load_prev_checkpoint("entity_extract")
                if not prev or not prev.get("data"):
                    from app.core.errors import AppError
                    raise AppError(
                        code="step.no_input",
                        message="entity_extract 결과 없음",
                        status_code=400,
                    )
                extract_data = prev["data"]

        from app.modules.pipeline.entity_reviewer import review_entities

        review_result = review_entities(
            entities=extract_data,
            fulltext=fulltext,
            project_config=self.project_config,
            opik_metadata=self.build_opik_metadata(),
        )

        # 불필요 항목 필터링 — 일시 비활성화 (review 결과만 기록, 제거 안 함)
        filtered = {
            "characters": list(extract_data.get("characters", [])),
            "locations": list(extract_data.get("locations", [])),
            "props": list(extract_data.get("props", [])),
            "review": review_result,
        }

        total = (
            len(filtered["characters"])
            + len(filtered["locations"])
            + len(filtered["props"])
        )
        return {
            "completed_count": 1,
            "applicable_count": 1,
            "failed_count": 0,
            "data": filtered,
        }


# ── Step 13.5: entity_merge ──


class EntityMergeStep(_EntityStepMixin, StepRunner):
    """요소 중복 병합 — 타입 간 중복/유사 요소를 LLM으로 판별하여 제거."""

    def _execute(self, mode="resume") -> Dict[str, Any]:
        # 3개 extract 결과 로드
        char_cp = self._load_prev_checkpoint("entity_extract_character")
        loc_cp = self._load_prev_checkpoint("entity_extract_location")
        prop_cp = self._load_prev_checkpoint("entity_extract_prop")

        characters = (char_cp or {}).get("data", {}).get("characters", [])
        locations = (loc_cp or {}).get("data", {}).get("locations", [])
        props = (prop_cp or {}).get("data", {}).get("props", [])

        total = len(characters) + len(locations) + len(props)
        if total == 0:
            return {"completed_count": 0, "applicable_count": 0, "failed_count": 0,
                    "data": {"characters": [], "locations": [], "props": [], "removed": []}}

        # visual_world_rules
        vwr_cp = self._load_prev_checkpoint("visual_world_rules")
        vwr_data = vwr_cp.get("data", {}) if vwr_cp else {}
        era = vwr_data.get("era", "")
        region = vwr_data.get("region", "")

        # scene_summary
        sum_cp = self._load_prev_checkpoint("scene_summary")
        summaries = sum_cp.get("data", {}).get("summaries", []) if sum_cp else []
        summary_text = "\n".join(
            f"씬{s.get('scene_index', '?')}: {s.get('scene_summary', '')}"
            for s in summaries
        )

        # 요소 목록 구성
        entity_lines = []
        for c in characters:
            entity_lines.append(f"{c.get('short_id','')} {c['name']} (character): {c.get('description','')}")
        for l in locations:
            entity_lines.append(f"{l.get('short_id','')} {l['name']} (location): {l.get('description','')}")
        for p in props:
            entity_lines.append(f"{p.get('short_id','')} {p['name']} (prop): {p.get('description','')}")

        user_prompt = (
            f"[세계관] 시대: {era}, 지역: {region}\n\n"
            f"[씬별 요약]\n{summary_text}\n\n"
            f"[전체 요소 목록]\n" + "\n".join(entity_lines) + "\n\n"
            "위 요소 목록에서 타입이 다르지만 실질적으로 같은 대상인 항목을 찾으세요.\n"
            "예: 같은 물체가 배경(location)과 소품(prop)에 동시 등록된 경우.\n"
            "중복이 있으면 제거할 short_id를 선택하세요. 없으면 빈 배열을 반환하세요.\n"
            "각 타입 쌍을 고르게 검토하되, 실질적으로 같은 대상이 서로 다른 타입으로 등록된 경우를 찾으세요."
        )

        remove_schema = {
            "type": "object",
            "properties": {
                "remove": {
                    "type": "array",
                    "items": {"type": "string"},
                    "description": "제거할 short_id 목록",
                },
            },
            "required": ["remove"],
            "additionalProperties": False,
        }

        removed = []
        try:
            result = call_structured(
                step="entity_merge",
                system_prompt="당신은 시나리오 요소 분석 전문가입니다. 중복 요소를 찾아 제거할 항목을 판별합니다.",
                user_prompt=user_prompt,
                response_schema=remove_schema,
                project_config=self.project_config,
                schema_name="entity_merge",
                opik_metadata=self.build_opik_metadata(),
            )
            removed = result.get("remove", [])
        except Exception as exc:
            logger.warning("Entity merge LLM call failed: %s", exc)

        # 제거 적용
        if removed:
            all_sids = {e.get("short_id", "") for e in characters + locations + props}
            unknown = [sid for sid in removed if sid not in all_sids]
            if unknown:
                logger.warning("Entity merge: LLM returned unknown short_ids: %s", unknown)
            remove_set = set(removed) & all_sids  # 존재하는 것만
            characters = [c for c in characters if c.get("short_id", "") not in remove_set]
            locations = [l for l in locations if l.get("short_id", "") not in remove_set]
            props = [p for p in props if p.get("short_id", "") not in remove_set]
            logger.info("Entity merge: removed %s", list(remove_set))

        new_total = len(characters) + len(locations) + len(props)
        return {
            "completed_count": new_total,
            "applicable_count": total,
            "failed_count": 0,
            "data": {
                "characters": characters,
                "locations": locations,
                "props": props,
                "removed": removed,
            },
        }


# ── Step 10: entity_detail ──


class EntityDetailStep(_EntityStepMixin, StepRunner):
    """Step 10: 요소 상세 — 시각적 상세 정보 + short_id 확정."""

    def _execute(self, mode="resume") -> Dict[str, Any]:
        fulltext = self._load_cleaned_text()

        # entity_filter 우선 → entity_merge fallback → entity_extract fallback
        filter_cp = self._load_prev_checkpoint("entity_filter")
        if filter_cp and filter_cp.get("data", {}).get("filtered_entities"):
            filtered = filter_cp["data"]["filtered_entities"]
        else:
            merge_cp = self._load_prev_checkpoint("entity_merge")
            if merge_cp and "removed" in merge_cp.get("data", {}):
                filtered = {
                    "characters": merge_cp["data"].get("characters", []),
                    "locations": merge_cp["data"].get("locations", []),
                    "props": merge_cp["data"].get("props", []),
                }
            else:
                char_cp = self._load_prev_checkpoint("entity_extract_character")
                loc_cp = self._load_prev_checkpoint("entity_extract_location")
                prop_cp = self._load_prev_checkpoint("entity_extract_prop")
                filtered = {
                    "characters": (char_cp or {}).get("data", {}).get("characters", []),
                    "locations": (loc_cp or {}).get("data", {}).get("locations", []),
                    "props": (prop_cp or {}).get("data", {}).get("props", []),
                }

        if not any(filtered.values()):
            from app.core.errors import AppError
            raise AppError(
                code="step.no_input",
                message="entity_extract 결과 없음",
                status_code=400,
            )

        # 엔티티 큐 구성 (name, type, short_id) — (name, type) 기준 중복 제거
        entity_queue: List[tuple] = []
        seen_keys: set = set()
        for etype_key, etype_label in [("characters", "character"), ("locations", "location"), ("props", "prop")]:
            for e in filtered.get(etype_key, []):
                name = e["name"]
                key = (name, etype_label)
                if key in seen_keys:
                    logger.warning("Entity queue: 중복 제거 — %s (%s)", name, etype_label)
                    continue
                seen_keys.add(key)
                entity_queue.append((name, etype_label, e.get("short_id", "")))

        from app.modules.pipeline.entity_extractor_v3 import _load_prompt, _load_schema

        # visual_world_rules 로드 — description 생성에 시대/의상 맥락 반영
        vwr_cp = self._load_prev_checkpoint("visual_world_rules")
        vwr_data = vwr_cp.get("data", {}) if vwr_cp else {}
        world_block = ""
        if vwr_data.get("era") or vwr_data.get("region"):
            world_block = f"\n[세계관]\n시대: {vwr_data.get('era', '')}\n지역: {vwr_data.get('region', '')}"
            for r in vwr_data.get("rules", []):
                if r.get("rule_type") == "costume":
                    world_block += f"\n의상규칙: {r.get('description', '')}"

        # 기획서 인물 정보 보강 (첫 에피소드만)
        from app.core.planning_doc_context import get_planning_context
        pctx = get_planning_context(self.project_id, self.episode_id, self.db)
        planning_block = pctx.inject_if_available("characters_text", "## 기획서 인물 참고 정보")

        detail_schema = _load_schema("turn1_7_detail_batch_schema.json")
        entity_list_text = "\n".join(
            f"- {name} ({etype})" for name, etype, *_ in entity_queue
        )
        detail_prompt = _load_prompt(
            "turn1_7_detail_batch",
            entity_list=entity_list_text,
            fulltext=fulltext,
        ) + world_block + planning_block

        # key = "name:type" (동명 다른 타입 구분)
        entity_details: Dict[str, Dict] = {}

        try:
            batch_result = call_structured(
                step="entity_detail_batch",
                system_prompt=self._load_prompt("entity_extractor_v2", "system") if hasattr(self, '_load_prompt') else load_prompt("entity_extractor_v2", "system"),
                user_prompt=detail_prompt,
                response_schema=detail_schema,
                project_config=self.project_config,
                schema_name="entity_detail_batch",
                opik_metadata=self.build_opik_metadata(),
            )
            for ent in batch_result.get("entities", []):
                ename = ent["name"]
                etype = ent.get("entity_type", "")
                key = f"{ename}:{etype}" if etype else ename
                entity_details[key] = {
                    "description": ent.get("description", ""),
                    "visual_traits": ent.get("visual_traits", []),
                    # ② P0a generic fix (2026-06-18): evidence-gated distinctive
                    # 외형. EntityT2iStep 이 source_quote + interpretation_kind 로
                    # 게이트하므로 여기선 raw 보존 (gate 는 build 시점). 누락 시 t2i
                    # 가 distinctive 를 못 봄.
                    "distinctive_visual_traits": ent.get("distinctive_visual_traits", []),
                }

            # 누락 재시도
            missing = [
                (n, t) for n, t, *_ in entity_queue
                if f"{n}:{t}" not in entity_details and n not in entity_details
            ]
            if missing:
                missing_text = "\n".join(
                    f"- {n} ({t})" for n, t in missing
                )
                retry_prompt = _load_prompt(
                    "turn1_7_detail_batch",
                    entity_list=missing_text,
                    fulltext=fulltext,
                )
                try:
                    retry_result = call_structured(
                        step="entity_detail_batch",
                        system_prompt=self._load_prompt("entity_extractor_v2", "system") if hasattr(self, '_load_prompt') else load_prompt("entity_extractor_v2", "system"),
                        user_prompt=retry_prompt,
                        response_schema=detail_schema,
                        project_config=self.project_config,
                        schema_name="entity_detail_retry",
                        opik_metadata=self.build_opik_metadata(["retry"]),
                    )
                    for ent in retry_result.get("entities", []):
                        ename = ent["name"]
                        etype = ent.get("entity_type", "")
                        key = f"{ename}:{etype}" if etype else ename
                        entity_details[key] = {
                            "description": ent.get("description", ""),
                            "visual_traits": ent.get("visual_traits", []),
                        }
                except Exception as exc:
                    logger.warning("Entity detail retry parse failed: %s", exc)
        except Exception as exc:
            logger.warning("Entity detail batch failed: %s", exc)

        return {
            "completed_count": len(entity_details),
            "applicable_count": len(entity_queue),
            "failed_count": max(0, len(entity_queue) - len(entity_details)),
            "data": {"entity_queue": entity_queue, "entity_details": entity_details},
        }


# ── Step 11: entity_t2i ──


class EntityT2iStep(_EntityStepMixin, StepRunner):
    """Step 11: T2I 프롬프트 병렬 생성."""

    def _execute(self, mode="resume") -> Dict[str, Any]:
        prev = self._load_prev_checkpoint("entity_detail")
        if not prev or not prev.get("data"):
            from app.core.errors import AppError
            raise AppError(
                code="step.no_input",
                message="entity_detail 결과 없음",
                status_code=400,
            )

        entity_queue = prev["data"]["entity_queue"]
        # 하위 호환: gpt_details → entity_details
        entity_details = prev["data"].get("entity_details") or prev["data"].get("gpt_details", {})

        # name→short_id 매핑 구축 (entity_queue에 3번째 요소로 포함)
        name_to_sid = {}
        for item in entity_queue:
            if len(item) >= 3 and item[2]:
                name_to_sid[item[0]] = item[2]

        # resume: 이미 완료된 것 스킵
        cp = self.load_checkpoint()
        done = (
            cp.get("data", {}).get("completed", {})
            if cp and mode == "resume"
            else {}
        )
        # ★현재 대상 밖의 완료분은 버린다 (2026-08-07 실측·Codex 3차 리뷰).
        #
        # 체크포인트가 archive 에서 자동 복원될 수 있고(상류 force 가 하류를
        # 지운 뒤에도 그랬다), 그때 복원본에는 **지금은 없는 엔티티**의 완료
        # 기록이 들어 있다. 그것을 그대로 받으면 `len(done)` 이 현재 대상
        # 수를 넘어 `failed_count = total - len(done)` 가 음수가 되고
        # (실측 154/97, failed=-57), 그 partial 이 entity sync 의 stale
        # cleanup 을 막아 UniqueViolation 으로 파이프라인이 멈췄다.
        #
        # 표식 층(`invalidate_downstream` 의 marker)과 **둘 다** 막는다 —
        # 한 층만 막으면 다른 경로로 같은 오염이 들어온다.
        if done:
            valid = {n for (n, _t, *_rest) in entity_queue}
            stale_names = [n for n in done if n not in valid]
            if stale_names:
                logger.warning(
                    "entity_t2i: 현재 대상에 없는 완료 기록 %d건 폐기 "
                    "(archive 복원 잔존 추정) — 표본 %s",
                    len(stale_names), stale_names[:8],
                )
                done = {n: v for n, v in done.items() if n in valid}

        remaining = [
            (i, n, t)
            for i, (n, t, *_) in enumerate(entity_queue)
            if n not in done
        ]
        total = len(entity_queue)

        from app.modules.pipeline.entity_extractor_v3 import (
            ENTITY_DETAIL_SCHEMA,
            _load_system,
            _load_prompt,
        )

        system_prompt = _load_system()

        # visual_world_rules에서 era/region + costume rules 로드
        vwr_cp = self._load_prev_checkpoint("visual_world_rules")
        vwr_data = vwr_cp.get("data", {}) if vwr_cp else {}
        era = vwr_data.get("era", "")
        region = vwr_data.get("region", "")
        world_context = ""
        if era or region:
            world_context = f"\n\n[세계관]\n시대: {era}\n지역: {region}"
            for r in vwr_data.get("rules", []):
                if r.get("rule_type") == "costume":
                    world_context += f"\n의상규칙: {r.get('description', '')}"
        t2i_ctx = vwr_data.get("t2i_context", "")
        if t2i_ctx:
            world_context += f"\n\n[T2I 시각 컨텍스트]\n{t2i_ctx}"

        def _gen_t2i(idx, ename, etype):
            # name:type key 우선, fallback name만 (하위 호환)
            ent_detail = entity_details.get(f"{ename}:{etype}") or entity_details.get(ename, {})
            # 덧붙임 블록을 둘로 나눠 둔다 — world_context 는 프로젝트·주행
            # 단위로 고정이고 detail_block 은 호출마다 변한다. 축 A1 재배열이
            # 그 경계를 쓴다. 이어 붙이는 순서는 예전과 같다(OFF=무변화).
            detail_block = ""
            if ent_detail:
                desc = ent_detail.get("description", "")
                traits = ", ".join(ent_detail.get("visual_traits", []))
                detail_block += (
                    f"\n\n[시나리오 기반 상세 정보]\n"
                    f"설명: {desc}\n시각적 특징: {traits}"
                )
                # ② P0a generic fix (2026-06-18): evidence-gated distinctive 외형.
                # literal_visual + source_quote 항목만 별도 블록으로 t2i context weave.
                _dblock, _excluded = build_distinctive_trait_block(
                    ent_detail.get("distinctive_visual_traits")
                )
                detail_block += _dblock
                if _excluded:
                    logger.info(
                        "entity_t2i %s: distinctive_visual_traits %d개 렌더 제외 "
                        "(evidence-gate, literal_visual+source_quote 아님): %s",
                        ename, len(_excluded),
                        [
                            {
                                "trait": d.get("trait"),
                                "kind": d.get("interpretation_kind"),
                                "has_quote": bool((d.get("source_quote") or "").strip()),
                                "reason": d.get("interpretation_reason"),
                            }
                            for d in _excluded
                        ],
                    )

            turn_msg = assemble_user_prompt(
                "entity_t2i",
                render=lambda **kw: _load_prompt("turn_entity_detail", **kw),
                fields={"entity_name": ename, "entity_type": etype},
                project_block=world_context,
                call_block=detail_block,
            )

            for attempt in range(3):
                try:
                    detail = call_structured(
                        step="entity_t2i",
                        system_prompt=system_prompt,
                        user_prompt=turn_msg,
                        response_schema=ENTITY_DETAIL_SCHEMA,
                        project_config=self.project_config,
                        schema_name="entity_t2i",
                        opik_metadata=self.build_opik_metadata(),
                    )
                    # description / visual_traits 는 source detail
                    # (entity_detail 결과) 에서 forward — LLM 의 unsourced
                    # trait 이 final output 에 도입되는 path 를 봉쇄. ent_detail
                    # 이 빈 dict 라도 LLM detail 로 silent fallback 안 함
                    # (feedback_no_silent_fallback). LLM 출력의 description /
                    # visual_traits 는 무시. t2i_prompt 만 LLM 출력에서 take —
                    # entity_t2i 의 단일 책임. short_description 은 후속 fix 로
                    # schema 및 출력 dict 에서 제거 (Codex MIN 1, consumer 0).
                    #
                    # D6 T2 (B4 #2): metadata_json 은 LLM 출력에서 take.
                    # entity_detail batch 가 만들지 않는 field 라 source-forward
                    # 정책 외 — location 의 space_profile 분류는 entity_t2i LLM
                    # 호출 시점에 system prompt 의 schema 가이드대로 LLM 이 출력.
                    # T2-fix (I3): schema required 라 항상 {} 또는 dict 보장.
                    # I7: character/prop 은 prompt 가 `{}` 출력 명시 (location 만
                    # 의미 있음). EntitySyncService 가 entity_type 으로 분기.
                    src = ent_detail or {}
                    md = detail.get("metadata_json") or {}

                    # Area B (Task 3 / C3) + D6 carry: location + prop 양쪽
                    # metadata SOT post-validation. 3 attempts 모두 fail 시 done
                    # 에 안 넣음 → failed_count 증가. character 는 semantic
                    # post-validation 부재 (기존 빈-marker 패턴 유지 — image gen
                    # preflight 가 catch).
                    #
                    # location 의 D6 space_profile 내부 검증은
                    # validate_entity_metadata_shape helper 안에서 통합 호출
                    # (Task 2 fix commit 8383c20). 본 site 별도 호출 불필요 —
                    # single point of validation. SpaceProfileError 는 raise 시
                    # attempt loop 의 except 가 catch → retry.
                    if etype in ("location", "prop"):
                        from app.core.entity_metadata import (
                            validate_entity_metadata_shape,
                        )
                        from app.core.errors import AppError
                        from app.core.bg_state_vocab import SpaceProfileError
                        try:
                            validate_entity_metadata_shape(
                                etype, md, short_id=name_to_sid.get(ename, "") or ename,
                            )
                        except (AppError, SpaceProfileError) as exc:
                            # Area B Task 3 quality review M2 fix:
                            # location 의 D6 space_profile 위반은 SpaceProfileError —
                            # `except AppError` 만으로는 catch 안 됨 → outer Exception
                            # 으로 떨어져 log warning 누락 (silent retry). 양쪽 모두
                            # catch + warn + raise 로 logging 대칭 보장.
                            exc_msg = getattr(exc, "message", None) or str(exc)
                            logger.warning(
                                "entity_t2i %s post-validate failed (attempt %d/3) "
                                "for %r: %s — retry",
                                etype, attempt + 1, ename, exc_msg,
                            )
                            raise

                    return (
                        idx,
                        ename,
                        etype,
                        {
                            "name": ename,
                            "short_id": name_to_sid.get(ename, ""),
                            "description": src.get("description", ""),
                            "visual_traits": src.get("visual_traits", []) if isinstance(src.get("visual_traits"), list) else [],
                            "t2i_prompt": detail.get("t2i_prompt", ""),
                            "metadata_json": md,
                        },
                    )
                except Exception:
                    if attempt < 2:
                        time.sleep(2 * (attempt + 1))
            # 3 attempts 모두 실패: source detail 보존 + t2i_prompt 빈 문자열
            # marker. downstream 이 빈 prompt 를 빈 image 로 흘리지 않도록
            # consumer 는 check 필요 (image gen 단계의 ref_image_pipeline 가
            # asset_readiness preflight 로 catch 가능).
            #
            # Area B (Task 3 / C3) + D6 carry: location + prop entity 는 failure
            # 시 fail-fast — SOT (space_profile / visual_identity) 손실은 silent
            # absorb 안 함. data=None 반환으로 caller 가 done 에 안 넣고
            # failed_count 증가. character 는 기존 빈-marker pattern 유지
            # (image gen preflight 가 catch).
            if etype in ("location", "prop"):
                return (idx, ename, etype, None)

            src = ent_detail or {}
            # Area B (Task 4 review C1 fix): character failure marker 도 normalized
            # closed shape `{"location": None, "visual_identity": None}` 출력 — Task 4
            # sync (Boundary 2) 의 strict shape contract 정합. 빈 dict `{}` 는 sync
            # 의 keys check 에서 fail-fast 발생 → character LLM 실패 시 entire
            # episode fail cascade 위험 차단. 의미는 동일 — t2i_prompt="" marker 가
            # image preflight 의 catch path 보존.
            return (
                idx,
                ename,
                etype,
                {
                    "name": ename,
                    "short_id": name_to_sid.get(ename, ""),
                    "description": src.get("description", ""),
                    "visual_traits": src.get("visual_traits", []) if isinstance(src.get("visual_traits"), list) else [],
                    "t2i_prompt": "",
                    "metadata_json": {"location": None, "visual_identity": None},
                },
            )

        if remaining:
            max_workers = min(10, len(remaining))
            with ThreadPoolExecutor(max_workers=max_workers) as executor:
                futures = {}
                for i, ename, etype in remaining:
                    if futures:
                        time.sleep(1)
                    futures[executor.submit(_gen_t2i, i, ename, etype)] = i

                for future in as_completed(futures):
                    result = future.result()
                    _, ename, etype, data = result
                    # T2-fix (review iter3 I1): location post-validation 실패 시
                    # data=None — done 에 안 넣음 → failed_count 증가 (line 851).
                    if data is None:
                        logger.error(
                            "entity_t2i: %r (%s) failed all 3 attempts — "
                            "excluded from done (failed_count++)",
                            ename, etype,
                        )
                        self.update_progress(len(done), total)
                        continue
                    data["entity_type"] = etype
                    done[ename] = data
                    self.update_progress(len(done), total)
                    # 증분 체크포인트 (crash resume용) — 분류 배열도 포함
                    _chars = [v for v in done.values() if v.get("entity_type") == "character"]
                    _locs = [v for v in done.values() if v.get("entity_type") == "location"]
                    _prps = [v for v in done.values() if v.get("entity_type") == "prop"]
                    self.save_checkpoint(
                        {
                            "status": "running",
                            "data": {
                                "completed": done,
                                "entity_queue": entity_queue,
                                "characters": _chars,
                                "locations": _locs,
                                "props": _prps,
                            },
                        }
                    )

        # 최종 분류
        characters = [v for v in done.values() if v.get("entity_type") == "character"]
        locations = [v for v in done.values() if v.get("entity_type") == "location"]
        props = [v for v in done.values() if v.get("entity_type") == "prop"]

        return {
            "completed_count": len(done),
            "applicable_count": total,
            "failed_count": total - len(done),
            "data": {
                "completed": done,
                "characters": characters,
                "locations": locations,
                "props": props,
            },
        }
