"""에피소드 관리 API 라우터."""

import asyncio
import json
import logging

from fastapi import APIRouter, Depends, Request, UploadFile, File, Form
from fastapi.responses import StreamingResponse
from sqlalchemy.orm import Session as OrmSession

from app.api.deps import api_endpoint, get_db, get_current_user, verify_project_access
from app.core.config import settings
from app.core.database import SessionLocal
from app.core.errors import AppError
from app.i18n.loader import t
from app.models.catalog import UserAccount
from app.models.project import Episode, PipelineProgress
from app.schemas.episode import EpisodeCreate, EpisodeUpdate, EpisodeResponse
from app.services.episode_service import EpisodeService

logger = logging.getLogger(__name__)

router = APIRouter(
    prefix="/api/v1/projects/{project_id}/episodes",
    tags=["episodes"],
)

# Phase 0 deprecation notice shared by /analyze, /reanalyze-scenes
_ANALYSIS_SERVICE_DEPRECATION_MSG = (
    "DEPRECATED endpoint %s called (episode=%s). "
    "AnalysisService 경로는 architecture refactor Phase 4에서 제거됩니다. "
    "공식 실행 경로는 StepRunner "
    "(POST /api/v1/projects/{project_id}/episodes/{episode_id}/steps/run-all) 입니다. "
    "참조: docs/architecture-refactor-final/02-final-roadmap.md §Phase 4"
)


def _make_service(
    project_id: str,
    db: OrmSession,
    user: UserAccount,
) -> EpisodeService:
    return EpisodeService(
        db=db,
        project_id=project_id,
        actor_id=user.id,
    )


@router.get("/", response_model=list[EpisodeResponse])
def list_episodes(
    project_id: str = Depends(verify_project_access),
    db: OrmSession = Depends(get_db),
    current_user: UserAccount = Depends(get_current_user),
):
    service = _make_service(project_id, db, current_user)
    return service.list_episodes()


@router.post("/", response_model=EpisodeResponse)
async def create_episode(
    episode_number: int = Form(...),
    title: str = Form(...),
    file: UploadFile = File(...),
    request: Request = None,
    project_id: str = Depends(verify_project_access),
    db: OrmSession = Depends(get_db),
    current_user: UserAccount = Depends(get_current_user),
):
    ip = request.client.host if request and request.client else None
    pdf_bytes = await file.read()
    filename = file.filename or "screenplay.pdf"
    service = _make_service(project_id, db, current_user)
    return service.create_episode(episode_number, title, pdf_bytes, filename, ip)


@router.get("/{episode_id}", response_model=EpisodeResponse)
def get_episode(
    episode_id: str,
    project_id: str = Depends(verify_project_access),
    db: OrmSession = Depends(get_db),
    current_user: UserAccount = Depends(get_current_user),
):
    service = _make_service(project_id, db, current_user)
    return service.get_episode(episode_id)


@router.get("/{episode_id}/fulltext")
def get_episode_fulltext(
    episode_id: str,
    project_id: str = Depends(verify_project_access),
    db: OrmSession = Depends(get_db),
    current_user: UserAccount = Depends(get_current_user),
):
    """에피소드 시나리오 전문 반환 (씬 세그먼트 표시용)."""
    from app.core.config import settings
    episode = db.query(Episode).filter(Episode.id == episode_id, Episode.project_id == project_id).first()
    if not episode:
        raise AppError(code="episode.not_found", message=t("episode.not_found"), status_code=404)
    return {
        "fulltext": episode.fulltext or "",
        "segment_context_chars": settings.scene_segment_context_chars,
    }


@router.patch("/{episode_id}", response_model=EpisodeResponse)
def update_episode(
    episode_id: str,
    body: EpisodeUpdate,
    request: Request,
    project_id: str = Depends(verify_project_access),
    db: OrmSession = Depends(get_db),
    current_user: UserAccount = Depends(get_current_user),
):
    ip = request.client.host if request.client else None
    updates = body.model_dump(exclude_none=True)
    service = _make_service(project_id, db, current_user)
    return service.update_episode(episode_id, updates, ip)


@router.delete("/{episode_id}")
def delete_episode(
    episode_id: str,
    request: Request,
    project_id: str = Depends(verify_project_access),
    db: OrmSession = Depends(get_db),
    current_user: UserAccount = Depends(get_current_user),
):
    ip = request.client.host if request.client else None
    service = _make_service(project_id, db, current_user)
    service.delete_episode(episode_id, ip)
    return {"ok": True}


@router.post("/{episode_id}/analyze", deprecated=True)
def analyze_episode(
    episode_id: str,
    request: Request,
    project_id: str = Depends(verify_project_access),
    db: OrmSession = Depends(get_db),
    current_user: UserAccount = Depends(get_current_user),
):
    """Start LLM analysis for an episode (background thread).

    DEPRECATED (Phase 4에서 StepRunner 경로로 내부 교체됨):
        공식 경로는 `/api/v1/projects/{project_id}/episodes/{episode_id}/steps/run-all?category=analysis`.
        본 엔드포인트는 dispatch_category_run(category="analysis")으로 내부 위임되며,
        호환성 유지 목적으로 한시 보존됨. 차기 릴리스에서 제거 예정.
        참조: docs/architecture-refactor-final/02-final-roadmap.md §Phase 4
    """
    from app.services.analysis_dispatch_service import dispatch_category_run

    logger.warning(_ANALYSIS_SERVICE_DEPRECATION_MSG, "/analyze", episode_id)

    # W1-F08: preflight(row lock, status 전이, fulltext/openai_key 검증, 실패 rollback)는
    # dispatch_category_run이 category='analysis'일 때 내부에서 처리.
    result = dispatch_category_run(
        project_id=project_id,
        episode_id=episode_id,
        category="analysis",
        mode="resume",
        db=db,
    )
    return {"ok": True, "message": t("analysis.started"), **result}


@router.post("/{episode_id}/reanalyze-scenes", deprecated=True)
def reanalyze_scenes(
    episode_id: str,
    request: Request,
    project_id: str = Depends(verify_project_access),
    db: OrmSession = Depends(get_db),
    current_user: UserAccount = Depends(get_current_user),
):
    """씬만 재분석 (요소 유지).

    DEPRECATED (Phase 4에서 StepRunner 경로로 내부 교체됨):
        scene_save + downstream(category=analysis, active) force 실행.
        공식 경로는 `/steps/{step_id}?mode=force` 개별 호출.
        참조: docs/architecture-refactor-final/02-final-roadmap.md §Phase 4
    """
    from app.services.analysis_dispatch_service import dispatch_scene_reanalysis

    logger.warning(_ANALYSIS_SERVICE_DEPRECATION_MSG, "/reanalyze-scenes", episode_id)

    # Lock the row to prevent race condition between concurrent requests
    episode = (
        db.query(Episode)
        .filter(Episode.id == episode_id, Episode.project_id == project_id)
        .with_for_update()
        .first()
    )
    if not episode:
        raise AppError(code="episode.not_found", message=t("episode.not_found"), status_code=404)
    if episode.status == "analyzing":
        raise AppError(code="analysis.already_running", message=t("analysis.already_running"), status_code=409)
    if not episode.fulltext:
        raise AppError(code="analysis.no_text", message=t("analysis.no_text"), status_code=400)

    # Set status to "analyzing" immediately to close the race window
    prior_status = episode.status
    episode.status = "analyzing"
    episode.analysis_error = None
    db.commit()

    # dispatch 실패 시 status 원복 (Codex Review Important #1)
    try:
        result = dispatch_scene_reanalysis(
            project_id=project_id,
            episode_id=episode_id,
            db=db,
        )
    except Exception:
        episode.status = prior_status
        db.commit()
        raise
    return {"ok": True, "message": "씬 재분석이 시작되었습니다", **result}


@router.get("/{episode_id}/segment-preview")
def segment_preview(
    episode_id: str,
    threshold: int = 600,
    project_id: str = Depends(verify_project_access),
    db: OrmSession = Depends(get_db),
    current_user: UserAccount = Depends(get_current_user),
):
    """정규식 세그먼테이션 미리보기 — threshold별 예상 씬 수."""
    episode = db.query(Episode).filter(Episode.id == episode_id, Episode.project_id == project_id).first()
    if not episode or not episode.fulltext:
        return {"base_scenes": 0, "split_candidates": 0, "estimated_total": 0, "scenes": []}

    from app.modules.pipeline.scene_extractor_v2 import _segment_scenes
    segments = _segment_scenes(episode.fulltext, split_threshold=threshold)
    long_scenes = [s for s in segments if s["length"] >= threshold]

    return {
        "base_scenes": len(segments),
        "split_candidates": len(long_scenes),
        "estimated_total": len(segments),
        "threshold": threshold,
        "scenes": [
            {
                "index": s["scene_index"],
                "heading": s["heading"],
                "length": s["length"],
                "will_split": s["length"] >= threshold,
            }
            for s in segments
        ],
    }


_OPERATION_KEYS = ("analysis", "reference_image_generation", "image_generation", "webbook", "pdf_render")


@router.post("/{episode_id}/repair-projection")
@api_endpoint
def repair_projection(
    episode_id: str,
    project_id: str = Depends(verify_project_access),
    db: OrmSession = Depends(get_db),
    current_user: UserAccount = Depends(get_current_user),
):
    """W4 P3-3 — sync_status='failed' step_run의 DB projection만 재시도.

    checkpoint는 canonical이므로 건드리지 않음. LLM 재호출 없음.
    5-way sync 실행 + 같은 트랜잭션에서 failed→synced 전환 (Codex Medium: split commit 회피).
    sync_error가 우리가 캡처한 값과 일치할 때만 전환 — concurrent 프로세스의 새 실패는
    덮어쓰지 않음 (Codex High).

    Returns:
        {
            "ok": true,
            "sync_result": {entity/relation/scene_still/outlook/episode 결과},
            "repaired_steps": [repair 시도한 step_id 리스트]
        }

    예외:
        - 404 episode.not_found
        - sync 실패 시 AppError → @api_endpoint가 내부 오류로 정규화(500)
    """
    from sqlalchemy import text as sql_text
    from app.services.checkpoint_sync import orchestrate_full_sync

    # Episode existence check (권한은 verify_project_access에서 수행됨)
    ep = (
        db.query(Episode)
        .filter(Episode.id == episode_id, Episode.project_id == project_id)
        .first()
    )
    if not ep:
        raise AppError(
            code="episode.not_found",
            message=t("episode.not_found"),
            status_code=404,
        )

    # 현재 failed 상태 step 수집 — 각 step의 sync_error도 캡처 (guarded UPDATE용).
    failed_rows = db.execute(sql_text(
        "SELECT step_id, sync_error FROM step_run "
        "WHERE project_id = :pid AND episode_id = :eid AND sync_status = 'failed'"
    ), {"pid": project_id, "eid": episode_id}).fetchall()
    repair_map: dict[str, str | None] = {r[0]: r[1] for r in failed_rows}

    # 5-way sync + failed→synced 전환을 단일 트랜잭션에서 수행.
    # 실패 시 orchestrator가 raise → @api_endpoint가 AppError(internal_error, 500)로 정규화.
    try:
        result = orchestrate_full_sync(
            project_id, episode_id, db, repair_step_ids=repair_map or None,
        )
    except Exception as exc:
        logger.error(
            "repair-projection failed for episode %s: %s", episode_id[:8], exc
        )
        db.rollback()
        raise AppError(
            code="repair.sync_failed",
            message=t("repair.sync_failed"),
            status_code=500,
        )

    return {
        "ok": True,
        "sync_result": result,
        "repaired_steps": list(repair_map.keys()),
    }


@router.get("/{episode_id}/progress")
def get_progress(
    episode_id: str,
    project_id: str = Depends(verify_project_access),
    db: OrmSession = Depends(get_db),
    current_user: UserAccount = Depends(get_current_user),
):
    """Get pipeline progress for an episode."""
    all_progress = (
        db.query(PipelineProgress)
        .filter(
            PipelineProgress.project_id == project_id,
            PipelineProgress.episode_id == episode_id,
        )
        .all()
    )

    def _format_progress(row):
        return {
            "status": row.status,
            "current_step": row.current_step or "",
            "completed_steps": row.completed_steps or 0,
            "total_steps": row.total_steps or 0,
            "error_message": row.error_message,
            "started_at": row.started_at,
            "updated_at": row.updated_at,
            "completed_at": row.completed_at,
        }

    result = {}
    for key in _OPERATION_KEYS:
        # Pick the most recently started row for each operation
        matches = [p for p in all_progress if p.operation == key]
        if matches:
            best = max(matches, key=lambda p: p.started_at or p.updated_at)
            result[key] = _format_progress(best)
        else:
            result[key] = None
    return result


@router.get("/{episode_id}/progress/stream")
async def stream_progress(
    episode_id: str,
    project_id: str = Depends(verify_project_access),
    current_user: UserAccount = Depends(get_current_user),
):
    """SSE endpoint for real-time progress updates."""

    async def event_stream():
        db = SessionLocal()
        try:
            while True:
                progress_rows = (
                    db.query(PipelineProgress)
                    .filter(
                        PipelineProgress.project_id == project_id,
                        PipelineProgress.episode_id == episode_id,
                    )
                    .all()
                )

                data = {}
                for key in _OPERATION_KEYS:
                    matches = [p for p in progress_rows if p.operation == key]
                    if matches:
                        best = max(matches, key=lambda p: p.started_at or p.updated_at)
                        data[key] = {
                            "status": best.status,
                            "current_step": best.current_step or "",
                            "completed_steps": best.completed_steps or 0,
                            "total_steps": best.total_steps or 0,
                        }

                yield f"data: {json.dumps(data, ensure_ascii=False)}\n\n"

                if (
                    progress_rows
                    and all(p.status in ("completed", "error") for p in progress_rows)
                ):
                    yield 'data: {"done": true}\n\n'
                    break

                db.expire_all()
                await asyncio.sleep(2)
        finally:
            db.close()

    return StreamingResponse(
        event_stream(),
        media_type="text/event-stream",
        headers={"Cache-Control": "no-cache", "X-Accel-Buffering": "no"},
    )
