"""기획서 project-level 분석 service.

기획서는 project-level 자원이므로 분석 결과도 episode 와 무관하게
``projects/{project_id}/checkpoints/planning_doc_analysis/manifest.json`` 에 저장.

기존 ``PlanningDocAnalysisStep`` (StepRunner) 은 episode 단위 cascade 호환을
위해 유지하되, 본 service 가 만든 project-level checkpoint 가 있으면 그대로
mirror (LLM 재호출 없음).

업로드 직후 백그라운드 dispatch 로 사용:

    from app.services.planning_doc_analysis_service import (
        dispatch_planning_doc_analysis_async,
    )
    dispatch_planning_doc_analysis_async(project_id)

## race 안전 — per-project lock + multi-layer guard

핵심 race: 업로드 → 분석 시작 → 사용자가 즉시 delete + 재업로드 시, 이전
background thread (A) 가 늦게 끝나서 새 SOT (B) 를 덮어쓰는 시나리오.
status write / manifest promote / clear / dispatch queued write 가 모두
"check-then-act" 라 단순 atomic write 만으로는 race window 가 남음:

    A: is_current_job(A) → True (status=A)
    B: clear → dispatch → status=B
    A: atomic_write status={state:running, job:A}  ← B 의 도장 덮어씀

근본 해법은 critical section 직렬화. ``_project_lock(project_id)`` 가
fcntl.flock + threading.Lock combo 로 same/cross-process 모두 보호.
LLM call 은 길어서 lock 밖 — write/check/cleanup 만 lock 안.

추가 안전망 (lock 으로 닫지 못한 edge — 옛 schema, race-free 호출 등):

  L1. ``_write_status``/``_write_manifest`` 의 ``expected_job_id`` 가드 —
      lock 안에서 호출되면 redundant 한 double check (안전 차원).
  L2. manifest payload 에 ``job_id`` + ``source_hash`` 동봉 → reader 검증.
  L3. promote = staging file (``manifest.{job_id}.json``) 에 atomic write
      후 lock 안에서 ``os.replace`` — cross-job manifest.json 보호.
  L4. ``load_project_checkpoint`` 가 manifest.source_hash 와 현재 disk
      source hash 비교 → 다르면 reject (cleanup unlink 제거 — reader 의
      늦은 unlink 가 fresh manifest 를 죽이는 race 차단). 다음 dispatch
      promote 또는 clear 가 stale 자연 정리.
"""

from __future__ import annotations

import hashlib
import json
import logging
import os
import platform
import threading
import uuid
from contextlib import contextmanager
from datetime import datetime, timezone
from pathlib import Path
from typing import Dict, Iterator, Optional

from app.core.config import settings

logger = logging.getLogger(__name__)


_STATUS_FILE = "status.json"          # {state, started_at, finished_at, error, job_id, source_hash}
_MANIFEST_FILE = "manifest.json"       # {data, job_id, source_hash, generated_at}
_LOCK_FILE = ".lock"

# fcntl 은 POSIX (darwin/linux) 만. Windows 는 threading.Lock 단독 (single
# process 가정 — production deploy target 이 POSIX 라 우선순위 낮음).
_IS_POSIX = platform.system() != "Windows"
if _IS_POSIX:
    import fcntl


# per-project threading.Lock 캐시 — 같은 process 내 thread 직렬화.
# fcntl.flock 은 cross-process protect 하지만 BSD-style 은 same-process
# different-fd 도 block 해주므로 이중 안전. _thread_locks 는 single-process
# fast path (file open/close 비용 회피).
_thread_locks: Dict[str, threading.Lock] = {}
_thread_locks_guard = threading.Lock()


def _get_thread_lock(project_id: str) -> threading.Lock:
    with _thread_locks_guard:
        lk = _thread_locks.get(project_id)
        if lk is None:
            lk = threading.Lock()
            _thread_locks[project_id] = lk
        return lk


@contextmanager
def _project_lock(project_id: str) -> Iterator[None]:
    """per-project critical section — same/cross process 모두 직렬화.

    threading.Lock + fcntl.flock combo:
      - threading.Lock: 같은 process 내 thread 들 빠르게 직렬화.
      - fcntl.flock (POSIX): uvicorn workers / --reload 일시 중첩 등
        cross-process race 방어 (file descriptor advisory lock).

    LLM call 같은 long-running operation 은 lock 밖에서 호출 — write /
    check / cleanup 만 critical section. 분석 자체가 lock 잡으면 사용자
    재업로드 시 새 dispatch 가 이전 분석 끝날 때까지 막히는 UX 손해.
    """
    tlock = _get_thread_lock(project_id)
    tlock.acquire()
    f = None
    try:
        if _IS_POSIX:
            d = get_project_checkpoint_dir(project_id)
            d.mkdir(parents=True, exist_ok=True)
            lock_path = d / _LOCK_FILE
            f = open(lock_path, "w")
            fcntl.flock(f.fileno(), fcntl.LOCK_EX)
        yield
    finally:
        if f is not None:
            try:
                fcntl.flock(f.fileno(), fcntl.LOCK_UN)
            except OSError as exc:
                logger.warning(
                    "_project_lock fcntl unlock failed (project=%s): %s",
                    project_id[:8], exc,
                )
            try:
                f.close()
            except OSError:
                pass
        tlock.release()


def get_project_checkpoint_dir(project_id: str) -> Path:
    """project-level planning_doc_analysis checkpoint 디렉토리."""
    return (
        Path(settings.projects_dir)
        / project_id
        / "checkpoints"
        / "planning_doc_analysis"
    )


def _atomic_write_json(path: Path, payload: Dict) -> None:
    """temp file + os.replace 로 atomic JSON write.

    부분 쓰기/동시 read 에 의한 corrupt JSON 방지. tmp 파일에 uuid 접미사를
    붙여 동시 다중 thread 충돌을 회피 (현재는 job_id 가드로 단일 writer 가
    보장되지만 안전 차원 dual layer).
    """
    path.parent.mkdir(parents=True, exist_ok=True)
    tmp = path.with_suffix(path.suffix + f".{uuid.uuid4().hex[:8]}.tmp")
    try:
        tmp.write_text(
            json.dumps(payload, ensure_ascii=False, indent=2),
            encoding="utf-8",
        )
        os.replace(str(tmp), str(path))
    except Exception:
        if tmp.exists():
            try:
                tmp.unlink()
            except OSError:
                pass
        raise


def compute_source_hash(project_id: str, db=None) -> str:
    """기획서 source 의 결정적 hash.

    PDF bytes (streaming) + 'PDF_END|TEXT:' separator + planning_doc_text
    의 sha256 결합. PDF 가 없거나 text 가 비어 있어도 hex digest 가 항상
    return (empty source 도 결정적 값).

    reader/writer 가 같은 입력으로 같은 hash 를 만들도록 separator 명시.
    """
    h = hashlib.sha256()
    pdf_path = (
        Path(settings.projects_dir) / project_id / "assets" / "planning_doc.pdf"
    )
    if pdf_path.exists():
        h.update(b"PDF:")
        try:
            with pdf_path.open("rb") as f:
                for chunk in iter(lambda: f.read(65536), b""):
                    h.update(chunk)
        except OSError as exc:
            logger.warning(
                "compute_source_hash: PDF read failed (project=%s): %s",
                project_id[:8], exc,
            )
            h.update(b"<read-error>")
    h.update(b"|TEXT:")
    own_db = False
    if db is None:
        from app.core.database import SessionLocal
        db = SessionLocal()
        own_db = True
    try:
        from app.models.catalog import ProjectRegistry
        project = (
            db.query(ProjectRegistry)
            .filter(ProjectRegistry.id == project_id)
            .first()
        )
        text = (project.planning_doc_text or "") if project else ""
        h.update(text.encode("utf-8"))
    finally:
        if own_db:
            db.close()
    return h.hexdigest()


def project_has_planning_doc(project_id: str, db=None) -> bool:
    """기획서 존재 여부 단일 판정 — 어느 한 곳이라도 있으면 True.

    검사 순서 (우선순위):
      1. project-level checkpoint — ``load_project_checkpoint`` 로 검증된
         data (source_hash 일치) + ``available_sections`` 비어 있지 않을 때만.
         empty/stale/error manifest 는 false-positive 차단.
      2. PDF on disk (assets/planning_doc.pdf)
      3. ProjectRegistry.planning_doc_text (≥100 chars)
      4. legacy episode-level checkpoint (옛 프로젝트 호환)

    db: caller 세션. None 이면 SessionLocal 일회 사용.

    `_if_planning_doc` (applicability) 와 `_project_has_planning_doc`
    (dispatcher) 가 동일 판정을 공유하도록 단일 source.
    """
    # 1) project-level checkpoint — load 까지 가서 source_hash 검증 통과 +
    # 의미 있는 데이터일 때만 True
    cp_data = load_project_checkpoint(project_id, db=db)
    if cp_data is not None:
        if cp_data.get("available_sections") or cp_data.get("characters"):
            return True

    # 2) PDF on disk
    pdf_path = (
        Path(settings.projects_dir) / project_id / "assets" / "planning_doc.pdf"
    )
    if pdf_path.exists():
        return True

    # 3) DB text
    own_db = False
    if db is None:
        from app.core.database import SessionLocal
        db = SessionLocal()
        own_db = True
    try:
        from app.models.catalog import ProjectRegistry
        project = (
            db.query(ProjectRegistry)
            .filter(ProjectRegistry.id == project_id)
            .first()
        )
        if project and len((project.planning_doc_text or "").strip()) >= 100:
            return True
    finally:
        if own_db:
            db.close()

    # 4) legacy episode-level checkpoint
    legacy_base = (
        Path(settings.projects_dir) / project_id / "checkpoints" / "episodes"
    )
    if legacy_base.exists():
        for ep_dir in legacy_base.iterdir():
            if (ep_dir / "planning_doc_analysis" / "manifest.json").exists():
                return True

    return False


def load_project_checkpoint(project_id: str, db=None) -> Optional[Dict]:
    """저장된 project-level 분석 결과 — source_hash 검증 통과 시에만 반환.

    검증 로직 (L4 — stale resurrection 최종 안전망):
      - manifest.source_hash 가 있고 현재 disk source 의 hash 와 다르면
        stale (delete 후 다른 source 로 재업로드 race 등) → ``None`` reject.
        **cleanup unlink 는 의도적으로 하지 않음** — hash 계산 도중 다른
        thread 가 fresh manifest 를 promote 하면, reader 가 fresh 를 unlink
        하는 race 가 생김. disk 정리는 다음 dispatch promote (덮어쓰기) 또는
        ``clear_project_checkpoint`` 가 담당.
      - manifest.source_hash 가 없으면 (옛 schema) 검증 skip + 정상 read.

    옛 schema (source_hash 없는 manifest) 는 한 번 정상 read 된 후, 새 분석
    트리거 시 새 schema 로 덮어씀.
    """
    cp = get_project_checkpoint_dir(project_id) / _MANIFEST_FILE
    if not cp.exists():
        return None
    try:
        raw = json.loads(cp.read_text(encoding="utf-8"))
    except (OSError, json.JSONDecodeError) as exc:
        logger.warning("planning_doc project checkpoint read failed: %s", exc)
        return None
    if not isinstance(raw, dict):
        return None

    manifest_hash = raw.get("source_hash")
    if isinstance(manifest_hash, str) and manifest_hash:
        try:
            current_hash = compute_source_hash(project_id, db=db)
        except Exception as exc:
            logger.warning(
                "planning_doc source_hash compute failed (project=%s): %s — "
                "skipping reader hash check (read as-is)",
                project_id[:8], exc,
            )
            current_hash = manifest_hash  # fail-open: 값 못 구하면 read 허용
        if manifest_hash != current_hash:
            # reject 만 — cleanup unlink 는 의도적으로 안 함.
            #
            # cleanup race: hash 계산 중 (큰 PDF) 새 dispatch B 가 fresh
            # manifest 를 promote 하면, reader 는 "처음 읽은 stale A" 가
            # 아니라 "지금 disk 의 fresh B" 를 unlink 하게 됨. 그래서 reject
            # only — disk 정리는 다음 dispatch promote (덮어쓰기) 또는
            # ``clear_project_checkpoint`` 에 위임.
            logger.warning(
                "planning_doc manifest source_hash mismatch (stale): project=%s "
                "manifest=%s current=%s — rejecting (no unlink)",
                project_id[:8], manifest_hash[:12], current_hash[:12],
            )
            return None

    data = raw.get("data")
    return data if isinstance(data, dict) else None


def load_status(project_id: str) -> Dict:
    """현재 분석 진행 상태. {state: idle|queued|running|done|error, ...}

    corrupt JSON 은 silent ``idle`` 대신 ``error/status_unreadable`` 로
    surface — 운영자가 disk 손상을 인지하도록.
    """
    cp = get_project_checkpoint_dir(project_id) / _STATUS_FILE
    if not cp.exists():
        return {"state": "idle"}
    try:
        return json.loads(cp.read_text(encoding="utf-8"))
    except (OSError, json.JSONDecodeError) as exc:
        logger.warning(
            "planning_doc status read failed (project=%s): %s — surfacing as error",
            project_id[:8], exc,
        )
        return {"state": "error", "error": "status_unreadable", "detail": str(exc)}


def _read_status_job_id(project_id: str) -> Optional[str]:
    """현재 disk status.json 의 job_id (없거나 손상 시 None)."""
    cp = get_project_checkpoint_dir(project_id) / _STATUS_FILE
    if not cp.exists():
        return None
    try:
        raw = json.loads(cp.read_text(encoding="utf-8"))
    except (OSError, json.JSONDecodeError):
        return None
    if isinstance(raw, dict):
        v = raw.get("job_id")
        return v if isinstance(v, str) else None
    return None


def _is_current_job(project_id: str, job_id: Optional[str]) -> bool:
    """job_id 가 현재 status.json 의 job_id 와 일치하면 True.

    job_id=None 으로 호출되면 guard 를 비활성 (sync 직접 호출용).
    """
    if job_id is None:
        return True
    current = _read_status_job_id(project_id)
    return current == job_id


def _write_status(
    project_id: str,
    payload: Dict,
    *,
    expected_job_id: Optional[str] = None,
) -> bool:
    """status.json atomic write. expected_job_id 가 주어지면 현재 disk
    job_id 와 일치할 때만 쓰고 stale 이면 drop.

    Returns True iff 실제로 write 됨.
    """
    if expected_job_id is not None and not _is_current_job(project_id, expected_job_id):
        logger.info(
            "planning_doc status write dropped (stale job_id): project=%s "
            "self=%s current=%s",
            project_id[:8], expected_job_id, _read_status_job_id(project_id),
        )
        return False
    _atomic_write_json(get_project_checkpoint_dir(project_id) / _STATUS_FILE, payload)
    return True


def _write_manifest(
    project_id: str,
    data: Dict,
    *,
    expected_job_id: Optional[str] = None,
    source_hash: Optional[str] = None,
) -> Optional[Path]:
    """manifest 를 staging file → promote 패턴으로 atomic 게시.

    이전의 "write-then-conditional-unlink" 는 race window 가 있었음:
        A: status check OK → atomic replace manifest.json (A data)
        B: 그 사이 clear + 새 dispatch + write manifest.json (B data)
        A: status check FAIL → unconditional out.unlink() → **B 의 fresh
            manifest 까지 삭제**

    staging 패턴:
      1. ``manifest.{job_id}.json`` 에 atomic write (job-isolated)
      2. status.job_id 가 자기 것인지 재확인
         - 자기 것 → ``os.replace(staging, manifest.json)`` 으로 promote
         - 아니면 → 자기 staging 만 unlink (manifest.json 은 건드리지 않음)
      3. promote 후 동시 다른 job 이 새 promote 해도 자기 staging 은 이미
         소비된 후 → 안전. manifest.json 의 owner 는 마지막 promote 가 결정.

    L1: 1차 status 가드 (early bail-out — staging write 비용 절약).
    L2: payload 에 ``job_id`` + ``source_hash`` → reader 가 L4 검증 가능.
    L3: promote 직전 status 재확인 + cross-job manifest.json 보호.

    expected_job_id=None 호출 (sync 직접 caller) 은 race-free 가정 → 단순
    atomic write. caller 는 자기 책임으로 호출 직전에 status/source 정합성
    확인.

    Returns 실제로 promote 된 Path (drop / staging-only cleanup 시 None).
    """
    out = get_project_checkpoint_dir(project_id) / _MANIFEST_FILE
    payload = {
        "data": data,
        "job_id": expected_job_id,
        "source_hash": source_hash,
        "generated_at": datetime.now(timezone.utc).isoformat(),
    }

    # job_id 가드 비활성 호출 — 단순 atomic write.
    if expected_job_id is None:
        _atomic_write_json(out, payload)
        return out

    # L1 — 1차 가드 (early bail-out)
    if not _is_current_job(project_id, expected_job_id):
        logger.info(
            "planning_doc manifest write dropped (1차 stale): project=%s self=%s current=%s",
            project_id[:8], expected_job_id, _read_status_job_id(project_id),
        )
        return None

    staging = out.parent / f"manifest.{expected_job_id}.json"
    _atomic_write_json(staging, payload)

    # L3 — promote 직전 가드. 자기 job_id 가 여전히 status SOT 이면 promote.
    # 아니면 자기 staging 만 정리하고, manifest.json 은 절대 건드리지 않음
    # (다른 job 의 fresh SOT 일 수 있음).
    if not _is_current_job(project_id, expected_job_id):
        logger.info(
            "planning_doc manifest promote dropped (2차 stale): project=%s self=%s current=%s",
            project_id[:8], expected_job_id, _read_status_job_id(project_id),
        )
        try:
            staging.unlink()
        except FileNotFoundError:
            pass
        except OSError as exc:
            logger.warning(
                "planning_doc staging cleanup failed: %s — %s", staging, exc,
            )
        return None

    try:
        os.replace(str(staging), str(out))
    except OSError as exc:
        logger.warning(
            "planning_doc manifest promote failed: project=%s — %s",
            project_id[:8], exc,
        )
        try:
            staging.unlink()
        except (FileNotFoundError, OSError):
            pass
        return None

    return out


def clear_project_checkpoint(project_id: str) -> None:
    """기획서 교체/삭제 시 호출 — manifest+status+staging file 모두 제거.

    staging file (``manifest.{job_id}.json``) 는 in-flight 분석이 promote
    하기 전 임시 파일. clear 시 이전 dispatch 의 잔재를 같이 청소.

    ``_project_lock`` 안에서 수행 — 동시에 발생하는 dispatch / status write
    /  promote 와 직렬화되어 partial state (예: status 만 남고 manifest 잔존)
    가 노출되지 않음.
    """
    with _project_lock(project_id):
        d = get_project_checkpoint_dir(project_id)
        if not d.exists():
            return
        for name in (_MANIFEST_FILE, _STATUS_FILE):
            p = d / name
            if p.exists():
                try:
                    p.unlink()
                except OSError as exc:
                    logger.warning("clear planning checkpoint failed: %s — %s", p, exc)
        # staging files (manifest.{job_id}.json)
        for f in d.glob("manifest.*.json"):
            if f.name == _MANIFEST_FILE:
                continue
            try:
                f.unlink()
            except OSError as exc:
                logger.warning("clear planning staging failed: %s — %s", f, exc)


def run_planning_doc_analysis_sync(
    project_id: str,
    *,
    job_id: Optional[str] = None,
    queued_at: Optional[str] = None,
) -> Dict:
    """동기 실행 — 기획서를 LLM 으로 분석 + project-level checkpoint 저장.

    구조: write/check 만 ``_project_lock`` 안에서 직렬화, LLM call 은 lock
    밖. concurrent dispatch / clear 가 status SOT 를 뒤집어도 자기 phase
    시작 시 ``_is_current_job(job_id)`` 로 안전하게 bail out.

    Phase 흐름:
      1. **lock**: running status 마킹 + bail-out 체크 + LLM input 수집
         (PDF bytes, text, source_hash)
      2. **lock 밖**: LLM call (PDF multimodal → text fallback)
      3. **lock**: bail-out 체크 + manifest staging→promote + done/error
         status 마킹

    job_id: dispatch 시 발급된 식별자. ``None`` 이면 ``_is_current_job`` /
        ``_write_*`` 의 stale-write guard 만 비활성화됨 (caller 가 race-free
        보장한다는 가정). ``_project_lock`` 자체는 항상 유지 — 동시 sync
        호출이나 clear 와 SOT mutation 직렬화는 그대로 보장.
    queued_at: dispatch 시점 timestamp (status payload 에 보존).
    """
    # late-import 로 모듈 로드 비용 + 순환 import 회피
    from app.core.database import SessionLocal
    from app.core.steps.planning_doc_step import (
        _ANALYSIS_SCHEMA,
        _EMPTY_RESULT,
        _SYSTEM_PROMPT,
    )
    from app.models.catalog import ProjectRegistry
    from app.modules.llm.llm_client import call_structured
    from app.services.analysis_dispatch_service import load_project_llm_config

    started_at = datetime.now(timezone.utc).isoformat()

    try:
        return _run_planning_doc_analysis_phases(
            project_id,
            job_id=job_id,
            queued_at=queued_at,
            started_at=started_at,
            session_local=SessionLocal,
            schema=_ANALYSIS_SCHEMA,
            empty_result=_EMPTY_RESULT,
            system_prompt=_SYSTEM_PROMPT,
            project_registry_cls=ProjectRegistry,
            call_structured=call_structured,
            load_project_llm_config=load_project_llm_config,
        )
    except Exception as exc:
        # Phase 1/2/3 어디서든 unexpected 예외 → status=error 마킹.
        # job_id stale 이면 _write_status guard 가 drop. lock 안에서 atomic.
        try:
            with _project_lock(project_id):
                _write_status(
                    project_id,
                    {
                        "state": "error",
                        "job_id": job_id,
                        "queued_at": queued_at,
                        "started_at": started_at,
                        "finished_at": datetime.now(timezone.utc).isoformat(),
                        "error": str(exc),
                    },
                    expected_job_id=job_id,
                )
        except Exception as e2:
            logger.exception(
                "planning_doc_analysis error-status write failed (project=%s): %s",
                project_id[:8], e2,
            )
        logger.exception(
            "planning_doc_analysis (project=%s) failed: %s",
            project_id[:8], exc,
        )
        raise


def _run_planning_doc_analysis_phases(
    project_id: str,
    *,
    job_id: Optional[str],
    queued_at: Optional[str],
    started_at: str,
    session_local,
    schema,
    empty_result,
    system_prompt: str,
    project_registry_cls,
    call_structured,
    load_project_llm_config,
) -> Dict:
    """run_planning_doc_analysis_sync 의 본체 — outer try/except 가 status
    error 마킹을 보장하므로 본 함수는 정상 흐름만 책임.
    """
    # ── Phase 1 (lock): running status + bail-out + input 수집 ──────
    db = session_local()
    try:
        with _project_lock(project_id):
            if job_id is not None and not _is_current_job(project_id, job_id):
                logger.info(
                    "planning_doc_analysis bail at phase 1 (stale job): "
                    "project=%s self=%s current=%s",
                    project_id[:8], job_id, _read_status_job_id(project_id),
                )
                return empty_result

            project = (
                db.query(project_registry_cls)
                .filter(project_registry_cls.id == project_id)
                .first()
            )
            planning_text = project.planning_doc_text if project else None
            pdf_path = (
                Path(settings.projects_dir)
                / project_id
                / "assets"
                / "planning_doc.pdf"
            )
            has_pdf = pdf_path.exists()

            # 사용자 모델 override + 분석 시점 source_hash (lock 안에서 일관 read)
            project_llm_config = load_project_llm_config(db, project_id)
            source_hash = compute_source_hash(project_id, db=db)

            _write_status(
                project_id,
                {
                    "state": "running",
                    "job_id": job_id,
                    "queued_at": queued_at,
                    "started_at": started_at,
                },
                expected_job_id=job_id,
            )

            # 빈 source 인 경우 LLM call 없이 lock 안에서 즉시 종료
            if not has_pdf and (
                not planning_text or len(planning_text.strip()) < 100
            ):
                logger.info(
                    "planning_doc_analysis (project=%s): 기획서 없음 — 빈 결과 저장",
                    project_id[:8],
                )
                _write_manifest(
                    project_id, empty_result,
                    expected_job_id=job_id, source_hash=source_hash,
                )
                _write_status(
                    project_id,
                    {
                        "state": "done",
                        "job_id": job_id,
                        "queued_at": queued_at,
                        "started_at": started_at,
                        "finished_at": datetime.now(timezone.utc).isoformat(),
                        "empty": True,
                    },
                    expected_job_id=job_id,
                )
                return empty_result

            # PDF bytes 는 lock 안에서 읽음 (lock 밖에서 source 가 바뀌어도
            # 우리는 분석 시점 snapshot 으로 진행. 어차피 새 dispatch 가 들어
            # 오면 promote bail-out 됨).
            if has_pdf:
                pdf_bytes = pdf_path.read_bytes()
            else:
                pdf_bytes = None
    finally:
        db.close()

    # ── Phase 2 (lock 밖): LLM call ─────────────────────────────────
    result = None
    if pdf_bytes is not None:
        import base64
        pdf_b64 = base64.b64encode(pdf_bytes).decode("utf-8")
        user_prompt = [
            {
                "type": "text",
                "text": "이 기획서 PDF를 분석하여 구조화된 정보를 추출하세요.",
            },
            {
                "type": "image_url",
                "image_url": {
                    "url": f"data:application/pdf;base64,{pdf_b64}",
                },
            },
        ]
        logger.info(
            "planning_doc_analysis (project=%s): PDF multimodal (%d bytes)",
            project_id[:8], len(pdf_bytes),
        )
        try:
            result = call_structured(
                step="planning_doc_analysis",
                system_prompt=system_prompt,
                user_prompt=user_prompt,
                response_schema=schema,
                project_config=project_llm_config or None,
                temperature=0.2,
            )
        except Exception as exc:
            logger.warning(
                "PDF multimodal failed, falling back to text: %s", exc,
            )

    if (
        result is None
        and planning_text
        and len(planning_text.strip()) >= 100
    ):
        user_prompt_text = (
            "다음 기획서를 분석하여 구조화된 정보를 추출하세요.\n\n"
            f"## 기획서 전문\n\n{planning_text}"
        )
        logger.info(
            "planning_doc_analysis (project=%s): text fallback (%d chars)",
            project_id[:8], len(planning_text),
        )
        # text fallback 도 실패하면 raise — outer except 가 status=error 마킹.
        result = call_structured(
            step="planning_doc_analysis",
            system_prompt=system_prompt,
            user_prompt=user_prompt_text,
            response_schema=schema,
            project_config=project_llm_config or None,
            temperature=0.2,
        )

    # ── Phase 3 (lock): promote + done/error status ─────────────────
    with _project_lock(project_id):
        # bail-out: 분석 도중 다른 dispatch 가 들어왔으면 promote 안 함
        if job_id is not None and not _is_current_job(project_id, job_id):
            logger.info(
                "planning_doc_analysis bail at phase 3 (stale job): "
                "project=%s self=%s current=%s — discard result",
                project_id[:8], job_id, _read_status_job_id(project_id),
            )
            return result if result is not None else empty_result

        if result is None:
            logger.warning(
                "planning_doc_analysis (project=%s): both PDF and text failed",
                project_id[:8],
            )
            _write_manifest(
                project_id, empty_result,
                expected_job_id=job_id, source_hash=source_hash,
            )
            _write_status(
                project_id,
                {
                    "state": "error",
                    "job_id": job_id,
                    "queued_at": queued_at,
                    "started_at": started_at,
                    "finished_at": datetime.now(timezone.utc).isoformat(),
                    "error": "PDF 와 text 분석 모두 실패",
                },
                expected_job_id=job_id,
            )
            return empty_result

        # available_sections 코드에서 결정 (기존 step 과 동일 규칙).
        computed = []
        if result.get("characters"):
            computed.append("characters")
        if result.get("world_setting", "").strip():
            computed.append("world_setting")
        if result.get("tone_mood", "").strip():
            computed.append("tone_mood")
        if result.get("story_arc", "").strip():
            computed.append("story_arc")
        if result.get("visual_concepts", "").strip():
            computed.append("visual_concepts")
        if result.get("key_relationships"):
            computed.append("key_relationships")
        result["available_sections"] = computed

        wrote = _write_manifest(
            project_id, result,
            expected_job_id=job_id, source_hash=source_hash,
        )
        if wrote is None:
            return result
        _write_status(
            project_id,
            {
                "state": "done",
                "job_id": job_id,
                "queued_at": queued_at,
                "started_at": started_at,
                "finished_at": datetime.now(timezone.utc).isoformat(),
                "characters": len(result.get("characters", [])),
                "relationships": len(result.get("key_relationships", [])),
                "sections": computed,
            },
            expected_job_id=job_id,
        )
        logger.info(
            "planning_doc_analysis (project=%s) done: %d characters, %d "
            "relationships, sections=%s, job_id=%s",
            project_id[:8],
            len(result.get("characters", [])),
            len(result.get("key_relationships", [])),
            computed,
            job_id,
        )
        return result


def dispatch_planning_doc_analysis_async(project_id: str) -> str:
    """업로드 직후 백그라운드 thread 로 분석 실행 + ``job_id`` 반환.

    동작:
      - ``_project_lock`` 안에서 uuid ``job_id`` 발급 + ``status.json`` 에
        ``{state:queued, job_id}`` 도장. concurrent dispatch 와 직렬화되어
        마지막 dispatch 의 job_id 가 SOT 가 됨.
      - background thread 는 자기 ``job_id`` 를 들고 ``run_..._sync`` 호출.
        phase 마다 자기 job_id 가 still current 인지 lock 안에서 확인.

    응답 지연 회피 + 결과는 ``checkpoints/planning_doc_analysis/manifest.json``.
    진행 상태는 ``status.json`` 으로 polling 가능.
    """
    job_id = uuid.uuid4().hex
    queued_at = datetime.now(timezone.utc).isoformat()
    # queued status 도장 — lock 안에서 atomic. concurrent dispatch 가 둘
    # 동시에 들어와도 마지막 dispatch 의 job_id 가 status SOT 로 남음.
    with _project_lock(project_id):
        _write_status(project_id, {
            "state": "queued",
            "job_id": job_id,
            "queued_at": queued_at,
        })

    def _run():
        try:
            run_planning_doc_analysis_sync(
                project_id, job_id=job_id, queued_at=queued_at,
            )
        except Exception as exc:
            # status.json 은 _run_sync 안에서 이미 error 마킹 (자기 job_id 일 때).
            logger.error(
                "planning_doc_analysis background thread crashed: %s", exc,
            )

    t = threading.Thread(
        target=_run,
        name=f"planning_doc_analysis-{project_id[:8]}-{job_id[:8]}",
        daemon=True,
    )
    t.start()
    return job_id
