"""주행을 **멈추는 길** — 프로세스를 죽이지 않고 세우는 방법.

## 왜 있나

이 시스템에는 주행을 멈출 방법이 프로세스를 죽이는 것밖에 없었다. 그런데
프로세스를 죽이면 `step_run` 에 `running` 행이 남고, 락은 소유자의 죽음을
확인할 길이 없어 경과 시간(기본 3600초)을 기다렸다. **멈추는 길이 없어서
락이 망가지는 순환**이었다.

`step_run.cancel_requested_at` 은 **스텝 하나**를 세운다. 그것만으로는
부족하다 — `run_steps_batch` 는 `for sid in step_ids:` 로 도는지라
(`analysis_dispatch_service.py`), 지금 도는 스텝이 멈춰도 **다음 스텝으로
넘어간다.** 그래서 주행 전체를 덮는 표식이 따로 필요하다.

## 무엇을 하나

`run_cancel_request` 한 행이 「이 에피소드의 주행을 세워라」를 뜻한다.
읽는 자리는 셋이다:

1. **스텝 사이** — `run_steps_batch` 가 다음 `sid` 로 넘어가기 전
2. **단위 사이** — 팬아웃 스텝이 샷 하나를 끝낼 때 (`checkpoint_gate`)
3. **돈 쓰기 직전** — worker 가 provider 를 부르기 전

3번이 중요하다. 팬아웃은 작업을 executor 에 **한꺼번에 제출**하므로,
대기열에 이미 들어간 작업은 부모가 멈춰도 시작된다. 그 작업들이 유료
호출 앞에서 스스로 물러나야 돈이 안 나간다.

★멈춤은 **협조적**이다. 진행 중인 provider 호출을 중간에 끊지 않는다 —
 이미 산 것은 결과까지 가져온다. 끊어 봐야 돈은 이미 나갔고 결과만 잃는다.
"""

from __future__ import annotations

import logging
import threading
import uuid
from dataclasses import dataclass
from datetime import datetime, timezone
from typing import Optional

from sqlalchemy import text

logger = logging.getLogger(__name__)


CREATE_TABLE_SQL = """
CREATE TABLE IF NOT EXISTS run_cancel_request (
    id text PRIMARY KEY,
    project_id text NOT NULL,
    episode_id text NOT NULL,
    scope text NOT NULL DEFAULT 'episode',
    requested_at timestamptz NOT NULL DEFAULT CURRENT_TIMESTAMP,
    requested_by text,
    reason text,
    cleared_at timestamptz,
    CONSTRAINT run_cancel_request_scope_key UNIQUE(project_id, episode_id, scope)
)
"""


@dataclass(frozen=True)
class CancelState:
    """지금 이 에피소드에 정지 요청이 걸려 있나."""

    requested: bool
    requested_at: Optional[datetime] = None
    requested_by: Optional[str] = None
    reason: Optional[str] = None

    def describe(self) -> str:
        if not self.requested:
            return "정지 요청 없음"
        who = self.requested_by or "?"
        when = self.requested_at.isoformat() if self.requested_at else "?"
        why = self.reason or "(사유 없음)"
        return f"정지 요청 — {who} 가 {when} 에, 사유: {why}"


def request_cancel(
    db,
    project_id: str,
    episode_id: str,
    *,
    reason: str = "",
    requested_by: str = "",
    scope: str = "episode",
) -> CancelState:
    """정지를 요청한다. 이미 걸려 있으면 사유만 갱신한다 (멱등).

    ★시각은 DB 시계로 찍는다 — 읽는 쪽과 같은 시계여야 한다.
    """
    db.execute(text("""
        INSERT INTO run_cancel_request (
            id, project_id, episode_id, scope,
            requested_at, requested_by, reason, cleared_at
        )
        VALUES (
            :id, :pid, :eid, :scope,
            CURRENT_TIMESTAMP, :by, :reason, NULL
        )
        ON CONFLICT (project_id, episode_id, scope) DO UPDATE
          SET requested_at = CURRENT_TIMESTAMP,
              requested_by = EXCLUDED.requested_by,
              reason = EXCLUDED.reason,
              cleared_at = NULL
    """), {
        "id": str(uuid.uuid4()),
        "pid": project_id,
        "eid": episode_id,
        "scope": scope,
        "by": (requested_by or None),
        "reason": (reason or None),
    })
    db.commit()
    logger.warning(
        "[RUN-CANCEL] 정지 요청 project=%s episode=%s by=%s reason=%s",
        project_id, episode_id, requested_by or "?", reason or "-",
    )
    return read_cancel_state(db, project_id, episode_id, scope=scope)


def clear_cancel(
    db,
    project_id: str,
    episode_id: str,
    *,
    scope: str = "episode",
    expected_requested_at: Optional[datetime] = None,
) -> bool:
    """정지 요청을 내린다 — 다시 주행할 수 있게 한다.

    ★행을 지우지 않고 `cleared_at` 을 찍는다. 언제 세웠고 언제 풀었는지가
     남아야 나중에 읽을 수 있다.

    Args:
        expected_requested_at: 본 요청의 시각. 주면 **그 요청일 때만** 내린다.

    ★`expected_requested_at` 이 왜 있나 — 조회하고 내리기까지 사이에 누가
     **새 정지**를 요청했을 수 있다. 확인하지 않으면 내가 본 옛 요청을 푼다는
     것이 남의 새 요청까지 지운다. 「멈춰라」가 조용히 사라지는 것은
     「멈추지 못하는 것」과 같다.

    Returns: 실제로 내린 것이 있으면 True.
    """
    params: dict = {"pid": project_id, "eid": episode_id, "scope": scope}
    cas_clause = ""
    if expected_requested_at is not None:
        cas_clause = "AND requested_at = :expected_at"
        params["expected_at"] = expected_requested_at

    result = db.execute(text(f"""
        UPDATE run_cancel_request
           SET cleared_at = CURRENT_TIMESTAMP
         WHERE project_id = :pid AND episode_id = :eid AND scope = :scope
           AND cleared_at IS NULL
           {cas_clause}
    """), params)
    db.commit()
    cleared = bool(result.rowcount)
    if cleared:
        logger.info(
            "[RUN-CANCEL] 정지 요청 해제 project=%s episode=%s",
            project_id, episode_id,
        )
    return cleared


def read_cancel_state(
    db,
    project_id: str,
    episode_id: str,
    *,
    scope: str = "episode",
    strict: bool = False,
) -> CancelState:
    """지금 정지 요청이 걸려 있나.

    Args:
        strict: 읽기 실패를 **올릴지**. 기본 False.

    ★기본값은 읽기가 실패하면 **걸려 있지 않다**고 답한다. 표가 아직 없거나
     DB 가 잠깐 끊긴 것으로 정상 주행이 멈추는 쪽이 더 나쁘기 때문이다.

    ★그러나 **돈을 쓰기 직전**에는 이 관용이 거짓말이 된다. 호출자가
     「fail-closed 로 확인했다」고 믿는데 여기서 조용히 「정지 없음」을
     돌려주면, DB 장애 동안 정지가 통째로 안 들리면서 지출은 계속된다.
     그 자리에서는 `strict=True` 로 불러 실패를 올린다.
    """
    try:
        row = db.execute(text("""
            SELECT requested_at, requested_by, reason
              FROM run_cancel_request
             WHERE project_id = :pid AND episode_id = :eid AND scope = :scope
               AND cleared_at IS NULL
        """), {"pid": project_id, "eid": episode_id, "scope": scope}).fetchone()
    except Exception as exc:
        if strict:
            raise
        logger.debug("정지 요청 조회 실패 (없는 것으로 본다): %s", exc)
        return CancelState(requested=False)

    if row is None:
        return CancelState(requested=False)
    return CancelState(
        requested=True,
        requested_at=row.requested_at,
        requested_by=row.requested_by,
        reason=row.reason,
    )


class CancellationToken:
    """worker 가 유료 호출 직전에 물어보는 표.

    팬아웃은 작업을 한꺼번에 제출하므로, 부모가 멈춰도 대기열의 작업은
    시작된다. 그 작업들이 **돈을 쓰기 전에** 이 표를 보고 물러나야 한다.

    ★DB 를 매번 때리지 않는다. `ttl_seconds` 동안 답을 재사용한다 —
     샷 하나가 수십 초 걸리므로 몇 초 지연은 무해하고, 수십 개 worker 가
     동시에 조회하는 것은 유해하다.

    ★세션을 공유하지 않는다. 스레드마다 다른 세션이 필요하므로
     `session_factory` 를 받아 조회할 때마다 열고 닫는다.
    """

    def __init__(
        self,
        session_factory,
        project_id: str,
        episode_id: str,
        *,
        ttl_seconds: float = 3.0,
    ) -> None:
        self._session_factory = session_factory
        self._project_id = project_id
        self._episode_id = episode_id
        self._ttl = ttl_seconds
        self._lock = threading.Lock()
        self._cached: Optional[CancelState] = None
        self._cached_at: Optional[datetime] = None

    def state(self, *, force: bool = False) -> CancelState:
        now = datetime.now(timezone.utc)
        with self._lock:
            fresh = (
                not force
                and self._cached is not None
                and self._cached_at is not None
                and (now - self._cached_at).total_seconds() < self._ttl
            )
            if fresh and self._cached is not None:
                return self._cached

        session = self._session_factory()
        try:
            state = read_cancel_state(session, self._project_id, self._episode_id)
        finally:
            session.close()

        with self._lock:
            self._cached = state
            self._cached_at = now
        return state

    def is_cancelled(self, *, force: bool = False) -> bool:
        return self.state(force=force).requested

    def raise_if_cancelled(self, where: str = "") -> None:
        """멈추라는 말이 왔으면 즉시 멈춘다.

        Raises: AppError(step.cancelled)
        """
        state = self.state()
        if not state.requested:
            return
        from app.core.errors import AppError

        raise AppError(
            code="step.cancelled",
            message=f"{where or '주행'} 정지 — {state.describe()}",
            status_code=409,
        )


#: ★★**이 실패들은 「일 하나가 실패한 것」이 아니다.** 병렬 loop 안에서
#:  잡아 세면 나머지 worker 가 계속 돌아 **돈이 나가고**, 재시도가 다시 부르며,
#:  감사에는 provider 장애로 남는다. 만나면 **즉시 위로 올린다**.
#:
#:  ★한 곳에 둔다 — `detail_steps` 가 같은 목록을 **두 번 인라인**으로 적어
#:  뒀고, 그런 것은 한쪽만 고쳐진다 (2026-08-31).
ABORT_CODES: frozenset = frozenset({
    "step.cancelled",       # 사람이 멈추라고 했다
    "step.owner_lost",      # 이미 남의 주행이 됐다
    "step.gate_unreadable", # 게이트를 못 읽는다 — 「통과」로 볼 수 없다
})


def is_abort(exc: BaseException) -> bool:
    """이 예외가 **주행을 세워야 하는 것**인가. ★code 로 본다 — 문구가 아니라."""
    return getattr(exc, "code", "") in ABORT_CODES


__all__ = [
    "CREATE_TABLE_SQL",
    "CancelState",
    "CancellationToken",
    "request_cancel",
    "clear_cancel",
    "read_cancel_state",
    "ABORT_CODES",
    "is_abort",
]
