"""CaptureQueue — append/flush 로 중간 생성물을 image_asset 으로 영속화 (Task A4).

저수준 capture(sink)는 spool 파일 경로 + 메타만 append 한다. step/runner 의
generation_context 종료 시 finally 에서 flush() 가:
  1. 독립 SessionLocal() 세션을 열고
  2. 각 항목을 ImageAsset(is_intermediate=True, ...) 로 insert
  3. spool 파일을 최종 generated 경로로 rename
  4. commit / close
한다. 전체를 try/except 로 감싸 어떤 실패도 raise 하지 않는다(독립 non-fatal
세션 — business transaction rollback 무관). 실패 시 rollback + CAPTURE_DIAG 증가.

★ reference_image_ids 는 절대 건드리지 않는다(char/prop lineage 보존). 전체 입력
lineage 는 input_image_ids 컬럼에 JSON 으로 저장한다.
"""

from __future__ import annotations

import json
import logging
import re
import threading
import uuid
from datetime import datetime, timezone
from pathlib import Path
from typing import TYPE_CHECKING, Any, Dict, List, Tuple

from app.core.database import SessionLocal

if TYPE_CHECKING:
    from app.services.image_capture.context import GenerationContext

logger = logging.getLogger(__name__)

# 누락 감사용 diagnostic counter — silent drop 금지.
#   skipped_no_context     : capture scope 밖 호출(=최종물 경로, default-off)
#   skipped_no_flush_scope : 예약(Phase B — queue handle 기대했으나 부재한 worker 경로)
#   flush_errors           : flush/sink 내부 비치명 실패
CAPTURE_DIAG: Dict[str, int] = {
    "skipped_no_context": 0,
    "skipped_no_flush_scope": 0,
    #   skipped_capture_off : 신원·trace 만 얻으려고 capture=False 로 연 scope
    #                         (2026-08-23) — 자산을 일부러 안 만든 횟수
    "skipped_capture_off": 0,
    "flush_errors": 0,
}

# 경로 세그먼트 보안 필터(의미추출 아님 — path traversal 방지용 문자 화이트리스트).
_UNSAFE_SEG_RE = re.compile(r"[^A-Za-z0-9_.-]")


def _safe_segment(stage: str | None) -> str:
    """stage 를 안전한 단일 path 세그먼트로. DB 엔 원문 저장, 파일 경로에만 사용.

    허용 문자 ``[A-Za-z0-9_.-]`` 외는 ``_`` 치환. ``.``/``..``/빈 문자열 등
    상위 디렉터리 탈출 위험 값은 ``unknown`` 으로(generated root 밖 금지)."""
    seg = _UNSAFE_SEG_RE.sub("_", stage or "")
    if seg in ("", ".", "..") or set(seg) == {"."}:
        return "unknown"
    return seg


def _now() -> str:
    return datetime.now(timezone.utc).isoformat()


def _generated_dir(project_id: str, episode_id: str | None, stage: str) -> Path:
    """최종 generated 디렉터리(절대). episode 무관 자산은 project-level fallback.

    stage 는 ``_safe_segment`` 로 정화해 경로 탈출을 막는다(Codex BLOCKING 2)."""
    from app.core.config import settings

    base = Path(settings.projects_dir) / project_id
    if episode_id:
        base = base / "episodes" / episode_id
    return base / "images" / "generated" / _safe_segment(stage)


class CaptureQueue:
    """단일 capture scope 의 flush 큐. append=적재, flush=독립세션 일괄 영속화."""

    def __init__(self, ctx: "GenerationContext") -> None:
        self._ctx = ctx
        self._items: List[Tuple[str, Dict[str, Any]]] = []
        # 병렬 롤 생성(still-variants parallel_rolls)에서 여러 worker 가
        # 같은 scope 큐에 동시 append — 유실 없는 적재를 lock 으로 계약
        # (2026-07-17 Codex BLOCKING-3).
        self._lock = threading.Lock()

    def append(self, spool_path: str, meta: Dict[str, Any]) -> None:
        """spool 파일 경로 + 메타를 적재(DB 접근 없음). thread-safe."""
        with self._lock:
            self._items.append((spool_path, dict(meta or {})))

    def flush(self) -> int:
        """독립 세션으로 적재 항목을 ImageAsset 으로 insert + spool→final rename.

        Returns:
            insert 한 행 수. 어떤 실패도 raise 하지 않는다(non-fatal). 실패 시
            rollback + CAPTURE_DIAG["flush_errors"] 증가 후 0 반환. 큐는 항상 비운다.
        """
        with self._lock:
            if not self._items:
                return 0
            items = self._items
            self._items = []  # 재flush 방지 — 시도 즉시 큐 비움.

        ctx = self._ctx
        session = None
        inserted = 0
        try:
            # FK relation 해소(image_asset.project_id → project_registry 등). catalog 미등록
            # 상태(runtime script / worker / standalone image path)에서 flush 시
            # NoReferencedTableError 로 전부 drop되는 것 방지. register_models 는
            # idempotent / lock-free(DDL/DML 없음). (Codex Phase A BLOCKING 1)
            from app.core.database import register_models

            register_models()
            from app.models.project import ImageAsset

            session = SessionLocal()
            for spool_path, meta in items:
                aid = str(uuid.uuid4())
                stage = meta.get("stage") or ctx.stage
                gen_dir = _generated_dir(ctx.project_id, ctx.episode_id, stage)
                gen_dir.mkdir(parents=True, exist_ok=True)
                final_path = gen_dir / f"{aid}.png"

                # spool → final rename (cross-device 안전 fallback).
                sp = Path(spool_path)
                try:
                    sp.rename(final_path)
                except OSError:
                    final_path.write_bytes(sp.read_bytes())
                    try:
                        sp.unlink()
                    except OSError:
                        pass

                input_ids = meta.get("input_image_ids")
                pipeline_meta = meta.get("pipeline_metadata")
                # scene_index 는 ImageAsset 컬럼이 없다(shot_index 만 존재). still-less
                # 중간물(outdoor aerial base 등)도 캔버스가 씬으로 묶으려면 scene_index 가
                # 필요하므로 pipeline_metadata 에 fold 한다. non-overwrite — 사이트가 명시한
                # 값이 ctx 값보다 우선(setdefault). scene_index None 이면 키 부재(기존 동작
                # 보존). (Phase C Decision #3/#4 — Codex 합의)
                pm = dict(pipeline_meta) if pipeline_meta is not None else None
                if ctx.scene_index is not None:
                    if pm is None:
                        pm = {}
                    pm.setdefault("scene_index", ctx.scene_index)
                row = ImageAsset(
                    id=aid,
                    project_id=ctx.project_id,
                    episode_id=ctx.episode_id,
                    asset_type="generated",
                    is_intermediate=True,
                    stage=stage,
                    pipeline_role=meta.get("pipeline_role"),
                    input_image_ids=(
                        json.dumps(input_ids) if input_ids is not None else None
                    ),
                    generation_call_id=meta.get("generation_call_id"),
                    candidate_index=meta.get("candidate_index"),
                    disposition=meta.get("disposition"),
                    attempt_index=meta.get("attempt_index"),
                    pipeline_metadata_json=(
                        json.dumps(pm) if pm is not None else None
                    ),
                    entity_id=ctx.entity_id,
                    still_id=ctx.still_id,
                    shot_index=ctx.shot_index,
                    prompt_used=meta.get("prompt"),
                    # ImagePathType 가 bind 시 자동 상대화 — 절대 경로 그대로 넣어도 안전.
                    file_path=str(final_path),
                    status="generated",
                    created_at=_now(),
                )
                session.add(row)
                inserted += 1

            session.commit()
            return inserted
        except Exception:
            # non-fatal drop. 주의: spool→final rename 이 이미 끝난 항목은 DB row 없이
            # generated dir 에 orphan 파일로 남을 수 있다(후속 cleanup 후보 — Codex MINOR).
            CAPTURE_DIAG["flush_errors"] += 1
            logger.warning(
                "CaptureQueue.flush failed (non-fatal) — %d item(s) dropped "
                "(renamed files may be orphaned in generated dir)",
                len(items),
                exc_info=True,
            )
            if session is not None:
                try:
                    session.rollback()
                except Exception:  # pragma: no cover
                    pass
            return 0
        finally:
            if session is not None:
                try:
                    session.close()
                except Exception:  # pragma: no cover
                    pass
