"""analysis_dispatch_service 단위 테스트 — Phase 4.1.

공식 dispatch 진입점(run-all / analyze / reanalyze-scenes)이 공유하는
step 선택 + background job 제출 경로를 검증.
"""
from __future__ import annotations

from unittest.mock import MagicMock, patch

import pytest

from app.core.errors import AppError
from app.services import analysis_dispatch_service as ds


# ── select_steps_for_category ───────────────────────────────


def test_select_steps_for_category_analysis_returns_active_analysis_only():
    db = MagicMock()
    # _project_has_planning_doc은 ProjectRegistry query — 기획서 없는 상태 mock
    db.query.return_value.filter.return_value.first.return_value = None

    step_ids = ds.select_steps_for_category(db, "p1", "analysis")

    assert len(step_ids) > 0
    # 모든 리턴 step은 analysis category & active
    from app.core.step_catalog import STEP_CATALOG
    for sid in step_ids:
        entry = STEP_CATALOG[sid]
        assert entry.category == "analysis"
        assert entry.lifecycle == "active"
        assert entry.applicability not in ("on_demand", "disabled")

    # 기획서 없는 경우 if_planning_doc step은 제외되어야 함
    for sid in step_ids:
        assert STEP_CATALOG[sid].applicability != "if_planning_doc"


def test_select_steps_for_category_planning_doc_inclusion():
    """planning_doc_text가 있으면 if_planning_doc step도 포함."""
    db = MagicMock()
    mock_project = MagicMock()
    # project_has_planning_doc() 의 DB-text path 는 strip 길이 >=100 자만 present 로
    # 인정 (planning_doc_analysis_service.py) — fixture 도 그 threshold 를 충족해야 함.
    mock_project.planning_doc_text = (
        "This is a generic planning document describing the project's characters, "
        "locations, visual tone, and overall narrative direction in sufficient detail."
    )
    db.query.return_value.filter.return_value.first.return_value = mock_project

    step_ids = ds.select_steps_for_category(db, "p1", "analysis")

    from app.core.step_catalog import STEP_CATALOG
    planning_steps = [
        sid for sid in step_ids
        if STEP_CATALOG[sid].applicability == "if_planning_doc"
    ]
    # manifest에 if_planning_doc step이 최소 1개 존재해야 (baseline 전제)
    assert len(planning_steps) >= 1


def test_select_steps_for_category_invalid_raises():
    db = MagicMock()
    with pytest.raises(AppError) as exc_info:
        ds.select_steps_for_category(db, "p1", "unknown_cat")
    assert exc_info.value.code == "dispatch.invalid_category"


def test_select_steps_for_category_all_includes_image():
    db = MagicMock()
    db.query.return_value.filter.return_value.first.return_value = None

    step_ids_all = ds.select_steps_for_category(db, "p1", "all")
    step_ids_analysis = ds.select_steps_for_category(db, "p1", "analysis")

    from app.core.step_catalog import STEP_CATALOG
    image_sids = [sid for sid in step_ids_all if STEP_CATALOG[sid].category == "image"]
    assert len(image_sids) > 0
    assert len(step_ids_all) >= len(step_ids_analysis)


def test_select_steps_for_category_image_no_cross_category_transitive_when_episode_none():
    """Fix 3.3 (M4 해소): category='image' 호출이 transitive deps를 자동 포함하지 않는다.

    이전 동작 (Phase 9.2 가드): image step의 analysis dep을 transitive로 끌어들임 →
    사용자가 "image만 force"를 의도해도 43 step 폭발하는 사고 trigger (2026-05-01 13:13).
    현재 동작 (Fix 3.3): 직접 image step만 반환. silent miss 방지는 prerequisite 검증으로 대체.

    `episode_id=None` (preview/list 호출) 시 prerequisite 검증 skip — direct만 반환.
    """
    db = MagicMock()
    db.query.return_value.filter.return_value.first.return_value = None

    step_ids = ds.select_steps_for_category(db, "p1", "image")

    # 모든 리턴 step은 image category여야 함 — analysis step 자동 포함 X.
    from app.core.step_catalog import STEP_CATALOG
    for sid in step_ids:
        assert STEP_CATALOG[sid].category == "image", (
            f"image dispatch에 cross-category step이 포함되면 안 됨: {sid} ({STEP_CATALOG[sid].category})"
        )


# ── select_scene_reanalysis_steps ────────────────────────────


def test_select_scene_reanalysis_steps_contains_scene_save_and_downstream():
    db = MagicMock()
    db.query.return_value.filter.return_value.first.return_value = None

    step_ids = ds.select_scene_reanalysis_steps(db, "p1")

    assert "scene_save" in step_ids
    # scene_save의 downstream 일부가 포함되어야 함
    from app.core.step_catalog import get_all_downstream_recursive
    downstream = set(get_all_downstream_recursive("scene_save"))
    # 최소 1개 이상의 downstream이 리턴에 포함
    overlap = [s for s in step_ids if s in downstream]
    assert len(overlap) >= 1


def test_select_scene_reanalysis_returns_ordered_by_step_order():
    db = MagicMock()
    db.query.return_value.filter.return_value.first.return_value = None

    step_ids = ds.select_scene_reanalysis_steps(db, "p1")

    from app.core.step_catalog import STEP_CATALOG
    orders = [STEP_CATALOG[s].order for s in step_ids]
    assert orders == sorted(orders), "reanalyze step은 order 오름차순이어야 한다"


# ── load_project_llm_config ───────────────────────────────────


def test_load_project_llm_config_returns_empty_when_no_settings():
    db = MagicMock()
    db.query.return_value.filter.return_value.first.return_value = None

    config = ds.load_project_llm_config(db, "p1")

    assert config == {}


def test_load_project_llm_config_parses_valid_json():
    db = MagicMock()
    mock_ps = MagicMock()
    mock_ps.llm_config_json = '{"model_override": {"scene_detail": "gpt-5"}}'
    db.query.return_value.filter.return_value.first.return_value = mock_ps

    config = ds.load_project_llm_config(db, "p1")

    assert config == {"model_override": {"scene_detail": "gpt-5"}}


def test_load_project_llm_config_returns_empty_on_corrupt_json():
    db = MagicMock()
    mock_ps = MagicMock()
    mock_ps.llm_config_json = "{not valid json}"
    db.query.return_value.filter.return_value.first.return_value = mock_ps

    config = ds.load_project_llm_config(db, "p1")

    assert config == {}


# ── dispatch_category_run ────────────────────────────────────


def test_dispatch_category_run_raises_when_already_running():
    db = MagicMock()
    db.execute.return_value.fetchone.return_value = ("proj_name", "ep_title")
    db.query.return_value.filter.return_value.first.return_value = None

    with patch.object(ds, "is_task_running", create=True):
        # is_task_running가 import 시점에만 필요하므로 module 내부로 patch
        pass

    with patch("app.core.task_registry.is_task_running", return_value=True):
        with pytest.raises(AppError) as exc_info:
            ds.dispatch_category_run("p1", "e1", "analysis", "resume", db)
        assert exc_info.value.code == "step.already_running"


def test_dispatch_category_run_submits_job_and_returns_job_key():
    db = MagicMock()
    db.execute.return_value.fetchone.return_value = ("proj_name", "ep_title")
    db.query.return_value.filter.return_value.first.return_value = None

    with patch("app.core.task_registry.is_task_running", return_value=False), \
         patch("app.services.analysis_dispatch_service.submit_background_job", return_value=True) as mock_submit, \
         patch.object(ds, "_is_step_satisfied", return_value=True):
        result = ds.dispatch_category_run("p1", "e1", "analysis", "resume", db)

    assert result["ok"] is True
    assert result["status"] == "started"
    assert result["job_key"] == "run_all:p1:e1:analysis"
    assert isinstance(result["steps"], list) and len(result["steps"]) > 0
    assert mock_submit.called


def test_dispatch_category_run_job_key_conflict_raises():
    db = MagicMock()
    db.execute.return_value.fetchone.return_value = ("proj_name", "ep_title")
    db.query.return_value.filter.return_value.first.return_value = None

    with patch("app.core.task_registry.is_task_running", return_value=False), \
         patch("app.services.analysis_dispatch_service.submit_background_job", return_value=False), \
         patch.object(ds, "_is_step_satisfied", return_value=True):
        with pytest.raises(AppError) as exc_info:
            ds.dispatch_category_run("p1", "e1", "analysis", "resume", db)
        assert exc_info.value.code == "step.already_running"


# ── dispatch_scene_reanalysis ────────────────────────────────


def test_dispatch_scene_reanalysis_uses_force_mode():
    """run_steps_batch가 run_mode='force'로 호출되도록 args가 구성되는지 확인."""
    import inspect

    db = MagicMock()
    db.execute.return_value.fetchone.return_value = ("proj_name", "ep_title")
    db.query.return_value.filter.return_value.first.return_value = None

    with patch("app.core.task_registry.is_task_running", return_value=False), \
         patch("app.services.analysis_dispatch_service.submit_background_job", return_value=True) as mock_submit:
        ds.dispatch_scene_reanalysis("p1", "e1", db)

    # target + args를 signature에 bind하여 positional-index 의존 제거
    target = mock_submit.call_args.kwargs["target"]
    args_tuple = mock_submit.call_args.kwargs["args"]
    sig = inspect.signature(target)
    bound = sig.bind(*args_tuple)
    bound.apply_defaults()
    assert bound.arguments["run_mode"] == "force"


def test_dispatch_category_run_and_reanalysis_block_each_other():
    """reanalyze job이 실행 중이면 dispatch_category_run도 차단되어야 한다 (Review Important #2)."""
    db = MagicMock()
    db.execute.return_value.fetchone.return_value = ("proj_name", "ep_title")
    db.query.return_value.filter.return_value.first.return_value = None

    def is_running(key):
        return key.endswith(":reanalyze")

    with patch("app.core.task_registry.is_task_running", side_effect=is_running):
        with pytest.raises(AppError) as exc_info:
            ds.dispatch_category_run("p1", "e1", "analysis", "resume", db)
        assert exc_info.value.code == "step.already_running"

    def is_running_analysis(key):
        return key.endswith(":analysis")

    with patch("app.core.task_registry.is_task_running", side_effect=is_running_analysis):
        with pytest.raises(AppError) as exc_info:
            ds.dispatch_scene_reanalysis("p1", "e1", db)
        assert exc_info.value.code == "step.already_running"


# ── run_steps_batch: episode.status 복구 (Review Critical #1) ─────


def test_recover_episode_status_resets_analyzing_to_error():
    """파이프라인 실패 시 Episode.status 'analyzing' → 'error'로 복구."""
    db = MagicMock()
    mock_episode = MagicMock()
    mock_episode.status = "analyzing"
    db.query.return_value.filter.return_value.first.return_value = mock_episode

    ds._recover_episode_status_on_failure(db, "e1", "Step scene_director failed")

    assert mock_episode.status == "error"
    assert mock_episode.analysis_error == "Step scene_director failed"
    assert db.commit.called


def test_recover_episode_status_noop_if_not_analyzing():
    """이미 analyzed/error인 에피소드는 건드리지 않는다."""
    db = MagicMock()
    mock_episode = MagicMock()
    mock_episode.status = "analyzed"
    db.query.return_value.filter.return_value.first.return_value = mock_episode

    ds._recover_episode_status_on_failure(db, "e1", "Ignored")

    assert mock_episode.status == "analyzed"
    assert not db.commit.called


def test_recover_episode_status_truncates_long_message():
    """analysis_error 필드 보호 — 2000자 초과 메시지 잘림."""
    db = MagicMock()
    mock_episode = MagicMock()
    mock_episode.status = "analyzing"
    db.query.return_value.filter.return_value.first.return_value = mock_episode

    long_msg = "x" * 5000
    ds._recover_episode_status_on_failure(db, "e1", long_msg)

    assert len(mock_episode.analysis_error) == 2000


def test_recover_episode_status_handles_missing_episode():
    """Episode 없어도 예외 없이 종료."""
    db = MagicMock()
    db.query.return_value.filter.return_value.first.return_value = None

    # 예외 없어야 함
    ds._recover_episode_status_on_failure(db, "e1", "anything")


# ── run_steps_batch 직접 테스트 (Codex Review Minor #1) ────────


def test_run_steps_batch_step_failed_triggers_recovery():
    """step이 failed 상태 리턴 시 episode.status 복구."""
    mock_db = MagicMock()
    mock_episode = MagicMock()
    mock_episode.status = "analyzing"
    mock_db.query.return_value.filter.return_value.first.return_value = mock_episode

    mock_runner = MagicMock()
    mock_runner.run.return_value = {"status": "failed"}

    with patch("app.services.analysis_dispatch_service.SessionLocal", return_value=mock_db), \
         patch("app.services.analysis_dispatch_service.get_step_runner", return_value=mock_runner), \
         patch("app.services.analysis_dispatch_service.orchestrate_full_sync"):
        ds.run_steps_batch("p1", "e1", ["text_cleanup"], "resume", {}, {})

    assert mock_episode.status == "error"
    assert "text_cleanup" in (mock_episode.analysis_error or "")


def test_run_steps_batch_step_crashed_triggers_recovery():
    """step이 예외 던지면 episode.status 복구."""
    mock_db = MagicMock()
    mock_episode = MagicMock()
    mock_episode.status = "analyzing"
    mock_db.query.return_value.filter.return_value.first.return_value = mock_episode

    mock_runner = MagicMock()
    mock_runner.run.side_effect = RuntimeError("boom")

    with patch("app.services.analysis_dispatch_service.SessionLocal", return_value=mock_db), \
         patch("app.services.analysis_dispatch_service.get_step_runner", return_value=mock_runner), \
         patch("app.services.analysis_dispatch_service.orchestrate_full_sync"):
        ds.run_steps_batch("p1", "e1", ["text_cleanup"], "resume", {}, {})

    assert mock_episode.status == "error"
    assert "boom" in (mock_episode.analysis_error or "")


def test_run_steps_batch_gate_blocked_continues():
    """AppError 'gate.blocked'는 해당 step skip하고 다음 step 실행."""
    mock_db = MagicMock()
    mock_db.query.return_value.filter.return_value.first.return_value = None

    mock_runner_blocked = MagicMock()
    mock_runner_blocked.run.side_effect = AppError(code="gate.blocked", message="skip me", status_code=409)
    mock_runner_ok = MagicMock()
    mock_runner_ok.run.return_value = {"status": "done"}

    call_log = []

    def _factory(sid, *a, **kw):
        call_log.append(sid)
        return {"step_a": mock_runner_blocked, "step_b": mock_runner_ok}[sid]

    with patch("app.services.analysis_dispatch_service.SessionLocal", return_value=mock_db), \
         patch("app.services.analysis_dispatch_service.get_step_runner", side_effect=_factory), \
         patch("app.services.analysis_dispatch_service.orchestrate_full_sync") as mock_sync:
        ds.run_steps_batch("p1", "e1", ["step_a", "step_b"], "resume", {}, {})

    assert call_log == ["step_a", "step_b"]
    # all_ok=True 경로 → final orchestrate_full_sync 1회
    assert mock_sync.call_count == 1


def test_run_steps_batch_partial_strict_step_breaks_batch():
    """partial + manifest allow_partial_downstream=False → batch 중단 (problems.md #11)."""
    mock_db = MagicMock()
    mock_episode = MagicMock()
    mock_episode.status = "analyzing"
    mock_db.query.return_value.filter.return_value.first.return_value = mock_episode

    mock_runner_partial = MagicMock()
    mock_runner_partial.run.return_value = {"status": "partial"}
    mock_runner_next = MagicMock()
    mock_runner_next.run.return_value = {"status": "done"}

    call_log = []

    def _factory(sid, *a, **kw):
        call_log.append(sid)
        return {"strict_step": mock_runner_partial, "next_step": mock_runner_next}[sid]

    def _meta(sid):
        return {
            "strict_step": {"depends_on": [], "allow_partial_downstream": False},
            "next_step": {"depends_on": []},
        }.get(sid, {})

    with patch("app.services.analysis_dispatch_service.SessionLocal", return_value=mock_db), \
         patch("app.services.analysis_dispatch_service.get_step_runner", side_effect=_factory), \
         patch("app.services.analysis_dispatch_service.orchestrate_full_sync"), \
         patch("app.core.step_manifest.get_manifest_dict", side_effect=_meta):
        ds.run_steps_batch("p1", "e1", ["strict_step", "next_step"], "resume", {}, {})

    # strict_step partial → break, next_step 미실행
    assert call_log == ["strict_step"]
    assert mock_episode.status == "error"
    assert "strict_step" in (mock_episode.analysis_error or "")


def test_run_steps_batch_partial_default_continues_backward_compat():
    """partial + 기본 (allow_partial_downstream 미지정) → 이전 동작대로 cascade 진행."""
    mock_db = MagicMock()
    mock_db.query.return_value.filter.return_value.first.return_value = None

    mock_runner_partial = MagicMock()
    mock_runner_partial.run.return_value = {"status": "partial"}
    mock_runner_next = MagicMock()
    mock_runner_next.run.return_value = {"status": "done"}

    call_log = []

    def _factory(sid, *a, **kw):
        call_log.append(sid)
        return {"flexible_step": mock_runner_partial, "next_step": mock_runner_next}[sid]

    def _meta(sid):
        # 둘 다 allow_partial_downstream 미지정 → default True
        return {
            "flexible_step": {"depends_on": []},
            "next_step": {"depends_on": []},
        }.get(sid, {})

    with patch("app.services.analysis_dispatch_service.SessionLocal", return_value=mock_db), \
         patch("app.services.analysis_dispatch_service.get_step_runner", side_effect=_factory), \
         patch("app.services.analysis_dispatch_service.orchestrate_full_sync"), \
         patch("app.core.step_manifest.get_manifest_dict", side_effect=_meta):
        ds.run_steps_batch("p1", "e1", ["flexible_step", "next_step"], "resume", {}, {})

    # 둘 다 실행됨 (backward-compat)
    assert call_log == ["flexible_step", "next_step"]


def test_run_steps_batch_mid_sync_before_scene_director():
    """scene_director 는 진입 **직전과 직후** 모두 orchestrate_full_sync 호출.

    2026-08-04: post-sync 추가로 기대값 2 → 3. flagged step 이 실행 후에도
    동기화하지 않으면, 그 step 이 마지막 flagged step 일 때 산출이 DB 에
    영영 반영되지 않는다 (실측: `scene_detail` 255/255 completed 인데
    `scene_still` 0행 — pre-sync 시점엔 CP 가 실패 상태였고 그 뒤 트리거가
    소진됐다). 단일 step 경로는 이미 post-sync 를 한다.
    """
    mock_db = MagicMock()
    mock_db.query.return_value.filter.return_value.first.return_value = None

    mock_runner = MagicMock()
    mock_runner.run.return_value = {"status": "done"}

    with patch("app.services.analysis_dispatch_service.SessionLocal", return_value=mock_db), \
         patch("app.services.analysis_dispatch_service.get_step_runner", return_value=mock_runner), \
         patch("app.services.analysis_dispatch_service.orchestrate_full_sync") as mock_sync:
        ds.run_steps_batch(
            "p1", "e1", ["text_cleanup", "scene_director"], "resume", {}, {}
        )

    # pre-sync(scene_director 직전) 1 + post-sync(직후) 1 + final 1 = 3회
    assert mock_sync.call_count == 3
    # post-sync 는 어느 step 이 냈는지 남기려고 step_id 를 넘긴다.
    assert any(
        c.kwargs.get("step_id") == "scene_director"
        for c in mock_sync.call_args_list
    ), "post-sync 가 step_id=scene_director 로 호출되지 않았다"


def test_run_steps_batch_final_sync_failure_triggers_recovery():
    """모든 step done인데 final sync가 실패해도 status 복구."""
    mock_db = MagicMock()
    mock_episode = MagicMock()
    mock_episode.status = "analyzing"
    mock_db.query.return_value.filter.return_value.first.return_value = mock_episode

    mock_runner = MagicMock()
    mock_runner.run.return_value = {"status": "done"}

    with patch("app.services.analysis_dispatch_service.SessionLocal", return_value=mock_db), \
         patch("app.services.analysis_dispatch_service.get_step_runner", return_value=mock_runner), \
         patch(
             "app.services.analysis_dispatch_service.orchestrate_full_sync",
             side_effect=RuntimeError("sync broke"),
         ):
        ds.run_steps_batch("p1", "e1", ["text_cleanup"], "resume", {}, {})

    assert mock_episode.status == "error"
    assert "sync broke" in (mock_episode.analysis_error or "")


def test_run_steps_batch_blocked_cascades_to_dependents():
    """Phase 9.2 회귀 가드: blocked step의 downstream(depends_on에 포함된 step)이
    silent skip되지 않고 cascade로 자동 skip 되어야 한다.

    이전 동작: gate.blocked → continue (silent skip)
                 → downstream이 부재 ref/data로 비정상 결과 생성
    수정 동작: blocked_steps 추적 → 다음 step의 depends_on에 blocked가 있으면
              cascade-skip + warning 로그
    """
    mock_db = MagicMock()
    mock_db.query.return_value.filter.return_value.first.return_value = None

    mock_runner_blocked = MagicMock()
    mock_runner_blocked.run.side_effect = AppError(
        code="gate.blocked", message="prereq missing", status_code=409
    )
    mock_runner_dependent = MagicMock()
    mock_runner_dependent.run.return_value = {"status": "done"}
    mock_runner_unrelated = MagicMock()
    mock_runner_unrelated.run.return_value = {"status": "done"}

    runners = {
        "step_a": mock_runner_blocked,
        "step_b": mock_runner_dependent,  # depends on step_a → cascade
        "step_c": mock_runner_unrelated,  # 독립 → 정상 실행
    }

    def _factory(sid, *a, **kw):
        return runners[sid]

    def _manifest(sid):
        manifests = {
            "step_a": {"depends_on": []},
            "step_b": {"depends_on": ["step_a"]},
            "step_c": {"depends_on": []},
        }
        return manifests.get(sid, {})

    with patch("app.services.analysis_dispatch_service.SessionLocal", return_value=mock_db), \
         patch("app.services.analysis_dispatch_service.get_step_runner", side_effect=_factory), \
         patch("app.services.analysis_dispatch_service.orchestrate_full_sync"), \
         patch("app.core.step_manifest.get_manifest_dict", side_effect=_manifest):
        ds.run_steps_batch("p1", "e1", ["step_a", "step_b", "step_c"], "resume", {}, {})

    # step_a: blocked → run 호출 시 AppError → blocked_steps에 추가
    assert mock_runner_blocked.run.called
    # step_b: cascade skip — runner.run 호출 없음
    assert not mock_runner_dependent.run.called
    # step_c: 의존 없음 → 정상 실행
    assert mock_runner_unrelated.run.called


def test_run_steps_batch_blocked_cascade_chain():
    """blocked cascade가 transitive하게 전파되는지: A blocked → B(A 의존) skip
    → C(B 의존) skip.
    """
    mock_db = MagicMock()
    mock_db.query.return_value.filter.return_value.first.return_value = None

    mock_a = MagicMock()
    mock_a.run.side_effect = AppError(code="gate.blocked", message="root", status_code=409)
    mock_b = MagicMock()
    mock_b.run.return_value = {"status": "done"}
    mock_c = MagicMock()
    mock_c.run.return_value = {"status": "done"}

    runners = {"step_a": mock_a, "step_b": mock_b, "step_c": mock_c}

    def _factory(sid, *a, **kw):
        return runners[sid]

    def _manifest(sid):
        return {
            "step_a": {"depends_on": []},
            "step_b": {"depends_on": ["step_a"]},
            "step_c": {"depends_on": ["step_b"]},
        }.get(sid, {})

    with patch("app.services.analysis_dispatch_service.SessionLocal", return_value=mock_db), \
         patch("app.services.analysis_dispatch_service.get_step_runner", side_effect=_factory), \
         patch("app.services.analysis_dispatch_service.orchestrate_full_sync"), \
         patch("app.core.step_manifest.get_manifest_dict", side_effect=_manifest):
        ds.run_steps_batch("p1", "e1", ["step_a", "step_b", "step_c"], "resume", {}, {})

    assert mock_a.run.called  # 시도됨 → blocked
    assert not mock_b.run.called  # cascade
    assert not mock_c.run.called  # cascade transitive


def test_select_scene_reanalysis_starts_at_scene_segmentation():
    """Codex Review Important #2 — reanalysis는 scene_segmentation부터 시작."""
    db = MagicMock()
    db.query.return_value.filter.return_value.first.return_value = None

    step_ids = ds.select_scene_reanalysis_steps(db, "p1")

    assert "scene_segmentation" in step_ids
    assert "scene_save" in step_ids
    # scene_segmentation이 scene_save보다 먼저
    assert step_ids.index("scene_segmentation") < step_ids.index("scene_save")


# ── get_step_runner ───────────────────────────────────────────


def test_get_step_runner_returns_instance_for_valid_step_id():
    db = MagicMock()
    runner = ds.get_step_runner("text_cleanup", "p1", "e1", db, {})
    assert runner is not None
    assert runner.step_id == "text_cleanup"


def test_get_step_runner_raises_for_unknown_step_id():
    db = MagicMock()
    with pytest.raises(AppError) as exc_info:
        ds.get_step_runner("nonexistent_step_xyz", "p1", "e1", db, {})
    assert exc_info.value.code == "step.not_found"


# ── R3 cascade fix (2026-05-05): analysis cross-category fail-fast + preflight auto-recover ──


def test_select_steps_for_category_analysis_cross_category_deps_fail_fast():
    """analysis 카테고리 force 시 image category prerequisite 미완료면 즉시 거절.

    Phase 7 도입 후 analysis step (background_prompt, scene_detail 등) 이 image step
    (floor_plan_render, background_render) 에 depend. analysis 단독 force 시 silent
    cascade skip + episode.status='analyzing' stuck 사고 방지 (R3 cascade fix).
    Error code 는 image 분기와 동일한 `dispatch.deps_incomplete` (review M1).
    """
    db = MagicMock()
    db.query.return_value.filter.return_value.first.return_value = None

    with patch.object(ds, "_is_step_satisfied", return_value=False):
        with pytest.raises(AppError) as exc_info:
            ds.select_steps_for_category(db, "p1", "analysis", episode_id="e1")

    assert exc_info.value.code == "dispatch.deps_incomplete"
    msg = str(exc_info.value.message)
    # 가이드 문구 포함 — 사용자에게 category='all' 안내
    assert "category='all'" in msg or "category=all" in msg.replace("'", "")


def test_select_steps_for_category_analysis_episode_none_skips_cross_validation():
    """episode_id=None (preview/list 호출) 일 때는 cross-category 검증 skip.

    image 분기 prerequisite 검증과 동일한 backward-compat 정책.
    """
    db = MagicMock()
    db.query.return_value.filter.return_value.first.return_value = None

    with patch.object(ds, "_is_step_satisfied", return_value=False):
        # episode_id 누락 → 검증 skip → 정상 반환
        step_ids = ds.select_steps_for_category(db, "p1", "analysis")

    assert len(step_ids) > 0


def test_select_steps_for_category_analysis_cross_deps_satisfied_passes():
    """cross-category prerequisite 모두 satisfied 면 정상 반환."""
    db = MagicMock()
    db.query.return_value.filter.return_value.first.return_value = None

    with patch.object(ds, "_is_step_satisfied", return_value=True):
        step_ids = ds.select_steps_for_category(db, "p1", "analysis", episode_id="e1")

    assert len(step_ids) > 0
    from app.core.step_catalog import STEP_CATALOG
    for sid in step_ids:
        assert STEP_CATALOG[sid].category == "analysis"


def test_preflight_analysis_start_stale_lock_auto_recovers(caplog):
    """status='analyzing' + task_registry 에 active job 없으면 auto-recover (stale lock).

    Review I1: stale recovery 시 prior_status 가 'error' 로 normalize 되어, 후속 dispatch
    가 fail 하더라도 rollback 시 다시 'analyzing' stuck 으로 복원되지 않아야 함.
    Review I2: 4 RUN_ALL_CATEGORY_TOKENS 모두 검사 + warning log 검증.
    """
    db = MagicMock()
    mock_episode = MagicMock()
    mock_episode.id = "e1"
    mock_episode.status = "analyzing"
    mock_episode.fulltext = "some scenario text"
    db.query.return_value.filter.return_value.with_for_update.return_value.first.return_value = mock_episode

    with patch("app.core.task_registry.is_task_running", return_value=False) as mock_is_running, \
         patch("app.core.config.settings") as mock_settings, \
         caplog.at_level("WARNING", logger="app.services.analysis_dispatch_service"):
        mock_settings.openai_api_key = "test-key"
        prior = ds.preflight_analysis_start(db, "p1", "e1")

    # auto-recover → 다시 'analyzing' 으로 set + commit (race window 종료)
    assert mock_episode.status == "analyzing"
    # I1: stale recovery 시 prior_status 가 'error' 로 normalize (rollback 안전)
    assert prior == "error"
    # 4 RUN_ALL_CATEGORY_TOKENS 모두 검사 (analysis/image/all/reanalyze)
    assert mock_is_running.call_count == 4
    db.commit.assert_called()
    # warning log 발생
    assert any("auto-recovering stale lock" in rec.message for rec in caplog.records)


def test_preflight_analysis_start_active_job_blocks():
    """status='analyzing' + task_registry 에 active job 있으면 409 거절 (정상 lock)."""
    db = MagicMock()
    mock_episode = MagicMock()
    mock_episode.status = "analyzing"
    db.query.return_value.filter.return_value.with_for_update.return_value.first.return_value = mock_episode

    with patch("app.core.task_registry.is_task_running", return_value=True):
        with pytest.raises(AppError) as exc_info:
            ds.preflight_analysis_start(db, "p1", "e1")

    assert exc_info.value.code == "analysis.already_running"
    assert exc_info.value.status_code == 409


def test_preflight_analysis_start_normal_pending_status_unchanged_prior():
    """status='pending' (정상) 시 prior_status='pending' 보존 — stale recovery 와 분리."""
    db = MagicMock()
    mock_episode = MagicMock()
    mock_episode.id = "e1"
    mock_episode.status = "pending"
    mock_episode.fulltext = "some scenario text"
    db.query.return_value.filter.return_value.with_for_update.return_value.first.return_value = mock_episode

    with patch("app.core.config.settings") as mock_settings:
        mock_settings.openai_api_key = "test-key"
        prior = ds.preflight_analysis_start(db, "p1", "e1")

    # 정상 path: prior_status 그대로 ('pending')
    assert prior == "pending"
    assert mock_episode.status == "analyzing"


# ── Fix D (cascade fix 2026-05-05): image+force 시 cross-category downstream 검증 ──


def test_select_steps_for_category_image_force_cross_category_downstream_blocks():
    """image+force 시 force-invalidate 가 cross-category(analysis) downstream 무효화 → 거절.

    Phase 7 cascade: floor_plan_render(image) → background_prompt(analysis) →
    scene_detail(analysis) ... force chain. mode=force 로 image batch 시작 시
    background_prompt 같은 analysis step 의 manifest 가 cascade-delete 됨 →
    image batch 가 못 채움 → background_render 시점 fail. fix D 가 시작 시점에 거절.
    """
    db = MagicMock()
    db.query.return_value.filter.return_value.first.return_value = None

    # image deps 모두 satisfied 가정 — fix D 가 force-cascade 만 catch
    with patch.object(ds, "_is_step_satisfied", return_value=True):
        with pytest.raises(AppError) as exc_info:
            ds.select_steps_for_category(
                db, "p1", "image", episode_id="e1", mode="force"
            )

    assert exc_info.value.code == "dispatch.force_cascade_cross_category"
    msg = str(exc_info.value.message)
    # 가이드 문구 — category=all&mode=force 또는 category=image&mode=resume
    assert "category=all" in msg.replace("'", "") or "mode=resume" in msg.replace("'", "")


def test_select_steps_for_category_image_resume_skips_force_cascade_check():
    """mode=resume 시 force-cascade 검증 skip — invalidate cascade 일어나지 않음.

    review M3: vacuous pass 회피 — floor_plan_render 가 step_ids 에 포함됐고
    그 step 의 transitive downstream 이 cross-category 를 포함한다는 production
    조건 명시 검증 (manifest refactor 시 fix D 가 silent regress 안 하도록).
    """
    db = MagicMock()
    db.query.return_value.filter.return_value.first.return_value = None

    with patch.object(ds, "_is_step_satisfied", return_value=True):
        # resume 모드 → fix D 검증 skip → 정상 반환
        step_ids = ds.select_steps_for_category(
            db, "p1", "image", episode_id="e1", mode="resume"
        )

    assert len(step_ids) > 0
    from app.core.step_catalog import STEP_CATALOG
    for sid in step_ids:
        assert STEP_CATALOG[sid].category == "image"
    # M3: floor_plan_render 가 batch 에 존재 + 그 transitive downstream 에
    # cross-category step 존재 → fix D 가 force 시 거절할 조건 충족.
    assert "floor_plan_render" in step_ids
    from app.core.step_catalog import get_all_downstream_recursive
    fpr_downstream = set(get_all_downstream_recursive("floor_plan_render"))
    cross_in_downstream = [
        d for d in fpr_downstream
        if d in STEP_CATALOG and STEP_CATALOG[d].category != "image"
    ]
    assert len(cross_in_downstream) > 0, (
        "floor_plan_render transitive downstream 에 cross-category step 이 "
        "있어야 fix D test 가 의미 있음"
    )


def test_select_steps_for_category_image_force_episode_none_skips_force_cascade():
    """episode_id=None 일 때는 force-cascade 검증 skip (preview/list 호출 backward-compat)."""
    db = MagicMock()
    db.query.return_value.filter.return_value.first.return_value = None

    with patch.object(ds, "_is_step_satisfied", return_value=True):
        # episode_id 누락 → 검증 skip → 정상 반환
        step_ids = ds.select_steps_for_category(
            db, "p1", "image", mode="force"
        )

    assert len(step_ids) > 0


def test_dispatch_category_run_image_force_fails_with_cascade_guidance():
    """dispatch_category_run(category='image', mode='force') 시 fix D 가 거절 + 가이드.

    W20E5 budget gate 가 cascade 검사보다 먼저 평가되므로(fail-closed by design),
    approve+cap 을 전달해 gate 를 통과시킨다 — 이 테스트의 대상은 Fix D cascade.
    (미전달 시 background_render_reference_mode=shot_aware_plan 환경(.env)에서
    step.image_generation_not_approved 가 먼저 발사돼 환경 의존 실패 — 실측.)
    """
    db = MagicMock()
    db.execute.return_value.fetchone.return_value = ("proj_name", "ep_title")
    db.query.return_value.filter.return_value.first.return_value = None

    with patch("app.core.task_registry.is_task_running", return_value=False), \
         patch.object(ds, "_is_step_satisfied", return_value=True):
        with pytest.raises(AppError) as exc_info:
            ds.dispatch_category_run(
                "p1", "e1", "image", "force", db,
                image_call_cap=10, approve_image_generation=True,
            )

    assert exc_info.value.code == "dispatch.force_cascade_cross_category"


def test_select_steps_for_category_analysis_force_cross_category_downstream_blocks():
    """analysis+force 도 같은 cascade risk — Fix D iter1 (Codex BLOCKING / Claude I2).

    analysis step 의 transitive downstream 이 image step (scene_image_pipeline
    등) 을 포함하면 force-invalidate 가 image manifest 무효화 → image batch
    later 시점 stale state. 대칭 검증 (image 분기와 동일).
    """
    db = MagicMock()
    db.query.return_value.filter.return_value.first.return_value = None

    with patch.object(ds, "_is_step_satisfied", return_value=True):
        with pytest.raises(AppError) as exc_info:
            ds.select_steps_for_category(
                db, "p1", "analysis", episode_id="e1", mode="force"
            )

    assert exc_info.value.code == "dispatch.force_cascade_cross_category"


def test_select_steps_for_category_invalid_mode_raises():
    """mode enum strict — typo / None 등 silent skip 차단 (Claude I1)."""
    db = MagicMock()
    db.query.return_value.filter.return_value.first.return_value = None

    for bad_mode in ("FORCE", "Force", "frorce", None, ""):
        with pytest.raises(AppError) as exc_info:
            ds.select_steps_for_category(
                db, "p1", "image", episode_id="e1", mode=bad_mode
            )
        assert exc_info.value.code == "dispatch.invalid_mode", (
            f"bad mode {bad_mode!r} 거절 안됨"
        )


# ── 2026-08-26 Codex 2차 재리뷰 BLOCK-4 ──────────────────────────────


def test_되돌리기가_실패해도_원래_예외가_올라온다():
    """★정리하다 난 예외가 **진짜 실패 사유를 가리면 안 된다.**

    종전에는 `rollback_episode_status` 가 실패하면 그 예외가 그대로 올라가
    맨 아래 bare raise 에 도달하지 못했다. 호출자는 「왜 실패했나」 대신
    「되돌리다 실패했다」만 봤다.
    """
    from app.core.task_registry import episode_run_holder

    db = MagicMock()
    db.execute.return_value.fetchone.return_value = ("proj", "ep")
    db.query.return_value.filter.return_value.first.return_value = None

    with patch("app.core.task_registry.is_task_running", return_value=False), \
         patch.object(ds, "preflight_analysis_start", return_value="idle"), \
         patch.object(ds, "select_steps_for_category",
                      side_effect=RuntimeError("진짜 실패 사유")), \
         patch.object(ds, "rollback_episode_status",
                      side_effect=RuntimeError("되돌리기도 실패")):
        with pytest.raises(RuntimeError) as ei:
            ds.dispatch_category_run("p1", "e1", "analysis", "resume", db)

    assert str(ei.value) == "진짜 실패 사유", (
        "정리하다 난 예외가 원래 사유를 가렸다")
    assert episode_run_holder("run_all:p1:e1") is None, (
        "되돌리기가 실패했다고 에피소드 자리까지 붙들고 있으면 그 에피소드는 "
        "영영 막힌다")


def test_세션이_죽었으면_새_세션으로_되돌린다():
    """`db.rollback()` 마저 실패하면 그 세션으로는 아무것도 못 한다."""
    from app.core.task_registry import episode_run_holder

    db = MagicMock()
    db.execute.return_value.fetchone.return_value = ("proj", "ep")
    db.query.return_value.filter.return_value.first.return_value = None
    db.rollback.side_effect = RuntimeError("세션이 죽었다")

    with patch("app.core.task_registry.is_task_running", return_value=False), \
         patch.object(ds, "preflight_analysis_start", return_value="idle"), \
         patch.object(ds, "select_steps_for_category",
                      side_effect=RuntimeError("진짜 실패 사유")), \
         patch.object(ds, "_rollback_episode_status_fresh") as fresh, \
         patch.object(ds, "rollback_episode_status") as same_session:
        with pytest.raises(RuntimeError) as ei:
            ds.dispatch_category_run("p1", "e1", "analysis", "resume", db)

    assert str(ei.value) == "진짜 실패 사유"
    fresh.assert_called_once_with("e1", "idle")
    same_session.assert_not_called()   # 죽은 세션을 그대로 다시 쓰면 안 된다
    assert episode_run_holder("run_all:p1:e1") is None
