"""CSV 작업을 하나씩 처리하는 교육용 Worker. Python 3.10+, 외부 패키지 없음."""
import csv
import json
import sys
import time
from datetime import datetime, timezone
from pathlib import Path

JOBS = Path(__file__).resolve().parent / "jobs"
STATES = ("pending", "processing", "completed", "failed", "inputs", "results")


def utc_now():
    return datetime.now(timezone.utc).isoformat()


def ensure_directories():
    for name in STATES:
        (JOBS / name).mkdir(parents=True, exist_ok=True)


def write_json(path, value):
    """완성된 JSON만 보이도록 임시 파일을 쓴 뒤 교체한다."""
    temporary = path.with_suffix(".tmp")
    temporary.write_text(
        json.dumps(value, ensure_ascii=False, indent=2) + "\n", encoding="utf-8"
    )
    temporary.replace(path)


def summarize_csv(path):
    """첫 번째 비어 있지 않은 행을 헤더로 사용한다."""
    with path.open(encoding="utf-8-sig", newline="") as stream:
        reader = csv.reader(stream, strict=True)
        header = next((row for row in reader if row), None)
        if header is None:
            raise ValueError("빈 파일입니다. 헤더가 있는 CSV를 사용하세요.")
        columns = [value.strip() for value in header]
        if any(not name for name in columns):
            raise ValueError("열 이름은 비어 있을 수 없습니다.")
        if len(set(columns)) != len(columns):
            raise ValueError("열 이름은 중복될 수 없습니다.")

        missing = dict.fromkeys(columns, 0)
        count = 0
        for row in reader:
            if not row:  # 실제 빈 줄만 제외. , , 같은 데이터 행은 집계한다.
                continue
            if len(row) != len(columns):
                raise ValueError(
                    f"CSV {reader.line_num}행의 열 수가 다릅니다: "
                    f"예상 {len(columns)}, 실제 {len(row)}"
                )
            count += 1
            for column, value in zip(columns, row):
                if not value.strip():
                    missing[column] += 1
        return {"row_count": count, "columns": columns, "missing_values": missing}


def process_job(pending):
    job_id = pending.stem
    processing = JOBS / "processing" / pending.name
    pending.rename(processing)
    record = {"job_id": job_id, "status": "processing"}
    try:
        record = json.loads(processing.read_text(encoding="utf-8"))
        if not isinstance(record, dict):
            raise ValueError("작업 JSON은 객체여야 합니다.")
        record.update(job_id=job_id, status="processing", started_at=utc_now())
    except (OSError, UnicodeError, ValueError) as error:
        record = {"job_id": job_id, "status": "failed", "error": str(error)}
        write_json(processing, record)
        processing.rename(JOBS / "failed" / pending.name)
        print(f"[failed] {job_id} -> {error}", flush=True)
        return
    write_json(processing, record)
    print(f"[processing] {job_id}", flush=True)
    try:
        # 외부 경로 대신 접수 시 복사한 작업 전용 파일만 읽는다.
        summary = summarize_csv(JOBS / "inputs" / f"{job_id}.csv")
    except (OSError, UnicodeError, ValueError, csv.Error) as error:
        record.update(
            status="failed",
            failed_at=utc_now(),
            error=f"{type(error).__name__}: {error}",
        )
    else:
        result = {
            "job_id": job_id,
            "status": "completed",
            "source_file": record.get("source_file"),
            **summary,
            "completed_at": utc_now(),
        }
        write_json(JOBS / "results" / f"{job_id}.json", result)
        record.update(
            status="completed",
            completed_at=result["completed_at"],
            result=f"jobs/results/{job_id}.json",
        )
    write_json(processing, record)
    processing.rename(JOBS / record["status"] / pending.name)
    if record["status"] == "completed":
        print(f"[completed] {job_id} -> {record['result']}", flush=True)
    else:
        print(f"[failed] {job_id} -> {record['error']}", flush=True)


def run():
    ensure_directories()
    interrupted = list((JOBS / "processing").glob("*.json"))
    if interrupted:
        print(
            f"[notice] 처리 중 기록 {len(interrupted)}개가 남아 있습니다. "
            "자동 재시도하지 않습니다. 결과를 확인한 뒤 필요하면 다시 접수하세요.",
            flush=True,
        )
    print("[worker] 대기 중입니다. 다른 터미널에서 접수하세요. 종료: Ctrl+C", flush=True)
    try:
        while True:
            pending = sorted(
                (JOBS / "pending").glob("*.json"),
                key=lambda item: (item.stat().st_mtime_ns, item.name),
            )
            for job in pending:
                process_job(job)
            time.sleep(1)
    except KeyboardInterrupt:
        print("\n[worker] 종료했습니다. 작업과 결과 파일은 그대로 남습니다.")
        return 0
    except OSError as error:
        # 기록을 저장할 수 없는 디스크 오류는 조용히 삼키지 않는다.
        print(f"[worker error] 파일 저장 상태를 확인하세요: {error}", file=sys.stderr)
        return 1


if __name__ == "__main__":
    sys.exit(run())
