#!/usr/bin/env python3
"""reve 2.1 판 주행 감시 — **달라진 것만** 한 줄로 뱉는다.

Monitor 의 이벤트 원천으로 쓴다. 매 바퀴 전부 찍으면 알림이 홍수가 되고,
반대로 좋은 소식만 찍으면 **죽은 것과 조용한 것이 같아 보인다.** 그래서
「진행이 한 칸 나아갔다」와 「무언가 잘못됐다」를 **둘 다** 뱉되, 직전 바퀴와
같은 것은 안 뱉는다.

보는 것:
  · 러너 로그의 새 줄 중 실패·정지·완주·건너뜀
  · 스텝 상태 변화 (완료 수 · 실패 수)
  · Opik 새 span 중 **오류**, 그리고 **스텝 태그가 안 붙은 것**
    (과제 #58 chat.completion 미분류 — 몇 건인지 세어야 고칠 값이 나온다)
  · 최종 변환 기록 — 제공자·관문 기각·접수 미결

usage: reve_watch_probe.py [--interval 180]
"""
import argparse
import datetime as dt
import json
import sys
import time
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")
RUNLOG = ROOT / "scratchpad" / "reve_e2e_run.log"
UVLOG = ROOT / "backend" / "uvicorn.log"

# ★백엔드 로그에서 **이번 판이 겨누는 사건**만 집는다. 러너 로그에는 이
#  줄들이 안 나오므로, 이것을 안 보면 마커 맵 후퇴가 발화했는지도, reve 첫
#  호출이 나갔는지도 모른 채 「조용하다」로 읽게 된다.
BACKEND_MARKERS = (
    "마커 맵 검사 소진",        # 콘티 lane 마커 맵 후퇴
    "맵 재투영이",              # 장소 마커 맵 후퇴
    "marker map check 소진",    # 사다리를 다 쓰고도 실패
    "cine_transform",           # 최종 변환 (reve)
    "접수한 변환",              # fal queue 회수
    "불변축을 바꿔 기각",       # 산출 관문 기각
    "검열로 포기",
)
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))

# 러너 로그에서 사람이 반응해야 하는 줄. ★성공만 보면 크래시가 침묵과
#  같아 보인다 — 실패 신호를 함께 넣는다.
RUN_MARKERS = ("FAILED", "TIMEOUT", "STOP at", "ALL STEPS DONE",
               "disabled", "retrying", "poll error")


def now() -> str:
    return dt.datetime.now(KST).strftime("%H:%M:%S")


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 step_snapshot():
    rows = db_rows(
        "SELECT status, count(*) FROM step_run WHERE episode_id=%s "
        "GROUP BY status", (EID,))
    return {st: n for st, n in rows}


def cine_snapshot():
    rec = (ROOT / "projects" / PID / "images" / EID / "scene" / "recipe"
           / "records.json")
    if not rec.exists():
        return None
    try:
        data = json.loads(rec.read_text())
    except Exception:  # noqa: BLE001 — 쓰는 중일 수 있다
        return None
    cine = {k: v for k, v in data.items() if k.endswith("::cine")}
    if not cine:
        return None
    return {
        "n": len(cine),
        "applied": sum(1 for v in cine.values() if v.get("applied")),
        "rejected": sum(1 for v in cine.values() if v.get("rejected")),
        "declined": sum(1 for v in cine.values() if v.get("declined")),
        "pending": sum(1 for v in cine.values() if v.get("pending")),
        "providers": sorted({str(v.get("provider") or "?")
                             for v in cine.values()}),
    }


def opik_snapshot(since):
    """★미분류는 **축 태그로만** 잰다.

    처음에 `step_of_call(sp, {})` 가 비었는지로 셌는데 그 값은 **절대 비지
    않는다** — 마지막 갈래가 span 이름으로 떨어지고 litellm span 이름
    (`{model}_{obj}_{created}`)은 언제나 있다. 그래서 「0건」이 나왔고 그것을
    「미분류가 없다」로 읽을 뻔했다. 재는 도구가 덜 잡으면 0 이 「없다」로
    읽힌다.

    제대로 재려면 `step:` 축 태그가 **자기에게도 부모 trace 에도** 없는
    것을 세야 한다 — 그래야 과제 #58 이 고칠 값이 나온다.
    """
    from tools.opik_prompt_audit.audit.fetch import (
        _step_tag, fetch_spans, fetch_trace_index,
    )

    base, ws, proj = _opik_env.opik_target()
    spans = fetch_spans(base, ws, proj, since)
    idx = fetch_trace_index(base, ws, proj, since)
    untagged, orphan, errs = 0, 0, []
    names = {}
    for sp in spans:
        if not _step_tag(sp):
            parent = idx.get(str(sp.get("trace_id")))
            if parent is None:
                orphan += 1          # 부모를 아예 못 찾는다 — 계보 끊김
            if parent is None or not _step_tag(parent):
                untagged += 1
                nm = str(sp.get("name") or "?")
                # 이름의 변동부(타임스탬프)를 뺀 모양으로 묶는다
                key = "_".join(nm.split("_")[:2]) if "_" in nm else nm
                names[key] = names.get(key, 0) + 1
        if sp.get("error_info"):
            errs.append((sp.get("start_time"), sp.get("name") or "?",
                         json.dumps(sp["error_info"],
                                    ensure_ascii=False)[:120]))
    return {"n": len(spans), "untagged": untagged, "orphan": orphan,
            "errs": errs, "kinds": sorted(names.items(),
                                          key=lambda kv: -kv[1])[:4]}


def main():
    ap = argparse.ArgumentParser(description=__doc__)
    ap.add_argument("--interval", type=int, default=180)
    args = ap.parse_args()

    since = dt.datetime.now(dt.timezone.utc).replace(
        hour=10, minute=0, second=0, microsecond=0).isoformat()
    seen_lines = 0
    # 백엔드 로그는 8MB 가 넘는다 — 매 바퀴 통째로 읽지 않고 **자란 만큼만**
    # 읽는다. 시작 시점의 끝에서 출발하므로 지난 줄을 다시 뱉지 않는다.
    uv_pos = UVLOG.stat().st_size if UVLOG.exists() else 0
    prev_steps, prev_cine, prev_untagged = {}, None, -1
    seen_errs = set()
    print(f"[{now()}] 감시 시작 — {args.interval}s 간격 · "
          f"eid={EID[:8]}", flush=True)

    while True:
        # ── 러너 로그: 반응이 필요한 새 줄만
        try:
            lines = RUNLOG.read_text(errors="replace").splitlines()
        except FileNotFoundError:
            lines = []
        for ln in lines[seen_lines:]:
            if any(m in ln for m in RUN_MARKERS):
                print(f"[{now()}] 러너 · {ln.strip()[:200]}", flush=True)
        seen_lines = len(lines)

        # ── 백엔드 로그: 이번 판이 겨누는 사건만 (자란 만큼만 읽는다)
        try:
            size = UVLOG.stat().st_size
            if size < uv_pos:      # 로그가 갈렸다 — 처음부터 다시
                uv_pos = 0
            if size > uv_pos:
                with UVLOG.open("r", errors="replace") as fh:
                    fh.seek(uv_pos)
                    chunk = fh.read()
                    uv_pos = fh.tell()
                for ln in chunk.splitlines():
                    if any(m in ln for m in BACKEND_MARKERS):
                        # json 로그면 message 만 꺼내 읽기 좋게
                        msg = ln
                        if ln.startswith("{"):
                            try:
                                msg = json.loads(ln).get("message") or ln
                            except Exception:  # noqa: BLE001
                                pass
                        print(f"[{now()}] 백엔드 · {msg.strip()[:220]}",
                              flush=True)
        except FileNotFoundError:
            pass

        # ── 스텝 진행: 칸이 움직였을 때만
        try:
            cur = step_snapshot()
        except Exception as exc:  # noqa: BLE001
            print(f"[{now()}] ✘ DB 조회 실패 {type(exc).__name__}: {exc}",
                  flush=True)
            cur = prev_steps
        if cur != prev_steps:
            done = cur.get("completed", 0)
            na = cur.get("not_applicable", 0) + cur.get("skipped", 0)
            bad = cur.get("failed", 0)
            print(f"[{now()}] 스텝 완료 {done} · 해당없음 {na} · "
                  f"실패 {bad}", flush=True)
            prev_steps = cur

        # ── 최종 변환: 이번 판의 핵심 — 제공자·관문이 실제로 발화하는가
        c = cine_snapshot()
        if c and c != prev_cine:
            print(f"[{now()}] 변환 {c['n']}건 — 적용 {c['applied']} · "
                  f"관문 기각 {c['rejected']} · 검열 포기 {c['declined']} · "
                  f"접수 미결 {c['pending']} · 제공자 {c['providers']}",
                  flush=True)
            prev_cine = c

        # ── Opik: 오류는 매건, 미분류는 수가 늘 때만
        try:
            o = opik_snapshot(since)
        except Exception as exc:  # noqa: BLE001
            print(f"[{now()}] ✘ Opik 조회 실패 {type(exc).__name__}: {exc}",
                  flush=True)
            o = None
        if o:
            for t, name, e in o["errs"]:
                key = f"{t}|{name}"
                if key not in seen_errs:
                    seen_errs.add(key)
                    print(f"[{now()}] ✘ Opik 오류 · {str(name)[:40]} · {e}",
                          flush=True)
            if o["untagged"] != prev_untagged:
                kinds = " ".join(f"{k}×{v}" for k, v in o["kinds"])
                print(f"[{now()}] Opik span {o['n']}건 · 축 태그 없음 "
                      f"{o['untagged']}건(부모 없음 {o['orphan']}) {kinds}",
                      flush=True)
                prev_untagged = o["untagged"]

        if lines and "ALL STEPS DONE" in lines[-1]:
            print(f"[{now()}] 완주 — 감시를 끝낸다", flush=True)
            return 0
        time.sleep(args.interval)


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