# -*- coding: utf-8 -*-
"""dispatch.anu_collector_result — task-2730 collector_result schema + atomic writer.

OS-level pickup deterministic closeout 의 산출물(``<task_id>.collector_result.json``)
스키마와 **원자적 writer**. ``pickup_once``(runner) 가 durability ordering
(**ledger → collector_result → done marker**) 안에서 이 helper 를 호출한다.

설계 anchor (수정 금지):
  - ANCHOR-A: closeout write OWNER = ``pickup_once`` (+ 이 helper). ``process_one``
    은 closeout 을 **소유하지 않는다** (verdict 분기만).
  - ANCHOR-B: durability order = ledger(dedupe·fsync) → collector_result(atomic
    os.replace) → done marker(terminal sentinel). 이 모듈의 writer 는 atomic
    write 만 책임지고, **순서 보장은 호출자(pickup_once)** 가 한다 (side-effect 최소).

★ raw key 0: ``owner_key`` / ANU key / argv literal 을 **절대 기록하지 않는다**.
  ``owner_proof`` 는 분류 라벨(outcome / verdict / schedule_id / query_ok)만 담는다.

Layer A / NO-CRON: cron register/remove 0, subprocess 0, merge 0, git/gh import 0.
"""
from __future__ import annotations

import json
import os
from dataclasses import dataclass, field
from typing import Optional

# ── 상수 ─────────────────────────────────────────────────────────────────────
COLLECTOR_RESULT_SCHEMA = "dispatch.anu_collector_result.v1"
RUNNER_VERSION = "task-2730.os-pickup-closeout.v1"

# closeout_action 값 (deterministic closeout 결과)
CLOSEOUT_DONE_ACKED = "done_acked"          # green-path: 결정론 closeout 완료(wake 0)
CLOSEOUT_RELAY_PENDING = "relay_pending"    # relay-path: 2-tier terminal relay 대기
CLOSEOUT_QUARANTINED = "quarantined"        # owner-proof L1 실패 → 격리(작업물 보존)

# collector_role 라벨 (raw key 아님 — 분류 라벨만)
COLLECTOR_ROLE_ANU = "ANU"

# task-2753: adopted_via 라벨 — collector_result 가 어떤 경로로 채택되었는지.
#   "envelope": 기존 경로(envelope 있는 self-authored/verify 경로).
#   "provenance": 신규 Option B provenance-path(envelope 부재 → provenance hard-fail 통과).
ADOPTED_VIA_ENVELOPE = "envelope"
ADOPTED_VIA_PROVENANCE = "provenance"

# dedupe / scope 기본 라벨
DEFAULT_SCOPE = "p0b_pickup_closeout"

# agent_relay producer-contract 필드 (result.json 의 relay_hints).
# runner 는 **이 4개 필드만** 읽어 판정한다 — 임의 추론 0.
RELAY_HINT_FIELDS = (
    "gemini_finding",
    "merge_ready_ambiguous",
    "critical7",
    "consolidated_report",
)


# ── owner_proof 분류 라벨 (raw key 0) ─────────────────────────────────────────
@dataclass
class OwnerProof:
    """owner-proof 분류 라벨. resolve_authoritative_owner / verify_collector_authoritative
    결과에서 파생한 라벨만 담는다 (raw key/argv 0)."""

    l1_outcome: str = ""     # OWNER_ANU | NOT_ANU_OWNED_OR_ACCESS_DENIED | QUERY_FAILED | PENDING | ""
    l2_verdict: str = ""     # AUTHORITATIVE | QUARANTINED | NON_AUTHORITATIVE | REJECTED | ""
    schedule_id: str = ""
    query_ok: bool = False


# ── agent_relay 판정 (producer contract) ──────────────────────────────────────
@dataclass
class AgentRelay:
    """agent_relay 판정 결과. required + reason(relay_hints 필드명 또는 quarantine_reason)."""

    required: bool = False
    reason: str = ""


# ── CollectorResult schema ────────────────────────────────────────────────────
@dataclass
class CollectorResult:
    """collector_result.json 스키마 (pure dataclass). raw key 미포함."""

    task_id: str
    source_result_sha256: str
    verdict: str
    closeout_action: str
    owner_proof: OwnerProof = field(default_factory=OwnerProof)
    agent_relay: AgentRelay = field(default_factory=AgentRelay)
    dedupe: bool = False
    scope: str = DEFAULT_SCOPE
    contract_violation: Optional[str] = None
    ts_kst: str = ""
    runner_version: str = RUNNER_VERSION
    schema: str = COLLECTOR_RESULT_SCHEMA
    collector_role: str = COLLECTOR_ROLE_ANU
    # task-2753: 채택 경로 라벨. 기본 "envelope" 로 하위호환(기존 pickup_once 가
    #   set 안 해도 envelope 로 동작). provenance-path 는 "provenance" 로 set.
    adopted_via: str = ADOPTED_VIA_ENVELOPE

    def to_json(self) -> dict:
        return {
            "schema": self.schema,
            "task_id": self.task_id,
            "source_result_sha256": self.source_result_sha256,
            "verdict": self.verdict,
            "owner_proof": {
                "l1_outcome": self.owner_proof.l1_outcome,
                "l2_verdict": self.owner_proof.l2_verdict,
                "schedule_id": self.owner_proof.schedule_id,
                "query_ok": bool(self.owner_proof.query_ok),
            },
            "dedupe": bool(self.dedupe),
            "scope": self.scope,
            "closeout_action": self.closeout_action,
            "agent_relay": {
                "required": bool(self.agent_relay.required),
                "reason": self.agent_relay.reason,
            },
            "contract_violation": self.contract_violation,
            "collector_role": self.collector_role,
            "adopted_via": self.adopted_via,
            "ts_kst": self.ts_kst,
            "runner_version": self.runner_version,
        }


# ── agent_relay 판정 (relay_hints producer contract, 추론 0) ──────────────────
def determine_agent_relay(result: object) -> AgentRelay:
    """result.json 의 ``relay_hints`` producer-contract 필드만 읽어 agent_relay 판정.

    - ``relay_hints`` 가 dict 가 아니거나 부재 → required=false(green).
    - ``RELAY_HINT_FIELDS`` 4종 중 하나라도 truthy → required=true,
      reason=처음 truthy 인 필드명.
    - **임의 추론 0**: 정의된 4 필드만 본다. 다른 키는 무시한다.
    """
    hints = result.get("relay_hints") if isinstance(result, dict) else None
    if not isinstance(hints, dict):
        return AgentRelay(required=False, reason="")
    for name in RELAY_HINT_FIELDS:
        if bool(hints.get(name)):
            return AgentRelay(required=True, reason=name)
    return AgentRelay(required=False, reason="")


# ── collector_result 경로 ─────────────────────────────────────────────────────
def collector_result_path(result_json_path: str, task_id: str) -> str:
    """collector_result.json 경로: result.json 과 같은 디렉토리,
    ``<task_id>.collector_result.json``."""
    result_dir = os.path.dirname(result_json_path)
    return os.path.join(result_dir, f"{task_id}.collector_result.json")


# ── parent-dir fsync (rename crash-durability) ───────────────────────────────
def _fsync_dir(dir_path: str) -> None:
    """디렉토리 엔트리(rename) 를 crash-durable 화. fsync 미지원 fs/플랫폼은 무시(fail-safe)."""
    try:
        fd = os.open(dir_path, os.O_RDONLY)
    except OSError:
        return
    try:
        os.fsync(fd)
    except OSError:
        pass  # 일부 fs 는 디렉토리 fsync 미지원 — best-effort.
    finally:
        os.close(fd)


# ── atomic writer ─────────────────────────────────────────────────────────────
def write_collector_result(result: CollectorResult, path: str) -> str:
    """``CollectorResult`` 를 **atomic(os.replace)** 으로 ``path`` 에 write. 경로 반환.

    side-effect 최소:
      - 부모 디렉토리 생성(makedirs exist_ok).
      - tmp 파일 write + flush + fsync.
      - ``os.replace(tmp, path)`` (같은 fs atomic rename).
    durability 순서(ledger→collector_result→marker)는 호출자(pickup_once)가 보장한다.
    실패 시 tmp 정리 후 OSError 를 **그대로 raise** — 호출자가 marker 미작성으로
    fail-closed 하여 다음 cycle 재시도하게 한다(collector_result before marker 불변 보존).

    raw key 미기록: ``result.to_json()`` 에 owner_key/ANU key/argv literal 0.
    """
    payload = json.dumps(result.to_json(), ensure_ascii=False, indent=2)
    parent = os.path.dirname(path) or "."
    os.makedirs(parent, exist_ok=True)
    tmp_path = f"{path}.tmp-{os.getpid()}"
    try:
        with open(tmp_path, "w", encoding="utf-8") as fh:
            fh.write(payload)
            fh.flush()
            os.fsync(fh.fileno())
        os.replace(tmp_path, path)
        _fsync_dir(parent)  # rename(dir entry) crash-durable화 — collector_result 유실 방지.
    except OSError:
        try:
            os.unlink(tmp_path)
        except OSError:
            pass
        raise
    return path


__all__ = [
    "COLLECTOR_RESULT_SCHEMA",
    "RUNNER_VERSION",
    "CLOSEOUT_DONE_ACKED",
    "CLOSEOUT_RELAY_PENDING",
    "CLOSEOUT_QUARANTINED",
    "COLLECTOR_ROLE_ANU",
    "ADOPTED_VIA_ENVELOPE",
    "ADOPTED_VIA_PROVENANCE",
    "DEFAULT_SCOPE",
    "RELAY_HINT_FIELDS",
    "OwnerProof",
    "AgentRelay",
    "CollectorResult",
    "determine_agent_relay",
    "collector_result_path",
    "write_collector_result",
]
