# -*- coding: utf-8 -*-
"""task-2945 — 봇 spawn(기동) 지연 실측 harness.

task-2942 가 dispatch 반환 직전 spawn 검증을 결선했으나, 총 대기 25초
(cron offset 10 + marker 예산 15)가 **실측 근거 없이** 정해졌다. 본 harness 는
더미 dispatch 를 반복 발사해 `dispatch() 호출 시각 ~ spawn-confirmed marker
생성 시각` 의 실제 분포를 수집한다.

측정 전용. 대상 코드(dispatch/spawn_verification)는 일절 수정하지 않는다.

관측 앵커
  t0        : dispatch() 호출 직전 (wall clock)
  t_ret     : dispatch() 반환 (verify 가 블로킹한 만큼 포함)
  t_marker  : watcher 스레드가 marker 파일을 **최초로 목격**한 시각
              (dispatch 가 25초에 포기해도 watcher 는 계속 관측하므로
               '진짜 지연' 과 '오판' 을 동시에 얻는다)

usage:
    python3 scripts/measure/task2945_spawn_latency.py --probes 8 --grace 240
"""
from __future__ import annotations

import argparse
import json
import os
import sys
import threading
import time
from datetime import datetime
from pathlib import Path

WORKSPACE = Path("/home/jay/workspace")
EVENTS_DIR = WORKSPACE / "memory" / "events"

# dev6(perun) = spawn 고장으로 사용금지(상시 규칙). dev1 = 본 harness 실행 주체(자기충돌 회피).
# dev2 = 파일럿 시점에 10:47 세션 점유 중(→ 스케줄 무기한 대기) 이므로 로테이션에서 제외.
BOT_ROTATION = ["dev3-team", "dev4-team", "dev5-team", "dev7-team", "dev8-team"]


def live_team_sessions() -> dict[str, int]:
    """현재 살아있는 봇 세션을 {team_dir_name: pid} 로 반환.

    ★ 파일럿에서 드러난 교란변수: 봇이 이미 세션을 점유 중이면 cokacdir 스케줄러가
    두 번째 세션을 띄우지 않아 스케줄이 무기한 pending 된다(= marker 영원히 없음).
    따라서 각 probe 의 t0 시점 점유 상태를 반드시 함께 기록한다.
    """
    out: dict[str, int] = {}
    for proc in Path("/proc").iterdir():
        if not proc.name.isdigit():
            continue
        try:
            cwd = os.readlink(str(proc / "cwd"))
        except Exception:
            continue
        if "/workspace/teams/dev" in cwd:
            out[Path(cwd).name] = int(proc.name)
    return out

DUMMY_DESC = """[task-2945 실측용 더미 태스크 — 코드 변경 없음]

목적: 봇 기동 지연 측정. **마커만 만들고 즉시 종료**하는 것이 전부입니다.

## 해야 할 일 (이것만)
프롬프트 최상단의 spawn 확인 마커 생성 명령을 **첫 명령으로 1회 실행**한 뒤,
아무것도 하지 말고 **즉시 세션을 종료**하세요.

## 하지 말 것 (전부 생략)
- task-timer start/end 호출 금지
- 보고서/이벤트 파일 작성 금지
- finish-task 실행 금지
- 콜백(ANU callback) 발사 금지
- 코드 읽기/수정/커밋/PR 금지
- 추가 조사·질문 금지

마커 생성 직후 "마커 생성 완료" 한 줄만 답하고 종료하면 됩니다.
"""


class MarkerWatcher(threading.Thread):
    """events 디렉토리를 고빈도 폴링해 marker 최초 목격 시각을 기록."""

    def __init__(self, events_dir: Path, interval: float = 0.2):
        super().__init__(daemon=True)
        self.events_dir = events_dir
        self.interval = interval
        self._stop = threading.Event()
        self.seen: dict[str, dict] = {}  # path -> {first_seen, task_id, payload}
        self.baseline: set[str] = set()

    def snapshot_baseline(self) -> None:
        self.baseline = {str(p) for p in self.events_dir.glob("*.spawn-confirmed-*.json")}

    def run(self) -> None:
        while not self._stop.is_set():
            now = time.time()
            try:
                for p in self.events_dir.glob("*.spawn-confirmed-*.json"):
                    sp = str(p)
                    if sp in self.baseline or sp in self.seen:
                        continue
                    task_id = p.name.split(".spawn-confirmed-")[0]
                    payload = None
                    try:
                        payload = json.loads(p.read_text(encoding="utf-8"))
                    except Exception:
                        pass
                    self.seen[sp] = {
                        "first_seen": now,
                        "task_id": task_id,
                        "path": sp,
                        "payload": payload,
                        "mtime": p.stat().st_mtime,
                    }
                    print(f"    [watch] marker 목격: {p.name} @ {datetime.fromtimestamp(now):%H:%M:%S.%f}",
                          flush=True)
            except Exception as exc:  # 관측이 실험을 죽이면 안 됨
                print(f"    [watch] 폴링 예외(무시): {exc}", flush=True)
            self._stop.wait(self.interval)

    def stop(self) -> None:
        self._stop.set()

    def find(self, task_id: str) -> dict | None:
        for rec in self.seen.values():
            if rec["task_id"] == task_id:
                return rec
        return None


def main() -> int:
    ap = argparse.ArgumentParser()
    ap.add_argument("--probes", type=int, default=8)
    ap.add_argument("--grace", type=int, default=240, help="마지막 dispatch 후 추가 관측 초")
    ap.add_argument("--start-id", type=int, default=99001)
    ap.add_argument("--out", default=str(WORKSPACE / "memory" / "reports" / "task-2945-raw.jsonl"))
    args = ap.parse_args()

    sys.path.insert(0, str(WORKSPACE))
    os.chdir(str(WORKSPACE))
    import dispatch as d  # noqa: E402

    watcher = MarkerWatcher(EVENTS_DIR)
    watcher.snapshot_baseline()
    watcher.start()
    print(f"[harness] watcher 시작 (baseline marker {len(watcher.baseline)}건)", flush=True)
    print(f"[harness] spawn_verify_enabled={d and __import__('dispatch.spawn_verification', fromlist=['x']).is_enabled()}",
          flush=True)

    records: list[dict] = []
    for i in range(args.probes):
        team = BOT_ROTATION[i % len(BOT_ROTATION)]
        task_id = f"task-{args.start_id + i}"
        occupancy = live_team_sessions()
        self_busy = team.replace("-team", "") in occupancy
        print(f"\n[probe {i+1}/{args.probes}] {task_id} → {team}  "
              f"({datetime.now():%H:%M:%S}) occupied={self_busy} live={occupancy}", flush=True)
        if self_busy:
            print("    [skip-wait] 대상 봇 점유 중 — 최대 90s 해제 대기", flush=True)
            _dl = time.time() + 90
            while time.time() < _dl and team.replace("-team", "") in live_team_sessions():
                time.sleep(2)
            occupancy = live_team_sessions()
            self_busy = team.replace("-team", "") in occupancy
        t0 = time.time()
        try:
            result = d.dispatch(
                team_id=team,
                task_desc=DUMMY_DESC,
                task_id=task_id,
                level="normal",
                task_type="coding",
                force=True,
                allow_no_scope=True,
                refresh_map=False,
            )
        except Exception as exc:
            result = {"status": "harness_exception", "message": repr(exc)}
        t1 = time.time()
        rec = {
            "probe": i + 1,
            "task_id": task_id,
            "team": team,
            "t0": t0,
            "t0_iso": datetime.fromtimestamp(t0).isoformat(),
            "target_bot_occupied_at_t0": self_busy,
            "live_sessions_at_t0": occupancy,
            "t_ret": t1,
            "dispatch_elapsed": round(t1 - t0, 3),
            "status": result.get("status") if isinstance(result, dict) else str(result),
            "spawn_verification": result.get("spawn_verification") if isinstance(result, dict) else None,
            "raw_result": {k: v for k, v in result.items() if k not in ("prompt",)}
            if isinstance(result, dict) else None,
        }
        records.append(rec)
        print(f"    dispatch 반환: status={rec['status']} verdict={rec['spawn_verification']} "
              f"elapsed={rec['dispatch_elapsed']}s", flush=True)

    print(f"\n[harness] 전 probe 발사 완료. grace {args.grace}s 추가 관측...", flush=True)
    deadline = time.time() + args.grace
    while time.time() < deadline:
        if all(watcher.find(r["task_id"]) for r in records):
            print("[harness] 전 probe marker 확보 — grace 조기 종료", flush=True)
            break
        time.sleep(1)
    watcher.stop()

    for rec in records:
        m = watcher.find(rec["task_id"])
        if m:
            rec["t_marker"] = m["first_seen"]
            rec["t_marker_iso"] = datetime.fromtimestamp(m["first_seen"]).isoformat()
            rec["marker_mtime"] = m["mtime"]
            rec["latency_from_call"] = round(m["first_seen"] - rec["t0"], 3)
            rec["latency_mtime_from_call"] = round(m["mtime"] - rec["t0"], 3)
            rec["marker_payload"] = m["payload"]
            rec["marker_path"] = m["path"]
        else:
            rec["t_marker"] = None
            rec["latency_from_call"] = None

    out = Path(args.out)
    out.parent.mkdir(parents=True, exist_ok=True)
    with out.open("w", encoding="utf-8") as f:
        for rec in records:
            f.write(json.dumps(rec, ensure_ascii=False) + "\n")
    print(f"\n[harness] raw 기록: {out}", flush=True)

    lat = [r["latency_from_call"] for r in records if r.get("latency_from_call") is not None]
    print(f"[harness] marker 확보 {len(lat)}/{len(records)}건")
    if lat:
        lat_s = sorted(lat)
        print(f"[harness] min={lat_s[0]} med={lat_s[len(lat_s)//2]} max={lat_s[-1]}")
    return 0


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