#!/usr/bin/env python3
"""reve 2.1 판 주행 관찰 — DB 상태 + Opik 기록을 한 화면에.

왜 둘 다 보는가: 스텝이 completed 인 것과 **모델에 실제로 나간 것**은 다르다.
기록이 빠지면 「0건」이 「안 돌았다」로 읽히고, 반대로 스텝이 열려 있는데
기록만 쌓이면 헛도는 것이다. 그래서 **DB 시각과 Opik 시각을 나란히** 찍는다.

★시각은 전부 KST 로 적는다(DB·Opik 은 UTC 로 저장된다).
★이번 판 수정이 실제로 발화했는지도 같이 본다 — 설정만 켜고 확인 안 하면
 「켠 줄 알았는데 안 실린」 판을 그대로 완주한다.

usage:
  backend/.venv/bin/python scratchpad/reve_run_watch.py            # 요약
  backend/.venv/bin/python scratchpad/reve_run_watch.py --spans    # 스텝별 발송량
"""
import argparse
import datetime as dt
import json
import sys
from collections import Counter, defaultdict
from pathlib import Path

sys.path.insert(0, "/Users/manta/Documents/Projects/TheRoad-I1/scratchpad")
import _opik_env  # noqa: E402  ★cwd 를 backend 로 고정 (설정을 읽는 유일한 조건)

ROOT = Path("/Users/manta/Documents/Projects/TheRoad-I1")
STATE = json.loads((ROOT / "scratchpad" / "minimal_e2e_state.json")
                   .read_text())
PID, EID = STATE["project_id"], STATE["episode_id"]
KST = dt.timezone(dt.timedelta(hours=9))


def kst(ts) -> str:
    """UTC 문자열/시각 → KST 표기. 못 읽으면 원문 그대로."""
    if not ts:
        return "-"
    try:
        s = str(ts).replace("Z", "+00:00")
        d = dt.datetime.fromisoformat(s)
        if d.tzinfo is None:
            d = d.replace(tzinfo=dt.timezone.utc)
        return d.astimezone(KST).strftime("%m-%d %H:%M:%S")
    except Exception:  # noqa: BLE001
        return str(ts)[:19]


def db_rows(sql, params=()):
    import psycopg2

    with psycopg2.connect(host="localhost", user="theroad",
                          password="theroad_dev_2026",
                          dbname="theroad") as conn:
        with conn.cursor() as cur:
            cur.execute(sql, params)
            return cur.fetchall()


def show_steps():
    rows = db_rows(
        "SELECT status, count(*) FROM step_run WHERE episode_id=%s "
        "GROUP BY status ORDER BY 2 DESC", (EID,))
    total = sum(n for _, n in rows)
    print(f"── 스텝 ({total}건 기록됨)")
    for st, n in rows:
        print(f"   {st:16s} {n}")
    cur = db_rows(
        "SELECT step_id, status, started_at, completed_at, "
        "completed_count, applicable_count, failed_count, error_message "
        "FROM step_run WHERE episode_id=%s AND status NOT IN "
        "('completed','not_applicable','skipped') "
        "ORDER BY started_at DESC NULLS LAST LIMIT 5", (EID,))
    if cur:
        print("   진행/미완:")
        for r in cur:
            err = (r[7] or "")[:80]
            print(f"     {r[0]:32s} {r[1]:10s} {kst(r[2])} "
                  f"{r[4]}/{r[5]} 실패{r[6]} {err}")
    last = db_rows(
        "SELECT step_id, completed_at FROM step_run WHERE episode_id=%s "
        "AND completed_at IS NOT NULL ORDER BY completed_at DESC LIMIT 3",
        (EID,))
    for r in last:
        print(f"   최근 완료  {r[0]:32s} {kst(r[1])}")


def show_calls():
    rows = db_rows(
        "SELECT model_name, status, count(*), max(created_at) "
        "FROM llm_call_log WHERE episode_id=%s GROUP BY 1,2 "
        "ORDER BY 3 DESC LIMIT 12", (EID,))
    if not rows:
        # ★「0건 = 기록이 샌다」로 읽지 말 것. 이 표는 **이미지·직접 호출**
        #  경로만 담는다(`log_llm_call` 을 부르는 것은 image client 와
        #  gemini_text/openai client 뿐). litellm 을 타는 텍스트 스텝은
        #  Opik 에만 남으므로, 분석 단계에서 여기가 비는 것이 정상이다.
        print("── DB llm_call_log 0건 — 이미지 단계 전이면 정상"
              "(텍스트 스텝은 litellm→Opik 경로)")
        return
    print("── 이미지·직접 호출 (DB llm_call_log)")
    for m, st, n, ts in rows:
        print(f"   {str(m)[:38]:38s} {st:8s} {n:4d}  최근 {kst(ts)}")
    errs = db_rows(
        "SELECT model_name, error_message, created_at FROM llm_call_log "
        "WHERE episode_id=%s AND status='error' "
        "ORDER BY created_at DESC LIMIT 5", (EID,))
    for m, e, ts in errs:
        print(f"   ✘ {kst(ts)} {str(m)[:28]:28s} {str(e)[:110]}")


def show_opik(since, detail=False):
    from tools.opik_prompt_audit.audit.fetch import fetch_spans, step_of_call

    base, ws, proj = _opik_env.opik_target()
    print(f"── Opik {base} · 프로젝트 {proj} · {kst(since)} 이후")
    try:
        spans = fetch_spans(base, ws, proj, since)
    except Exception as exc:  # noqa: BLE001
        print(f"   ✘ 조회 실패 — {type(exc).__name__}: {exc}")
        return
    if not spans:
        print("   기록 0건")
        return
    per = Counter()
    chars = defaultdict(int)
    errs = []
    for sp in spans:
        step = step_of_call(sp, {}) or "(태그 없음)"
        per[step] += 1
        inp = sp.get("input") or {}
        chars[step] += len(json.dumps(inp, ensure_ascii=False))
        if sp.get("error_info"):
            errs.append((sp.get("start_time"), step,
                         json.dumps(sp["error_info"],
                                    ensure_ascii=False)[:140]))
    print(f"   span {len(spans)}건 · 스텝 {len(per)}종")
    order = sorted(per.items(), key=lambda kv: -chars[kv[0]])
    for step, n in (order if detail else order[:10]):
        print(f"   {step[:34]:34s} {n:4d}회  {chars[step]:>9,}자")
    ts = [sp.get("start_time") for sp in spans if sp.get("start_time")]
    if ts:
        print(f"   기록 구간  {kst(min(ts))} ~ {kst(max(ts))}")
    for t, step, e in errs[:6]:
        print(f"   ✘ {kst(t)} {step[:26]:26s} {e}")


def show_this_round():
    """이번 판 수정이 **실제로 발화했는가** — 설정만 켜고 안 보면 모른다."""
    from app.core.config import settings
    from app.modules.pipeline.cine_provider import cine_provider_identity
    from app.modules.pipeline.still_recipe import (
        cine_pack_selector, resolve_prompt_version,
    )

    print("── 이번 판 설정 (도는 프로세스와 같은 .env)")
    print(f"   provider      {cine_provider_identity()}")
    print(f"   cine 팩       {resolve_prompt_version(cine_pack_selector())}")
    print(f"   연출 재료     {settings.still_cine_stage_direction_enabled}")
    print(f"   산출 관문     {settings.still_cine_verify_enabled}")
    rec = (ROOT / "projects" / PID / "images" / EID / "scene" / "recipe"
           / "records.json")
    if not rec.exists():
        print("   records.json 아직 없음 (이미지 단계 전)")
        return
    data = json.loads(rec.read_text())
    cine = {k: v for k, v in data.items() if k.endswith("::cine")}
    if not cine:
        print("   변환 기록 아직 없음")
        return
    prov = Counter(str(v.get("provider") or "(없음)") for v in cine.values())
    applied = sum(1 for v in cine.values() if v.get("applied"))
    rejected = sum(1 for v in cine.values() if v.get("rejected"))
    pending = sum(1 for v in cine.values() if v.get("pending"))
    print(f"   변환 {len(cine)}건 — 적용 {applied} · 관문 기각 {rejected} "
          f"· 접수 미결 {pending} · 제공자 {dict(prov)}")
    for k, v in list(cine.items())[:3]:
        print(f"     {k:16s} fp v{v.get('fingerprint_version','1')} "
              f"{v.get('endpoint','-')} {str(v.get('rejected_reason') or '')[:40]}")


def main():
    ap = argparse.ArgumentParser(description=__doc__)
    ap.add_argument("--since", default="",
                    help="Opik 조회 시작 (기본: 오늘 10:00 UTC = 19:00 KST)")
    ap.add_argument("--spans", action="store_true", help="스텝 전부 펼침")
    args = ap.parse_args()
    since = args.since or dt.datetime.now(dt.timezone.utc).replace(
        hour=10, minute=0, second=0, microsecond=0).isoformat()

    print(f"═══ reve 2.1 판 · pid={PID[:8]} eid={EID[:8]} · "
          f"지금 {dt.datetime.now(KST).strftime('%m-%d %H:%M:%S')} KST")
    show_this_round()
    show_steps()
    show_calls()
    show_opik(since, detail=args.spans)
    return 0


if __name__ == "__main__":
    raise SystemExit(main())
