"""GitHub 저장소 하나로 만드는 원격 작업 대기열 (작업을 받는 PC 쪽 실행기).

폴더 구조 (저장소 안):
    queue/<기기이름>/inbox/    할 일 파일(.cmd / .ps1 / .sh)을 여기에 넣는다
    queue/<기기이름>/done/     끝난 작업 파일과 실행 기록(.log)이 옮겨진다
    queue/<기기이름>/results/  작업이 만든 결과 파일
    queue/<기기이름>/heartbeat.txt  마지막으로 살아 있던 시각

사용법 (PC에서 5분마다 실행):
    python remote_queue.py --repo C:\\work\\myrepo --device pc

작업 파일 앞부분에 "rem timeout=600" 또는 "# timeout=600" 을 쓰면 제한 시간(초)이 된다.

주의: 이 저장소는 대기열 전용으로 쓴다. 실행기는 자기 대기열 폴더 밖의
로컬 변경과 git이 모르는 새 파일을 매 회차 지운다.
"""
from __future__ import annotations

import argparse
import datetime as dt
import os
import re
import shutil
import subprocess
from pathlib import Path

DEFAULT_TIMEOUT = 1800
# Windows에서 5분마다 git·cmd 창이 번쩍이지 않게 한다
NO_WINDOW = {"creationflags": 0x08000000} if os.name == "nt" else {}
# git 메시지를 영어로 고정하고(한글 메시지 해석 오류 방지), 로그인 창을 띄우지 말고 실패하게 한다
GIT_ENV = {**os.environ, "LC_ALL": "C", "LANG": "C", "GIT_TERMINAL_PROMPT": "0",
           "GCM_INTERACTIVE": "never"}


def git(repo: Path, *args: str) -> subprocess.CompletedProcess:
    return subprocess.run(["git", *args], cwd=repo, capture_output=True, text=True,
                          encoding="utf-8", errors="replace", env=GIT_ENV, **NO_WINDOW)


def lock(repo: Path):
    """같은 실행기가 겹쳐 돌지 않게 잠근다. 프로그램이 끝나거나 죽으면 잠금은 저절로 풀린다."""
    f = open(repo / ".git" / "remote_queue.lock", "a+")
    try:
        if os.name == "nt":
            import msvcrt
            msvcrt.locking(f.fileno(), msvcrt.LK_NBLCK, 1)
        else:
            import fcntl
            fcntl.flock(f, fcntl.LOCK_EX | fcntl.LOCK_NB)
    except OSError:
        return None
    return f


def timeout_of(job: Path) -> int:
    head = job.read_text(encoding="utf-8", errors="replace")[:500]
    m = re.search(r"(?:rem|#)\s*timeout\s*=\s*(\d+)", head, re.I)
    return int(m.group(1)) if m else DEFAULT_TIMEOUT


def command_for(job: Path) -> list[str]:
    if job.suffix == ".cmd":
        return ["cmd", "/c", str(job)]
    if job.suffix == ".ps1":
        return ["powershell", "-NoProfile", "-ExecutionPolicy", "Bypass", "-File", str(job)]
    return ["bash", str(job)]


def kill_tree(p: subprocess.Popen) -> None:
    if os.name == "nt":
        subprocess.run(["taskkill", "/PID", str(p.pid), "/T", "/F"], capture_output=True, **NO_WINDOW)
    else:
        import signal
        try:
            os.killpg(p.pid, signal.SIGKILL)
        except ProcessLookupError:
            pass


def publish(repo: Path, paths: list[str], msg: str) -> None:
    """자기 대기열 폴더만 커밋해서 올린다. 지난번에 못 올린 커밋이 있으면 그것도 올린다."""
    git(repo, "add", "-A", "--", *paths)
    if git(repo, "diff", "--cached", "--quiet", "--", *paths).returncode != 0:
        if git(repo, "commit", "-q", "-m", msg, "--", *paths).returncode != 0:
            print("commit 실패")
            return
    ahead = git(repo, "rev-list", "--count", "@{u}..HEAD").stdout.strip()
    if ahead in ("", "0"):
        return
    if git(repo, "push", "-q").returncode != 0:
        # 다른 기기가 먼저 올렸으면 합치고 다시. 합치다 막히면 멈춘 상태로 두지 않고 취소한다(다음 회차 시작 때 복구).
        if git(repo, "pull", "--rebase", "--autostash", "-q").returncode != 0:
            git(repo, "rebase", "--abort")
            print("합치기 실패 — 다음 회차에 복구")
            return
        print("push", "ok" if git(repo, "push", "-q").returncode == 0 else "실패 — 다음 회차에 다시")


def recover(repo: Path, base: Path) -> bool:
    """받기가 막혔을 때: 원격 상태로 맞춘 뒤, 이 기기가 만든 결과(done·results·heartbeat)만 되살린다."""
    keep = repo / ".git" / "rq_backup"
    shutil.rmtree(keep, ignore_errors=True)
    for name in ("done", "results", "heartbeat.txt"):
        src = base / name
        if src.is_dir():
            shutil.copytree(src, keep / name)
        elif src.exists():
            keep.mkdir(parents=True, exist_ok=True)
            shutil.copy2(src, keep / name)
    if git(repo, "fetch", "-q").returncode != 0:
        return False   # 네트워크 문제면 다음 회차에
    git(repo, "reset", "-q", "--hard", "@{u}")
    for item in keep.iterdir() if keep.exists() else []:
        if item.is_dir():
            shutil.copytree(item, base / item.name, dirs_exist_ok=True)
        else:
            shutil.copy2(item, base / item.name)
    shutil.rmtree(keep, ignore_errors=True)
    print("받기 막힘 — 원격 상태로 맞추고 이 기기의 결과를 되살림")
    return True


def run_job(repo: Path, job: Path, log: Path) -> int:
    start = dt.datetime.now()
    limit = timeout_of(job)
    out_file = repo / ".git" / f"rq_{job.stem}.out"   # 작업마다 따로: 앞 작업이 남긴 프로그램의 출력이 섞이지 않게
    extra = NO_WINDOW if os.name == "nt" else {"start_new_session": True}
    env = {**os.environ, "PYTHONUTF8": "1"}
    # 출력은 파일로 받는다: 작업이 띄운 프로그램이 출력 통로를 붙잡고 있어도 실행기가 멈추지 않는다
    with open(out_file, "wb") as out:
        p = subprocess.Popen(command_for(job), cwd=repo, stdin=subprocess.DEVNULL, stdout=out,
                             stderr=subprocess.STDOUT, env=env, **extra)
        try:
            code = p.wait(timeout=limit)
            note = ""
        except subprocess.TimeoutExpired:
            kill_tree(p)   # cmd만 끄면 cmd가 띄운 프로그램이 남아 계속 돈다
            try:
                p.wait(timeout=30)
            except subprocess.TimeoutExpired:
                pass
            code, note = -1, f"\nTIMEOUT after {limit}s (하위 프로세스까지 종료)\n"
    text = out_file.read_bytes().decode("utf-8", errors="replace").replace("\r\n", "\n") + note
    try:
        out_file.unlink()
    except OSError:
        pass   # 남은 프로그램이 아직 붙잡고 있으면(Windows) 그대로 둔다
    end = dt.datetime.now()
    log.write_text(f"job: {job.name}\nstart: {start:%Y-%m-%d %H:%M:%S}\nend: {end:%Y-%m-%d %H:%M:%S}\n"
                   f"exit: {code}\n{'=' * 40}\n{text}", encoding="utf-8")
    return code


def main() -> None:
    ap = argparse.ArgumentParser()
    ap.add_argument("--repo", required=True)
    ap.add_argument("--device", required=True)
    a = ap.parse_args()
    repo = Path(a.repo).resolve()
    held = lock(repo)
    if held is None:
        print("다른 실행이 아직 도는 중 — 이번 회차는 건너뜀")
        return
    base = repo / "queue" / a.device
    inbox, done, results = base / "inbox", base / "done", base / "results"
    for d in (inbox, done, results):
        d.mkdir(parents=True, exist_ok=True)
        (d / ".keep").touch()   # 빈 폴더는 git에 올라가지 않으므로 표시 파일을 둔다

    # 0) 이 기기는 자기 대기열 폴더(queue/<기기이름>) 밖은 고치지 않는다.
    #    그 밖에 생긴 로컬 변경·새 파일(줄바꿈 변환, 실수로 고친 파일, 지난번 충돌 흔적)은 버려서 받기가 막히지 않게 한다.
    own = base.relative_to(repo).as_posix()
    git(repo, "rebase", "--abort")
    git(repo, "reset", "-q")                                     # 스테이징 비우기
    git(repo, "checkout", "-q", "--", ".", f":(exclude){own}")    # 자기 폴더 밖 변경 되돌리기
    git(repo, "clean", "-fdq", "--", ".", f":(exclude){own}")     # 자기 폴더 밖 새 파일 지우기

    # 1) 새 작업 받기. 실패하면 이번 회차는 건너뛰고 다음에 다시 시도한다.
    if git(repo, "pull", "--rebase", "--autostash", "-q").returncode != 0:
        git(repo, "rebase", "--abort")
        if not recover(repo, base):
            print("pull 실패 — 다음 회차에 다시 시도")
            return
    for d in (inbox, done, results):
        d.mkdir(parents=True, exist_ok=True)
        (d / ".keep").touch()

    paths = [own]   # 올리는 것은 자기 대기열 폴더뿐이다 (다른 폴더의 옛 사본이 섞여 올라가지 않게)
    # 2) 이름 순서대로 실행 (010_, 020_ 처럼 번호를 붙이면 순서가 정해진다)
    jobs = sorted(p for p in inbox.iterdir() if p.suffix in (".cmd", ".ps1", ".sh"))
    for job in jobs:
        if not job.exists():   # 실행 도중 받은 변경으로 보내는 쪽이 지운 작업
            continue
        if (done / job.name).exists():   # 이미 시작한 작업(시작 기록을 못 올린 채 복구된 경우)은 다시 돌리지 않는다
            job.unlink()
            continue
        # 실행 전에 먼저 done으로 옮기고 올려 둔다. 작업이 PC를 끄거나 실행기가 죽어도 같은 작업이 다시 돌지 않는다.
        moved = done / job.name
        shutil.move(str(job), str(moved))
        log = done / (job.stem + ".log")
        log.write_text(f"job: {job.name}\nstart: {dt.datetime.now():%Y-%m-%d %H:%M:%S}\nexit: running\n",
                       encoding="utf-8")
        publish(repo, paths, f"{a.device}: start {job.name}")
        code = run_job(repo, moved, log)
        print(f"{job.name}: exit={code}")

    # 살아 있다는 표시: 작업을 했거나, 마지막 표시가 1시간 넘었을 때만 (빈 커밋이 5분마다 쌓이지 않게)
    hb = base / "heartbeat.txt"
    stale = not hb.exists() or (dt.datetime.now().timestamp() - hb.stat().st_mtime) > 3600
    if jobs or stale:
        hb.write_text(f"{dt.datetime.now():%Y-%m-%d %H:%M:%S}\n", encoding="utf-8")

    # 3) 결과를 올린다.
    publish(repo, paths, f"{a.device}: {len(jobs)} jobs")


if __name__ == "__main__":
    main()
