"""run-scoped state 의 file-based 저장/조회.

설계 원칙:
  - 한 run 의 입력 (`task-type`, `brief-path`, `directive`, `workers`, `models`,
    `related-tasks`, `approved-plan`, `clarification-response`) 은
    `<run-dir>/manifests/run-inputs-<seq>.json` 에 박힌다.
  - 한 run 의 계산된 paths/seqs/timestamp 는
    `<run-dir>/manifests/run-context-<seq>.json` 에 박힌다.
  - 두 파일이 한 번 디스크에 자리잡으면, 이후 모든 reader 는 환경 변수에 의존
    하지 않고 같은 값을 돌려준다 → 같은 claude 세션의 두 병렬 호출이 stale
    snapshot 을 보지 않는다.

부수효과 함수는 모두 per-task 락 (`~/.okstra/.locks/<task-key>.lock`) 안에서만
seq advance 와 파일 쓰기를 수행한다.
"""
from __future__ import annotations

import fcntl
import json
import os
from contextlib import contextmanager
from datetime import datetime, timezone
from pathlib import Path
from typing import Iterator, Optional

from okstra_project.dirs import okstra_home

from .path_hints import compact_run_context, hydrate_run_context
from .json_boundary import JsonBoundaryError, load_owned_object, write_owned_object_atomic
from .paths import compute_run_paths, task_runs_dir


def latest_run_inputs(
    project_root: object,
    task_group: str,
    task_id: str,
    *,
    phase_segment: str = "",
) -> dict:
    """가장 최근 ``run-inputs-*.json`` 의 ``inputs`` dict 를 반환 (부재 시 {}).

    ``phase_segment`` 가 주어지면 그 phase 의 ``runs/<seg>/manifests/`` 만 스캔한다
    (resume 는 같은 phase 를 재실행하므로 그 phase 의 직전 입력이 기준). 비어 있으면
    task 의 모든 phase 를 통틀어 가장 최근(mtime) 파일을 고른다.

    페이로드 구조는 ``{"schemaVersion", "writtenAt", "inputs": {...}}`` 이며 실제
    입력은 top-level 이 아니라 ``inputs`` 아래에 있다. run-inputs 의 위치/구조 지식은
    이 함수 한 곳에만 둔다(wizard·okstra.sh resume 가 공유).
    """
    if not (project_root and task_group and task_id):
        return {}
    runs_base = task_runs_dir(Path(project_root), task_group, task_id)
    if not runs_base.is_dir():
        return {}
    if phase_segment:
        glob_iter = (runs_base / phase_segment / "manifests").glob(
            "run-inputs-*.json")
    else:
        glob_iter = runs_base.glob("*/manifests/run-inputs-*.json")
    candidates: list[tuple[float, Path]] = []
    for inp in glob_iter:
        try:
            candidates.append((inp.stat().st_mtime, inp))
        except OSError:
            continue
    if not candidates:
        return {}
    candidates.sort(key=lambda x: -x[0])
    try:
        data = load_owned_object(candidates[0][1], artifact="run context")
    except JsonBoundaryError:
        return {}
    inputs = data.get("inputs")
    return inputs if isinstance(inputs, dict) else {}


def _now_iso() -> str:
    return datetime.now(timezone.utc).strftime("%Y-%m-%dT%H:%M:%SZ")


def _now_task_date() -> str:
    return datetime.now(timezone.utc).strftime("%Y-%m-%d %H:%M:%S")


def _task_lock_path(task_key: str) -> Path:
    """task-key 별 mutex 파일. central index 와는 별개로 task 단위 직렬화."""
    home = okstra_home()
    locks = home / ".locks"
    locks.mkdir(parents=True, exist_ok=True)
    safe = task_key.replace("/", "_").replace(":", "_")
    return locks / f"{safe}.lock"


@contextmanager
def task_mutex(task_key: str) -> Iterator[None]:
    """task-key per-process mutex. 동시 호출은 락 안에서 직렬화된다."""
    path = _task_lock_path(task_key)
    path.touch(exist_ok=True)
    with path.open("r+") as f:
        fcntl.flock(f.fileno(), fcntl.LOCK_EX)
        try:
            yield
        finally:
            fcntl.flock(f.fileno(), fcntl.LOCK_UN)


@contextmanager
def dir_flock(dir_path: Path, lock_filename: str) -> Iterator[None]:
    """dir_path 아래 lock_filename 파일 기반 exclusive flock.

    lock 은 보호 대상 파일과 같은 디렉토리에 두어 디렉토리마다 1:1 로
    격리한다 (마지막 세그먼트만 키로 쓰면 다른 task 의 동일 seq 가 같은
    lock 을 공유하므로 금지)."""
    dir_path.mkdir(parents=True, exist_ok=True)
    path = dir_path / lock_filename
    path.touch(exist_ok=True)
    with path.open("r+") as f:
        fcntl.flock(f.fileno(), fcntl.LOCK_EX)
        try:
            yield
        finally:
            fcntl.flock(f.fileno(), fcntl.LOCK_UN)


@contextmanager
def consumers_mutex(plan_run_root: Path) -> Iterator[None]:
    """plan run-root 별 consumers.jsonl append mutex."""
    with dir_flock(plan_run_root, ".consumers.lock"):
        yield


def _atomic_write_json(path: Path, payload: dict) -> None:
    write_owned_object_atomic(path, payload, artifact="run context")


def _run_context_filename(task_type_segment: str, seq: str) -> str:
    return f"run-context-{task_type_segment}-{seq}.json"


def _run_inputs_filename(task_type_segment: str, seq: str) -> str:
    return f"run-inputs-{task_type_segment}-{seq}.json"


def compute_and_write_run_context(
    *,
    workspace_root: Path,
    project_root: Path,
    project_id: str,
    task_group: str,
    task_id: str,
    task_type: str,
    run_seq_override: Optional[int] = None,
    stage: Optional[int] = None,
) -> dict:
    """task per-mutex 안에서 run paths 를 계산하고 디스크에 박는다.

    기존 next_run_seq 가 디스크 스캔 기반이라 두 호출이 동시에 진행되면 같은
    seq 를 받을 수 있다. mutex 안에서 (a) seq 계산 (b) run-context.json 저장
    이 한 트랜잭션으로 묶이면, 다음 호출은 새로 박힌 파일을 보고 다음 seq 를
    돌려준다. 락은 OKSTRA_HOME/.locks/<task-key>.lock.

    반환값: paths dict + 부가 메타(`runTimestamp`, 저장된 파일 경로 등).
    """
    project_root = Path(project_root)
    workspace_root = Path(workspace_root)
    task_key = f"{project_id}:{task_group}:{task_id}"

    with task_mutex(task_key):
        ctx = compute_run_paths(
            project_root=project_root,
            workspace_root=workspace_root,
            project_id=project_id,
            task_group=task_group,
            task_id=task_id,
            task_type=task_type,
            run_seq_override=run_seq_override,
            stage=stage,
        )
        ctx["RUN_TIMESTAMP_ISO"] = _now_iso()
        ctx["TASK_DATE"] = _now_task_date()
        run_manifests_dir = Path(ctx["RUN_MANIFESTS_DIR"])
        ctx_path = run_manifests_dir / _run_context_filename(
            ctx["TASK_TYPE_SEGMENT"], ctx["RUN_MANIFESTS_SEQ"])
        ctx["RUN_CONTEXT_FILE"] = str(ctx_path)
        ctx["RUN_CONTEXT_RELATIVE_PATH"] = (
            str(ctx_path.relative_to(project_root.resolve()))
            if ctx_path.resolve().is_relative_to(project_root.resolve())
            else str(ctx_path))
        _atomic_write_json(ctx_path, compact_run_context(ctx))
    return ctx


def refresh_run_context_snapshot(ctx: dict) -> None:
    """Rewrite a run context after prepare resolves run-scoped inputs."""
    payload = compact_run_context(ctx)
    payload["analysis"] = {
        "sourceCommit": ctx.get("ANALYSIS_SOURCE_COMMIT", ""),
        "target": json.loads(ctx.get("ANALYSIS_TARGET_JSON", "{}")),
        "evidenceInputs": json.loads(ctx.get("EVIDENCE_INPUTS_JSON", "[]")),
    }
    payload["ANALYSIS_SOURCE_COMMIT"] = ctx.get("ANALYSIS_SOURCE_COMMIT", "")
    payload["ANALYSIS_TARGET_JSON"] = ctx.get("ANALYSIS_TARGET_JSON", "{}")
    payload["EVIDENCE_INPUTS_JSON"] = ctx.get("EVIDENCE_INPUTS_JSON", "[]")
    _atomic_write_json(Path(ctx["RUN_CONTEXT_FILE"]), payload)


def write_run_inputs(
    *,
    project_root: Path,
    run_manifests_dir: Path,
    task_type_segment: str,
    seq: str,
    inputs: dict,
) -> Path:
    """사용자 입력(brief, directive, workers, models, ...) 을 run-inputs 파일에
    박는다. 호출자가 미리 계산한 seq 를 그대로 사용한다(같은 트랜잭션 내).

    inputs schema (모든 키 optional):
      taskBriefPath, directive, workers, leadModel, claudeModel, codexModel,
      antigravityModel, reportWriterModel, relatedTasks, approvedPlanPath,
      clarificationResponsePath, analysisTarget, evidenceInputs, renderOnly
    """
    run_manifests_dir = Path(run_manifests_dir)
    path = run_manifests_dir / _run_inputs_filename(task_type_segment, seq)
    payload = {
        "schemaVersion": "1.0",
        "writtenAt": _now_iso(),
        "inputs": dict(inputs),
    }
    _atomic_write_json(path, payload)
    return path


def read_run_context(run_manifests_dir: Path, task_type_segment: str,
                     seq: str) -> Optional[dict]:
    """run-context-<task-type>-<seq>.json 을 dict 로 돌려준다. 부재 시 None."""
    path = Path(run_manifests_dir) / _run_context_filename(task_type_segment, seq)
    if not path.is_file():
        return None
    payload = load_owned_object(path, artifact="run context")
    context = hydrate_run_context(payload)
    analysis = payload.get("analysis")
    if isinstance(analysis, dict):
        target_json = payload.get("ANALYSIS_TARGET_JSON")
        if not isinstance(target_json, str):
            target_json = json.dumps(analysis.get("target", {}), ensure_ascii=False)
        evidence_json = payload.get("EVIDENCE_INPUTS_JSON")
        if not isinstance(evidence_json, str):
            evidence_json = json.dumps(
                analysis.get("evidenceInputs", []), ensure_ascii=False,
            )
        context.update({
            "ANALYSIS_SOURCE_COMMIT": str(
                payload.get("ANALYSIS_SOURCE_COMMIT", analysis.get("sourceCommit", ""))
            ),
            "ANALYSIS_TARGET_JSON": target_json,
            "EVIDENCE_INPUTS_JSON": evidence_json,
        })
    return context


def read_run_inputs(run_manifests_dir: Path, task_type_segment: str,
                    seq: str) -> Optional[dict]:
    """run-inputs-<task-type>-<seq>.json 을 dict 로. 부재 시 None."""
    path = Path(run_manifests_dir) / _run_inputs_filename(task_type_segment, seq)
    if not path.is_file():
        return None
    return load_owned_object(path, artifact="run inputs")
