"""스텝 락의 소유자 신원과 살아있음 — 죽음을 추측하지 않고 확인한다.

## 왜 있나

`step_run.status='running'` 한 칸이 락이었다. 그 칸은 **누가 잡았는지**도
**그가 살아있는지**도 적지 않아서, 시스템은 죽음을 **경과 시간으로 추측**했다
(`step_running_timeout_seconds`, 기본 3600초). 그래서:

- 프로세스를 강제 종료하면 `running` 행이 남고 한 시간을 기다려야 했다
  (2026-08-26 새벽 실측 손실 40분)
- `force` 로도 못 뺏었다 (`_try_claim_running` 의 `WHERE status != 'running'`)
- 다섯 달치 시체가 쌓였다 (2026-03-25 부터 14행)

이 모듈은 락에 **소유자 신원**과 **하트비트**를 실어, 죽음을 추측 대신
**확인**하게 한다.

## 판정 순서 (`judge_owner`)

위에서부터 첫 일치에서 멈춘다.

| # | 조건 | 판정 |
|---|---|---|
| 1 | boot_id 가 이 프로세스 **and** 등록부에 없다 | DEAD |
| 2 | boot_id 가 이 프로세스 **and** 등록부에 있다 | ALIVE |
| 3 | host 가 이 호스트 **and** PID 가 죽은 프로세스 | DEAD |
| 4 | host 가 이 호스트 **and** PID 는 살아있으나 boot_id 다름 | 하트비트로 |
| 5 | 다른 호스트 / 신원 미상 | 하트비트 lease 로 |
| 6 | 신원·하트비트 둘 다 없음 (구 행) | UNKNOWN — 호출자가 경과 시간으로 |

`UNKNOWN` 은 「모른다」다. 호출자는 이를 **살아있음으로 취급**하거나 기존
경과 시간 경로로 떨어뜨린다 — 절대 죽음으로 읽지 않는다 (fail-closed).
"""

from __future__ import annotations

import logging
import os
import socket
import threading
import uuid
from dataclasses import dataclass
from datetime import datetime, timezone
from enum import Enum
from typing import Dict, Optional, Tuple

logger = logging.getLogger(__name__)


# ── 이 프로세스의 신원 ──
#
# boot_id 는 프로세스마다 다른 uuid4 다. PID 는 재사용되지만 uuid4 는 아니라서,
# 같은 host+pid 로 다른 프로세스가 떠도 구별된다.
#
# ★import 시점에 고정하지 않는다. uvicorn 을 pre-fork 로 띄우면 자식들이
#  부모의 boot_id 를 **그대로 물려받아** 서로를 자기 자신으로 오인한다.
#  그래서 부를 때마다 PID 를 확인하고, 바뀌었으면 새로 만든다.

PROCESS_HOST: str = socket.gethostname()

_IDENTITY_LOCK = threading.Lock()
_IDENTITY_PID: int = os.getpid()
_IDENTITY_BOOT_ID: str = str(uuid.uuid4())


def process_identity() -> Tuple[str, int, str]:
    """(host, pid, boot_id). fork 로 PID 가 바뀌었으면 신원을 새로 만든다.

    ★신원을 새로 만들 때 **등록부도 비운다.** fork 직후 자식은 부모의 등록부를
     통째로 물려받는데, 그 내용은 전부 거짓이다 — 자식은 그 스텝들을 하고
     있지 않다. 안 비우면 자식이 부모의 락을 「내가 잡고 있다」고 읽는다.
    """
    global _IDENTITY_PID, _IDENTITY_BOOT_ID

    pid = os.getpid()
    with _IDENTITY_LOCK:
        if _IDENTITY_PID != pid:
            _IDENTITY_PID = pid
            _IDENTITY_BOOT_ID = str(uuid.uuid4())
            REGISTRY.reset()
            logger.warning(
                "프로세스 신원 재생성 (PID %s) — fork 로 보인다. 등록부를 비웠다.",
                pid,
            )
        return PROCESS_HOST, _IDENTITY_PID, _IDENTITY_BOOT_ID


def process_boot_id() -> str:
    return process_identity()[2]


def process_pid() -> int:
    return process_identity()[1]


class OwnerVerdict(str, Enum):
    """소유자가 살아있는가에 대한 판정."""

    ALIVE = "alive"
    DEAD = "dead"
    UNKNOWN = "unknown"


@dataclass(frozen=True)
class LockOwner:
    """`step_run` 행에서 읽은 소유자 신원.

    전부 nullable — 신원 칸이 생기기 전에 만들어진 행은 모두 None 이다.
    """

    host: Optional[str] = None
    pid: Optional[int] = None
    boot_id: Optional[str] = None
    heartbeat_at: Optional[str] = None
    run_id: Optional[str] = None

    @classmethod
    def from_row(cls, row: Dict) -> "LockOwner":
        """`_get_step_run()` 이 준 dict 에서 신원만 뽑는다."""
        raw_pid = row.get("owner_pid")
        try:
            pid = int(raw_pid) if raw_pid is not None else None
        except (TypeError, ValueError):
            pid = None
        return cls(
            host=row.get("owner_host"),
            pid=pid,
            boot_id=row.get("owner_boot_id"),
            heartbeat_at=row.get("heartbeat_at"),
            run_id=row.get("run_id"),
        )

    @property
    def is_this_process(self) -> bool:
        return bool(self.boot_id) and self.boot_id == process_boot_id()

    def describe(self) -> str:
        if not self.boot_id and not self.host:
            return "owner=<신원 없음>"
        return (
            f"owner={self.host or '?'}:{self.pid if self.pid is not None else '?'}"
            f":{(self.boot_id or '?')[:8]}"
        )


# ── 이 프로세스가 잡고 있는 락의 등록부 ──

LockKey = Tuple[str, str, str]  # (project_id, episode_id, step_id)


class HeldLockRegistry:
    """이 프로세스가 **지금 실제로 일하고 있는** 락의 목록.

    DB 가 「이 프로세스가 소유자」라고 말하는데 여기에 없으면, 그 일을 하던
    스레드는 사라진 것이다 — 프로세스 안의 일은 프로세스가 안다.

    ★등록은 **claim 보다 먼저** 한다. claim 성공 직후·등록 직전의 창에서
     다른 스레드가 판정하면 살아있는 락을 죽었다고 읽기 때문이다. claim 이
     실패하면 등록을 되돌린다.

    ★한 키에 run_id 를 **하나만** 담으면 안 된다. 같은 스텝을 노리는 두
     스레드 A·B 가 있을 때, 나중에 등록한 B 가 A 를 덮어쓰고, claim 에 진 B 가
     물러나며 지우면 **A 의 항목까지 사라진다.** 그 뒤 C 가 판정하면
     「DB 는 이 프로세스가 소유자라는데 등록부에 없다」 → 살아있는 A 를 죽었다고
     읽는다. 그래서 키마다 run_id **집합**을 담는다 — 진 쪽의 퇴장이 산 쪽을
     건드리지 않는다.
    """

    def __init__(self) -> None:
        self._lock = threading.Lock()
        self._held: Dict[LockKey, set] = {}  # key → {run_id, ...}

    def register(self, key: LockKey, run_id: str) -> None:
        with self._lock:
            self._held.setdefault(key, set()).add(run_id)

    def unregister(self, key: LockKey, run_id: str) -> None:
        """내 run_id 만 뺀다 — 같은 키의 다른 소유자를 건드리지 않는다."""
        with self._lock:
            run_ids = self._held.get(key)
            if not run_ids:
                return
            run_ids.discard(run_id)
            if not run_ids:
                del self._held[key]

    def holds(self, key: LockKey, run_id: Optional[str] = None) -> bool:
        with self._lock:
            run_ids = self._held.get(key)
            if not run_ids:
                return False
            return True if run_id is None else run_id in run_ids

    def snapshot(self) -> Dict[LockKey, set]:
        with self._lock:
            return {key: set(run_ids) for key, run_ids in self._held.items()}

    def reset(self) -> None:
        """등록부를 비운다. fork 로 신원이 바뀔 때와 테스트에서만 부른다."""
        with self._lock:
            self._held.clear()


REGISTRY = HeldLockRegistry()


# ── 살아있음 확인 ──


def is_pid_alive(pid: Optional[int]) -> Optional[bool]:
    """이 호스트에서 그 PID 가 살아있나.

    Returns:
        True  — 살아있다
        False — 없다 (ESRCH)
        None  — 판단 못 한다 (pid 가 없거나 이상하다)

    ★`PermissionError` 는 **살아있다**는 뜻이다 — 프로세스가 있는데 다른
     사용자 소유라 신호를 못 보내는 것이다. 죽음으로 읽으면 남의 일을
     빼앗는다 (fail-closed).
    """
    if pid is None or pid <= 0:
        return None
    try:
        os.kill(pid, 0)
    except ProcessLookupError:
        return False
    except PermissionError:
        return True
    except OSError as exc:  # 그 밖의 OS 오류는 「모른다」
        logger.debug("is_pid_alive(%s) OSError: %s", pid, exc)
        return None
    return True


def parse_iso(raw) -> Optional[datetime]:
    """시각 칸을 datetime 으로. 실패하면 None.

    `heartbeat_at` / `cancel_requested_at` 은 timestamptz 라 드라이버가 이미
    datetime 을 준다. `started_at` 같은 옛 칸은 text 다. 둘 다 받는다.

    tz 없는 값은 UTC 로 읽는다 (기존 `_evaluate_running_state` 와 같은 규칙).
    """
    if raw is None or raw == "":
        return None
    if isinstance(raw, datetime):
        parsed = raw
    else:
        try:
            parsed = datetime.fromisoformat(str(raw).replace("Z", "+00:00"))
        except (ValueError, TypeError, AttributeError):
            return None
    if parsed.tzinfo is None:
        parsed = parsed.replace(tzinfo=timezone.utc)
    return parsed


def judge_owner(
    owner: LockOwner,
    *,
    key: Optional[LockKey] = None,
    now: Optional[datetime] = None,
    lease_seconds: int = 90,
    registry: Optional[HeldLockRegistry] = None,
    allow_heartbeat_steal: bool = False,
) -> Tuple[OwnerVerdict, str]:
    """소유자가 살아있나 — 모듈 머리말의 6갈래 판정.

    Args:
        owner: `step_run` 행에서 읽은 신원.
        key: 등록부 조회용. None 이면 1·2번 갈래를 건너뛴다.
        now: 하트비트 나이 계산 기준 (테스트 주입용).
        lease_seconds: 하트비트가 이만큼 안 뛰면 멈춘 것으로 본다.
        registry: 테스트 주입용. None 이면 모듈 전역.
        allow_heartbeat_steal: 하트비트 만료를 **DEAD 로 쓸지**. 기본 False.

    ★`allow_heartbeat_steal` 이 왜 기본 꺼져 있나 — 하트비트 만료는 죽음의
     **증거가 아니라 정황**이다. DB 가 잠깐 끊겨도, 프로세스가 멈춰 있어도
     하트비트는 멈춘다. 그런데 락을 뺏는 순간 원래 일하던 쪽이 살아 돌아오면
     둘이 같은 결과물에 쓴다 — 지금 코드에는 결과물 쓰기를 막는 울타리가
     `checkpoint_gate` 하나뿐이라 안전 지점 사이의 쓰기는 막지 못한다.
     그래서 자동으로 풀기는 **확정 사망**(같은 프로세스 등록부 부재, 로컬 PID
     소멸)에만 허용하고, 하트비트 만료는 `UNKNOWN` 으로 내려 기존 경과 시간
     경로가 판단하게 둔다. 운영자는 그 정황을 `GET .../locks` 에서 보고
     `release` 로 명시 해제할 수 있다.

    Returns:
        (판정, 사람이 읽을 사유)
    """
    reg = registry if registry is not None else REGISTRY
    now = now or datetime.now(timezone.utc)

    def _stalled(age: float, detail: str) -> Tuple[OwnerVerdict, str]:
        """하트비트가 멈췄다 — 켜져 있을 때만 DEAD, 아니면 UNKNOWN."""
        verdict = OwnerVerdict.DEAD if allow_heartbeat_steal else OwnerVerdict.UNKNOWN
        return (
            verdict,
            f"{owner.describe()} 하트비트가 {age:.0f}초 멈췄다 "
            f"(lease={lease_seconds}s{detail}) — "
            f"{'자동으로 풀기' if allow_heartbeat_steal else '자동으로 풀기 안 함(확정 사망 아님)'}",
        )

    # 1·2. 이 프로세스가 소유자라고 적혀 있다 — 등록부가 진실이다.
    if owner.is_this_process:
        if key is None:
            return (
                OwnerVerdict.UNKNOWN,
                f"{owner.describe()} 이 프로세스지만 조회 키가 없다",
            )
        if reg.holds(key, owner.run_id):
            return (
                OwnerVerdict.ALIVE,
                f"{owner.describe()} 이 프로세스가 실제로 잡고 있다",
            )
        return (
            OwnerVerdict.DEAD,
            f"{owner.describe()} 이 프로세스 소유인데 등록부에 없다 "
            f"(일하던 스레드가 사라졌다)",
        )

    heartbeat = parse_iso(owner.heartbeat_at)

    # 3·4. 같은 호스트다 — PID 를 직접 물어본다.
    if owner.host and owner.host == PROCESS_HOST:
        alive = is_pid_alive(owner.pid)
        if alive is False:
            return (
                OwnerVerdict.DEAD,
                f"{owner.describe()} 같은 호스트인데 그 PID 가 없다",
            )
        if alive is True:
            # PID 는 살아있다. 그런데 boot_id 가 이 프로세스가 아니다 —
            # 다른 백엔드거나, 그 PID 를 재사용한 남의 프로세스다.
            # 하트비트로 가른다.
            if heartbeat is None:
                return (
                    OwnerVerdict.UNKNOWN,
                    f"{owner.describe()} PID 는 살아있으나 하트비트가 없다",
                )
            age = (now - heartbeat).total_seconds()
            if age >= lease_seconds:
                # PID 는 살아있는데 하트비트가 멈췄다 — 그 PID 를 재사용한
                # 남의 프로세스이거나, 우리 백엔드가 멈춰 있는 것이다.
                # 어느 쪽인지 확정 못 하므로 자동으로 풀기하지 않는다.
                return _stalled(age, ", PID 는 살아있음")
            return (
                OwnerVerdict.ALIVE,
                f"{owner.describe()} PID 살아있고 하트비트 {age:.0f}초 전",
            )
        # alive is None — pid 를 못 읽었다. 하트비트로 떨어진다.

    # 5. 다른 호스트거나 신원이 부실하다 — 하트비트 lease 만 남았다.
    if heartbeat is not None:
        age = (now - heartbeat).total_seconds()
        if age >= lease_seconds:
            return _stalled(age, "")
        return (
            OwnerVerdict.ALIVE,
            f"{owner.describe()} 하트비트 {age:.0f}초 전",
        )

    # 6. 신원도 하트비트도 없다 — 구 행. 호출자가 경과 시간으로 판단한다.
    return (
        OwnerVerdict.UNKNOWN,
        f"{owner.describe()} 신원·하트비트 둘 다 없다 (경과 시간 판정으로)",
    )


# ── 하트비트 ──

# ★시각은 **DB 시계**로 찍는다 (`CURRENT_TIMESTAMP`). lease 는 두 시각의 차를
#  재는 계약인데, 쓰는 쪽과 읽는 쪽의 시계가 다르면 그 차가 거짓이 된다.
#  `updated_at` 은 기존 text 계약이라 그대로 문자열을 받는다.
_HEARTBEAT_SQL = """
    UPDATE step_run
       SET heartbeat_at = CURRENT_TIMESTAMP, updated_at = :now
     WHERE project_id = :pid AND episode_id = :eid AND step_id = :sid
       AND run_id = :run_id
       AND owner_boot_id = :boot_id
       AND status = 'running'
"""


def beat_once(session_factory, registry: Optional[HeldLockRegistry] = None) -> int:
    """등록부에 있는 락의 `heartbeat_at` 을 한 번 갱신한다.

    ★갱신은 **소유자 일치 조건부**다 (`run_id` + `owner_boot_id` + `running`).
     이미 빼앗긴 락을 되살리지 않는다.

    Returns: 갱신된 행 수.
    """
    from sqlalchemy import text as _sql_text

    reg = registry if registry is not None else REGISTRY
    held = reg.snapshot()
    if not held:
        return 0

    now = datetime.now(timezone.utc).isoformat()
    boot_id = process_boot_id()
    updated = 0
    session = session_factory()
    try:
        for (project_id, episode_id, step_id), run_ids in held.items():
            for run_id in run_ids:
                result = session.execute(_sql_text(_HEARTBEAT_SQL), {
                    "now": now,
                    "pid": project_id,
                    "eid": episode_id,
                    "sid": step_id,
                    "run_id": run_id,
                    "boot_id": boot_id,
                })
                updated += result.rowcount or 0
        session.commit()
    finally:
        session.close()
    return updated


class HeartbeatThread:
    """프로세스당 하나. 주기적으로 `beat_once` 를 부른다.

    ★하트비트 실패는 로그만 남기고 계속한다 — 하트비트가 주행을 죽이면
     안 된다. DB 가 잠깐 끊겼다고 락을 잃는 것보다 lease 가 만료되는 쪽이
     낫고, lease 는 어차피 다시 뛰면 회복된다.

    ★일하는 스레드의 세션을 공유하지 않는다 — 매 회 새 세션을 열고 닫는다.
     SQLAlchemy 세션은 스레드 안전하지 않다.
    """

    def __init__(
        self,
        session_factory,
        *,
        interval_seconds: int = 30,
        registry: Optional[HeldLockRegistry] = None,
    ) -> None:
        self._session_factory = session_factory
        self._interval = interval_seconds
        self._registry = registry if registry is not None else REGISTRY
        self._stop = threading.Event()
        self._thread: Optional[threading.Thread] = None

    def start(self) -> None:
        if self._thread is not None and self._thread.is_alive():
            return
        self._stop.clear()
        self._thread = threading.Thread(
            target=self._loop, name="step-lock-heartbeat", daemon=True
        )
        self._thread.start()
        logger.info(
            "스텝 락 하트비트 시작 (간격=%ds, boot_id=%s)",
            self._interval, process_boot_id()[:8],
        )

    def stop(self, timeout: float = 5.0) -> None:
        self._stop.set()
        if self._thread is not None:
            self._thread.join(timeout=timeout)

    def _loop(self) -> None:
        while not self._stop.wait(self._interval):
            try:
                beat_once(self._session_factory, self._registry)
            except Exception as exc:  # 하트비트가 주행을 죽이지 않는다
                logger.warning("스텝 락 하트비트 실패 (계속 진행): %s", exc)


_HEARTBEAT: Optional[HeartbeatThread] = None


def start_heartbeat(session_factory, *, interval_seconds: int = 30) -> HeartbeatThread:
    """프로세스 전역 하트비트를 켠다 (여러 번 불러도 하나)."""
    global _HEARTBEAT
    if _HEARTBEAT is None:
        _HEARTBEAT = HeartbeatThread(
            session_factory, interval_seconds=interval_seconds
        )
    _HEARTBEAT.start()
    return _HEARTBEAT


def stop_heartbeat() -> None:
    global _HEARTBEAT
    if _HEARTBEAT is not None:
        _HEARTBEAT.stop()


# ── 기동 시 자기 락 풀기 ──

_RECLAIM_SELECT = """
    SELECT project_id, episode_id, step_id, run_id,
           owner_host, owner_pid, owner_boot_id, heartbeat_at, started_at
      FROM step_run
     WHERE status = 'running'
       AND owner_host = :host
"""

# ★풀기할 때 **소유 토큰을 버린다** (`run_id` 를 새 값으로).
#
#  상태만 바꾸고 `run_id` 를 두면, 소유권 검사가 `run_id` 일치 하나뿐이라
#  (`step_runner._update_step_run` 의 owner_clause) 옛 worker 의 갱신이
#  그대로 통과해 풀기한 락을 되살린다. 여기서는 PID 사망을 확인했으니
#  되살아날 worker 가 없어야 정상이지만, 확인이 틀렸을 때의 폭발 반경을
#  0으로 만드는 값이 한 칸이므로 같이 버린다.
_RECLAIM_UPDATE = """
    UPDATE step_run
       SET status = 'failed',
           run_id = :revoked_run_id,
           error_message = :reason,
           last_recovery_reason = :reason,
           recovery_count = COALESCE(recovery_count, 0) + 1,
           completed_at = :now,
           updated_at = :now
     WHERE project_id = :pid AND episode_id = :eid AND step_id = :sid
       AND status = 'running'
       AND run_id = :run_id
"""


def reclaim_dead_locks_on_startup(db) -> int:
    """백엔드가 뜰 때, **이 호스트의 죽은 프로세스**가 잡은 락을 푼다.

    풀기 = `running` → `failed`. **행을 지우지 않는다** — 죽은 것을 죽었다고
    적는 것이다. 다음 `resume` 이 바로 가져간다.

    ★이 하나가 「서버를 죽였는데 그 서버가 잡은 락이 남아 한 시간을
     기다리는」 사고를 없앤다. 다시 띄우는 순간 풀린다.

    ★대상은 **`owner_host` 가 이 호스트인 행뿐**이다. 신원 칸이 NULL 인
     구 행과 다른 호스트의 행은 건드리지 않는다 — 자동 대량 정리 금지.

    Returns: 풀기한 행 수.
    """
    from sqlalchemy import text as _sql_text

    rows = db.execute(_sql_text(_RECLAIM_SELECT), {"host": PROCESS_HOST}).fetchall()
    if not rows:
        return 0

    now = datetime.now(timezone.utc).isoformat()
    reclaimed = 0
    for row in rows:
        owner = LockOwner.from_row({
            "owner_host": row.owner_host,
            "owner_pid": row.owner_pid,
            "owner_boot_id": row.owner_boot_id,
            "heartbeat_at": row.heartbeat_at,
            "run_id": row.run_id,
        })

        # 기동 시점에 이 프로세스가 잡고 있는 락은 있을 수 없다 —
        # boot_id 가 이 프로세스면 그건 재사용된 uuid 가 아니라 버그다.
        if owner.is_this_process:
            logger.error(
                "startup reclaim: %s/%s 가 이 프로세스 boot_id 를 갖고 있다 "
                "— 건너뛴다 (조사 필요)",
                row.episode_id, row.step_id,
            )
            continue

        alive = is_pid_alive(owner.pid)
        if alive is not False:
            # 살아있거나(다른 백엔드) 판단 못 함 — 안 건드린다.
            logger.info(
                "startup reclaim 건너뜀: %s/%s %s (pid_alive=%s)",
                row.episode_id, row.step_id, owner.describe(), alive,
            )
            continue

        reason = (
            f"startup reclaim: 소유 프로세스가 없다 ({owner.describe()}, "
            f"started_at={row.started_at}, revoked_run_id={row.run_id})"
        )
        result = db.execute(_sql_text(_RECLAIM_UPDATE), {
            "reason": reason[:2000],
            "revoked_run_id": f"reclaimed-{uuid.uuid4()}",
            "now": now,
            "pid": row.project_id,
            "eid": row.episode_id,
            "sid": row.step_id,
            "run_id": row.run_id,
        })
        if result.rowcount:
            reclaimed += 1
            logger.warning(
                "[STARTUP-RECLAIM] %s/%s 풀기 — %s",
                row.episode_id, row.step_id, reason,
            )
    db.commit()
    return reclaimed


__all__ = [
    "PROCESS_HOST",
    "process_identity",
    "process_boot_id",
    "process_pid",
    "OwnerVerdict",
    "LockOwner",
    "LockKey",
    "HeldLockRegistry",
    "REGISTRY",
    "is_pid_alive",
    "parse_iso",
    "judge_owner",
    "beat_once",
    "HeartbeatThread",
    "start_heartbeat",
    "stop_heartbeat",
    "reclaim_dead_locks_on_startup",
]
