# -*- coding: utf-8 -*-
"""task-2730 — OS-level pickup deterministic closeout regression (통합).

corrected_spec_v2 §3·§4 / design-lock candidate §3 10항목 + 회장 regression 8 검증.
네트워크 0, 전부 mock/fixture, tmpdir(isolated) 사용. canonical events 0 touch.
ANU key literal 절대 노출 금지 — 모듈 상수 / sealed_key_loader fake 주입으로만.

커버리지:
  ① owner-proof L1/L2 (OWNER_ANU→CLOSEOUT/relay / NOT_ANU·QUERY_FAILED→fail-closed
     / PENDING→retry / self-key→refuse)
  ② agent_relay relay_hints 4종 분기 (각 true→WAKE_BUILT / 전부 false·부재→CLOSEOUT_DONE)
  ③ durability order ledger→collector_result→marker + crash 3케이스 단일 closeout
  ④ idempotency: dedupe ledger + done/acked marker + 직렬 2회 → wake 1회
  ⑤ terminal_relay static (shell=False·git/gh/dispatch/merge import 0·argv 화이트리스트)
  ⑥ green-path driver launcher_fn 호출 0 (wake 0) + VERDICT_CLOSEOUT_DONE + move_processed
  ⑦ collector_result raw key 0
  ⑧ green CLOSEOUT_DONE 도 driver_enabled flag OFF 시 미실행(dry-run isolated, canonical 0 touch)
"""
from __future__ import annotations

import ast
import importlib.util
import inspect
import json
import os
import shutil
import sys
import tempfile
import unittest
from datetime import datetime, timedelta, timezone
from pathlib import Path

_ROOT = Path(__file__).resolve().parent.parent
if str(_ROOT) not in sys.path:
    sys.path.insert(0, str(_ROOT))


def _load(modname: str, relpath: str):
    if modname in sys.modules:
        return sys.modules[modname]
    spec = importlib.util.spec_from_file_location(modname, _ROOT / relpath)
    assert spec is not None and spec.loader is not None
    mod = importlib.util.module_from_spec(spec)
    sys.modules[modname] = mod
    spec.loader.exec_module(mod)
    return mod


# 의존 모듈 실 로드 (2720 패턴 동일)
_load("dispatch.callback_owner_enforcer", "dispatch/callback_owner_enforcer.py")
_load("dispatch.normal_fallback_callback_helper",
      "dispatch/normal_fallback_callback_helper.py")
M_enf = _load("dispatch.anu_owned_callback_enforcement",
              "dispatch/anu_owned_callback_enforcement.py")
CR = _load("dispatch.anu_collector_result", "dispatch/anu_collector_result.py")
TR = _load("dispatch.anu_terminal_relay", "dispatch/anu_terminal_relay.py")
M = _load("dispatch.anu_result_pickup_runner", "dispatch/anu_result_pickup_runner.py")
_load("dispatch.anu_pickup_wake_launcher", "dispatch/anu_pickup_wake_launcher.py")
DRV = _load("dispatch.anu_pickup_driver", "dispatch/anu_pickup_driver.py")

_ANU_KEY = M_enf.ANU_KEY
_DEV_KEY = "7943afbe12c12f7d"
_NOW = datetime(2026, 6, 9, 3, 0, 0, tzinfo=timezone.utc)


def _clock():
    return _NOW


def _fresh_ts():
    return (_NOW - timedelta(minutes=5)).strftime("%Y-%m-%dT%H:%M:%SZ")


def _loader():
    return _ANU_KEY


def _make_probe(owned):
    owned = set(owned)

    def _probe(sid: str) -> dict:
        if sid in owned:
            return {"status": "ok", "id": sid, "count": 1,
                    "history": [{"status": "ok", "ts": "2026-06-09 03:00:00"}]}
        return {"status": "error", "message": f"not found or access denied: {sid}"}

    return _probe


def _write_result(rdir: str, task_id: str, *, relay_hints=None,
                  envelope=None) -> str:
    payload = {"task_id": task_id, "summary": "done", "completion_signal": "RESULT_JSON_WRITTEN"}
    if relay_hints is not None:
        payload["relay_hints"] = relay_hints
    if envelope is not None:
        payload["collector_envelope"] = envelope
    path = os.path.join(rdir, f"{task_id}.result.json")
    with open(path, "w", encoding="utf-8") as fh:
        json.dump(payload, fh, ensure_ascii=False)
    return path


class TestAgentRelayBranch(unittest.TestCase):
    """② agent_relay relay_hints 4종 분기."""

    def setUp(self):
        self._tmp = tempfile.mkdtemp(prefix="t2730-relay-")
        self.addCleanup(shutil.rmtree, self._tmp, ignore_errors=True)

    def _run(self, task_id, relay_hints):
        rdir = tempfile.mkdtemp(dir=self._tmp)
        ledger = os.path.join(rdir, "l.jsonl")
        path = _write_result(rdir, task_id, relay_hints=relay_hints)
        return M.pickup_once(path, executor_key=_DEV_KEY, clock=_clock,
                             sealed_key_loader=_loader, ledger_path=ledger)

    def test_each_hint_true_triggers_relay_wake_built(self):
        for field in CR.RELAY_HINT_FIELDS:
            res = self._run(f"task-relay-{field}", {field: True})
            self.assertEqual(res.verdict, M.PICKUP_WAKE_BUILT,
                             f"relay_hint {field}=true 는 WAKE_BUILT 여야 함")
            self.assertTrue(res.wake_built)
            self.assertTrue(res.argv)
            self.assertEqual(res.closeout_action, CR.CLOSEOUT_RELAY_PENDING)
            self.assertEqual(res.agent_relay_reason, field)

    def test_all_false_is_green_closeout_wake_zero(self):
        res = self._run("task-green-allfalse",
                        {k: False for k in CR.RELAY_HINT_FIELDS})
        self.assertEqual(res.verdict, M.PICKUP_CLOSEOUT_DONE)
        self.assertFalse(res.wake_built)
        self.assertIsNone(res.argv)
        self.assertEqual(res.closeout_action, CR.CLOSEOUT_DONE_ACKED)

    def test_absent_hints_is_green_closeout(self):
        res = self._run("task-green-absent", None)
        self.assertEqual(res.verdict, M.PICKUP_CLOSEOUT_DONE)
        self.assertFalse(res.wake_built)
        self.assertFalse(res.agent_relay_required)

    def test_unknown_hint_field_ignored_no_inference(self):
        # 정의되지 않은 필드는 무시(추론 0) → green.
        res = self._run("task-green-unknown", {"something_else": True})
        self.assertEqual(res.verdict, M.PICKUP_CLOSEOUT_DONE)


class TestOwnerProofL1L2(unittest.TestCase):
    """① owner-proof L1/L2 (pickup_once gh_probe 경로)."""

    def setUp(self):
        self._tmp = tempfile.mkdtemp(prefix="t2730-owner-")
        self.addCleanup(shutil.rmtree, self._tmp, ignore_errors=True)

    def _env(self, sid):
        return {"task_id": self._tid, "schedule_id": sid, "recorded_at": _fresh_ts()}

    def test_owner_anu_green_closeout(self):
        self._tid = "task-owner-anu"
        rdir = tempfile.mkdtemp(dir=self._tmp)
        ledger = os.path.join(rdir, "l.jsonl")
        path = _write_result(rdir, self._tid, envelope=self._env("ANUSID"))
        res = M.pickup_once(path, gh_probe=_make_probe({"ANUSID"}),
                            executor_key=_DEV_KEY, clock=_clock,
                            sealed_key_loader=_loader, ledger_path=ledger)
        self.assertEqual(res.verdict, M.PICKUP_CLOSEOUT_DONE)
        # owner_proof 라벨이 collector_result 에 기록됨(AUTHORITATIVE).
        cr = json.load(open(res.closeout_path, encoding="utf-8"))
        self.assertEqual(cr["owner_proof"]["l2_verdict"], M_enf.VERDICT_AUTHORITATIVE)
        self.assertTrue(cr["owner_proof"]["query_ok"])

    def test_not_anu_quarantine(self):
        self._tid = "task-owner-notanu"
        rdir = tempfile.mkdtemp(dir=self._tmp)
        ledger = os.path.join(rdir, "l.jsonl")
        path = _write_result(rdir, self._tid, envelope=self._env("SELFSID"))
        res = M.pickup_once(path, gh_probe=_make_probe({"OTHER"}),
                            executor_key=_DEV_KEY, clock=_clock,
                            sealed_key_loader=_loader, ledger_path=ledger)
        self.assertEqual(res.verdict, M.PICKUP_QUARANTINE)
        self.assertFalse(res.wake_built)
        # quarantine 은 ledger/marker/collector_result 미작성 (ANCHOR-B: 종결 sentinel 0).
        self.assertFalse(os.path.exists(os.path.join(rdir, f"{self._tid}.pickup.done")))
        self.assertFalse(os.path.exists(
            os.path.join(rdir, f"{self._tid}.collector_result.json")))

    def test_query_failed_fail_closed(self):
        self._tid = "task-owner-qf"
        rdir = tempfile.mkdtemp(dir=self._tmp)
        ledger = os.path.join(rdir, "l.jsonl")
        path = _write_result(rdir, self._tid, envelope=self._env("X"))

        def qf_probe(_sid):
            return {"status": "error", "message": "boom unexpected"}

        res = M.pickup_once(path, gh_probe=qf_probe, executor_key=_DEV_KEY,
                            clock=_clock, sealed_key_loader=_loader, ledger_path=ledger)
        self.assertEqual(res.verdict, M.PICKUP_FAIL)
        self.assertFalse(res.wake_built)

    def test_pending_retryable(self):
        self._tid = "task-owner-pending"
        rdir = tempfile.mkdtemp(dir=self._tmp)
        ledger = os.path.join(rdir, "l.jsonl")
        path = _write_result(rdir, self._tid, envelope=self._env("X"))

        def weird_probe(_sid):
            return {"status": "weird"}

        res = M.pickup_once(path, gh_probe=weird_probe, executor_key=_DEV_KEY,
                            clock=_clock, sealed_key_loader=_loader, ledger_path=ledger)
        self.assertEqual(res.verdict, M.PICKUP_PENDING)

    def test_self_key_refused_on_relay_path(self):
        # relay-path 에서 sealed key == executor self key → anu_runner refuse → FAIL.
        self._tid = "task-self-key"
        rdir = tempfile.mkdtemp(dir=self._tmp)
        ledger = os.path.join(rdir, "l.jsonl")
        path = _write_result(rdir, self._tid, relay_hints={"critical7": True})
        res = M.pickup_once(path, executor_key=_DEV_KEY, clock=_clock,
                            sealed_key_loader=lambda: _DEV_KEY,  # self-key as ANU key
                            ledger_path=ledger)
        self.assertEqual(res.verdict, M.PICKUP_FAIL)
        self.assertFalse(res.wake_built)


class TestDurabilityOrdering(unittest.TestCase):
    """③ durability order ledger→collector_result→marker + crash 3케이스."""

    def setUp(self):
        self._tmp = tempfile.mkdtemp(prefix="t2730-dur-")
        self.addCleanup(shutil.rmtree, self._tmp, ignore_errors=True)

    def test_order_ledger_then_collector_then_marker(self):
        rdir = tempfile.mkdtemp(dir=self._tmp)
        ledger = os.path.join(rdir, "l.jsonl")
        tid = "task-order"
        path = _write_result(rdir, tid)  # green
        done = os.path.join(rdir, f"{tid}.pickup.done")
        cr_path = os.path.join(rdir, f"{tid}.collector_result.json")

        order = []
        real_write = M.write_collector_result

        def spy_write(collector, p):
            # collector_result write 시점: ledger 이미 존재, marker 아직 없음.
            self.assertTrue(os.path.isfile(ledger), "ledger 가 collector_result 보다 먼저")
            self.assertFalse(os.path.exists(done), "marker 는 collector_result 보다 나중")
            order.append("collector")
            return real_write(collector, p)

        M.write_collector_result = spy_write
        try:
            res = M.pickup_once(path, executor_key=_DEV_KEY, clock=_clock,
                                sealed_key_loader=_loader, ledger_path=ledger)
        finally:
            M.write_collector_result = real_write
        self.assertEqual(res.verdict, M.PICKUP_CLOSEOUT_DONE)
        self.assertEqual(order, ["collector"])
        # 종료 후 3개 모두 존재.
        self.assertTrue(os.path.isfile(ledger))
        self.assertTrue(os.path.isfile(cr_path))
        self.assertTrue(os.path.isfile(done))

    def test_crash_after_ledger_recovers_single_closeout(self):
        # collector_result write 실패(crash 모사) → COLLECTOR_WRITE_FAILED, marker 0.
        rdir = tempfile.mkdtemp(dir=self._tmp)
        ledger = os.path.join(rdir, "l.jsonl")
        tid = "task-crash-ledger"
        path = _write_result(rdir, tid)

        def boom(collector, p):
            raise OSError("simulated crash after ledger")

        real = M.write_collector_result
        M.write_collector_result = boom
        try:
            res1 = M.pickup_once(path, executor_key=_DEV_KEY, clock=_clock,
                                 sealed_key_loader=_loader, ledger_path=ledger)
        finally:
            M.write_collector_result = real
        self.assertEqual(res1.verdict, M.PICKUP_COLLECTOR_WRITE_FAILED)
        self.assertIsNone(res1.marker_path)
        # 재처리: recovery → SKIP_DEDUPE, collector_result/marker 보강, ledger 1줄.
        res2 = M.pickup_once(path, executor_key=_DEV_KEY, clock=_clock,
                             sealed_key_loader=_loader, ledger_path=ledger)
        self.assertEqual(res2.verdict, M.PICKUP_SKIP_DEDUPE)
        self.assertFalse(res2.wake_built)
        self.assertTrue(os.path.isfile(os.path.join(rdir, f"{tid}.collector_result.json")))
        self.assertTrue(os.path.isfile(os.path.join(rdir, f"{tid}.pickup.done")))
        with open(ledger, encoding="utf-8") as fh:
            self.assertEqual(sum(1 for ln in fh if ln.strip()), 1)

    def test_crash_collector_done_marker_missing_recovers(self):
        # collector_result 성공 + marker(os.replace .pickup.done)만 실패 → CLOSEOUT_DONE
        # (marker 비치명), 재처리 시 recovery 로 marker 보강.
        rdir = tempfile.mkdtemp(dir=self._tmp)
        ledger = os.path.join(rdir, "l.jsonl")
        tid = "task-crash-marker"
        path = _write_result(rdir, tid)
        done = os.path.join(rdir, f"{tid}.pickup.done")
        real_replace = M.os.replace

        def selective(src, dst):
            if str(dst).endswith(".pickup.done"):
                raise OSError("marker-only failure")
            return real_replace(src, dst)

        M.os.replace = selective
        try:
            res1 = M.pickup_once(path, executor_key=_DEV_KEY, clock=_clock,
                                 sealed_key_loader=_loader, ledger_path=ledger)
        finally:
            M.os.replace = real_replace
        self.assertEqual(res1.verdict, M.PICKUP_CLOSEOUT_DONE)
        self.assertIsNone(res1.marker_path)  # marker 미작성
        self.assertTrue(os.path.isfile(os.path.join(rdir, f"{tid}.collector_result.json")))
        self.assertFalse(os.path.exists(done))
        # 재처리 → recovery → marker 보강 + 신규 wake 0.
        res2 = M.pickup_once(path, executor_key=_DEV_KEY, clock=_clock,
                             sealed_key_loader=_loader, ledger_path=ledger)
        self.assertEqual(res2.verdict, M.PICKUP_SKIP_DEDUPE)
        self.assertTrue(os.path.isfile(done))

    def test_recovery_completes_without_sealed_key(self):
        # crash-after-ledger recovery 는 sealed-key 부재여도 완결되어야 한다(stranding 방지).
        import hashlib
        rdir = tempfile.mkdtemp(dir=self._tmp)
        ledger = os.path.join(rdir, "l.jsonl")
        tid = "task-rec-nokey"
        path = _write_result(rdir, tid)
        sha = hashlib.sha256(open(path, "rb").read()).hexdigest()
        with open(ledger, "w", encoding="utf-8") as fh:
            fh.write(json.dumps({"event": "PICKUP_CLOSEOUT_DONE",
                                 "task_id": tid, "sha256": sha}) + "\n")
        res = M.pickup_once(path, executor_key=_DEV_KEY, clock=_clock,
                            sealed_key_loader=lambda: None,  # key 부재
                            ledger_path=ledger)
        self.assertEqual(res.verdict, M.PICKUP_SKIP_DEDUPE)
        self.assertTrue(os.path.isfile(os.path.join(rdir, f"{tid}.collector_result.json")))
        self.assertTrue(os.path.isfile(os.path.join(rdir, f"{tid}.pickup.done")))

    def test_fresh_still_requires_sealed_key(self):
        # fresh closeout 은 sealed-key 부재 시 여전히 SEALED_KEY_MISSING(fail-closed).
        rdir = tempfile.mkdtemp(dir=self._tmp)
        ledger = os.path.join(rdir, "l.jsonl")
        path = _write_result(rdir, "task-fresh-nokey")
        res = M.pickup_once(path, executor_key=_DEV_KEY, clock=_clock,
                            sealed_key_loader=lambda: None, ledger_path=ledger)
        self.assertEqual(res.verdict, M.PICKUP_SEALED_KEY_MISSING)

    def test_crash_after_marker_skip_terminal(self):
        rdir = tempfile.mkdtemp(dir=self._tmp)
        ledger = os.path.join(rdir, "l.jsonl")
        tid = "task-crash-after-marker"
        path = _write_result(rdir, tid)
        res1 = M.pickup_once(path, executor_key=_DEV_KEY, clock=_clock,
                             sealed_key_loader=_loader, ledger_path=ledger)
        self.assertEqual(res1.verdict, M.PICKUP_CLOSEOUT_DONE)
        # marker 존재 → 재처리 SKIP_TERMINAL.
        res2 = M.pickup_once(path, executor_key=_DEV_KEY, clock=_clock,
                             sealed_key_loader=_loader, ledger_path=ledger)
        self.assertEqual(res2.verdict, M.PICKUP_SKIP_TERMINAL)


class TestIdempotency(unittest.TestCase):
    """④ dedupe ledger + done/acked marker + 직렬 2회 → wake 1회."""

    def setUp(self):
        self._tmp = tempfile.mkdtemp(prefix="t2730-idem-")
        self.addCleanup(shutil.rmtree, self._tmp, ignore_errors=True)

    def test_serial_double_relay_wake_once(self):
        rdir = tempfile.mkdtemp(dir=self._tmp)
        ledger = os.path.join(rdir, "l.jsonl")
        tid = "task-idem-relay"
        path = _write_result(rdir, tid, relay_hints={"gemini_finding": True})
        r1 = M.pickup_once(path, executor_key=_DEV_KEY, clock=_clock,
                           sealed_key_loader=_loader, ledger_path=ledger)
        self.assertEqual(r1.verdict, M.PICKUP_WAKE_BUILT)
        r2 = M.pickup_once(path, executor_key=_DEV_KEY, clock=_clock,
                           sealed_key_loader=_loader, ledger_path=ledger)
        self.assertIn(r2.verdict, (M.PICKUP_SKIP_TERMINAL, M.PICKUP_SKIP_DEDUPE))
        self.assertFalse(r2.wake_built)
        self.assertIsNone(r2.argv)

    def test_acked_marker_skip(self):
        rdir = tempfile.mkdtemp(dir=self._tmp)
        ledger = os.path.join(rdir, "l.jsonl")
        tid = "task-idem-acked"
        path = _write_result(rdir, tid)
        with open(os.path.join(rdir, f"{tid}.pickup.acked"), "w") as fh:
            fh.write("{}")
        res = M.pickup_once(path, executor_key=_DEV_KEY, clock=_clock,
                            sealed_key_loader=_loader, ledger_path=ledger)
        self.assertEqual(res.verdict, M.PICKUP_SKIP_TERMINAL)


class TestTerminalRelayStatic(unittest.TestCase):
    """⑤ terminal_relay allowlist static enforcement."""

    def test_no_forbidden_imports(self):
        src = inspect.getsource(TR)
        tree = ast.parse(src)
        imported = set()
        for node in ast.walk(tree):
            if isinstance(node, ast.Import):
                for n in node.names:
                    imported.add(n.name.split(".")[0])
            elif isinstance(node, ast.ImportFrom):
                if node.module:
                    imported.add(node.module.split(".")[0])
        for bad in ("git", "gh", "github", "dispatch"):
            self.assertNotIn(bad, imported,
                             f"terminal_relay 가 {bad} 모듈을 import 함 (static 위반)")
        # merge/PR/dispatch 명령 문자열 + shell=True 정적 부재.
        self.assertNotIn("shell=True", src)
        self.assertNotIn("os.system", src)
        self.assertNotIn("Popen", src)

    def test_default_runner_shell_false(self):
        self.assertIn("shell=False", inspect.getsource(TR._default_runner))

    def test_send_report_fixed_argv(self):
        captured = {}
        TR.send_report("/tmp/report.md", chat_id="6937032012",
                       runner=lambda argv: captured.setdefault("argv", argv))
        self.assertEqual(captured["argv"],
                         [TR.COKACDIR_BINARY, "--sendfile", "/tmp/report.md",
                          "--chat", "6937032012"])

    def test_relay_to_anu_fixed_argv_key_runtime(self):
        captured = {}
        TR.relay_to_anu("envelope-prompt", anu_key=_ANU_KEY,
                        runner=lambda argv: captured.setdefault("argv", argv))
        argv = captured["argv"]
        self.assertEqual(argv[0], TR.COKACDIR_BINARY)
        self.assertEqual(argv[1], "--cron")
        self.assertIn("--key", argv)
        # key 는 런타임 주입값(모듈 literal 아님).
        self.assertEqual(argv[argv.index("--key") + 1], _ANU_KEY)

    def test_whitelist_rejects_bad_subcommand(self):
        with self.assertRaises(ValueError):
            TR._run_fixed_argv([TR.COKACDIR_BINARY, "--merge", "x"])
        with self.assertRaises(ValueError):
            TR._run_fixed_argv(["gh", "--sendfile", "x"])

    def test_no_anu_key_literal_in_module(self):
        self.assertNotIn(_ANU_KEY, inspect.getsource(TR),
                         "terminal_relay 모듈에 ANU key literal 존재")


class TestCollectorResultRawKeyZero(unittest.TestCase):
    """⑦ collector_result raw key 0."""

    def setUp(self):
        self._tmp = tempfile.mkdtemp(prefix="t2730-rawkey-")
        self.addCleanup(shutil.rmtree, self._tmp, ignore_errors=True)

    def test_no_raw_key_in_collector_result(self):
        for tid, hints in (("task-rk-green", None),
                           ("task-rk-relay", {"gemini_finding": True})):
            rdir = tempfile.mkdtemp(dir=self._tmp)
            ledger = os.path.join(rdir, "l.jsonl")
            path = _write_result(rdir, tid, relay_hints=hints)
            res = M.pickup_once(path, executor_key=_DEV_KEY, clock=_clock,
                                sealed_key_loader=_loader, ledger_path=ledger)
            raw = open(res.closeout_path, encoding="utf-8").read()
            self.assertNotIn(_ANU_KEY, raw, f"{tid}: collector_result 에 ANU key literal")
            self.assertNotIn(_DEV_KEY, raw, f"{tid}: collector_result 에 executor key literal")
            self.assertNotIn("--key", raw)
            self.assertNotIn("owner_key", raw)
            data = json.loads(raw)
            # 스키마 키 존재 확인.
            for k in ("task_id", "source_result_sha256", "verdict", "owner_proof",
                      "dedupe", "scope", "closeout_action", "agent_relay",
                      "contract_violation", "ts_kst", "runner_version"):
                self.assertIn(k, data)


def _drv_dirs(tmp: str):
    root = os.path.join(tmp, "root")
    for sub in ("memory/events", "memory/state", "memory/p0b_state/quarantine",
                "memory/p0b_state/processed"):
        os.makedirs(os.path.join(root, sub), exist_ok=True)
    return root


def _drv_verify_authoritative(*a, **k):
    import types
    return types.SimpleNamespace(
        verdict=M_enf.VERDICT_AUTHORITATIVE, ok=True, classification="",
        reasons=[], owner_resolution={"outcome": M_enf.OWNER_ANU,
                                      "schedule_id": "SID", "query_ok": True})


class TestDriverGreenPathLauncherZero(unittest.TestCase):
    """⑥ green-path driver: launcher_fn 호출 0 + VERDICT_CLOSEOUT_DONE + move_processed."""

    def setUp(self):
        self._tmp = tempfile.mkdtemp(prefix="t2730-drv-")
        self.addCleanup(shutil.rmtree, self._tmp, ignore_errors=True)

    def _aged_result(self, root, tid, relay_hints=None):
        rdir = os.path.join(root, "memory/events")
        env = {"task_id": tid, "schedule_id": "SID", "recorded_at": _fresh_ts()}
        path = _write_result(rdir, tid, relay_hints=relay_hints, envelope=env)
        old = _NOW.timestamp() - 3600
        os.utime(path, (old, old))
        return path

    def _pickup_with_loader(self, *a, **k):
        return M.pickup_once(*a, sealed_key_loader=_loader, clock=_clock, **k)

    def test_green_closeout_no_launcher_move_processed(self):
        root = _drv_dirs(self._tmp)
        tid = "task-drv-green"
        path = self._aged_result(root, tid)  # green
        launcher_calls = []
        rec = DRV.process_one(
            path, root=root, pickup_fn=self._pickup_with_loader,
            verify_fn=_drv_verify_authoritative,
            launcher_fn=lambda *a, **k: launcher_calls.append((a, k)),
            ledger_path=os.path.join(root, "l.jsonl"),
            clock=lambda: _NOW, sleep_fn=lambda *a, **k: None,
        )
        self.assertEqual(rec.verdict, DRV.VERDICT_CLOSEOUT_DONE)
        self.assertEqual(launcher_calls, [], "green-path 에서 launcher 가 호출됨 (wake 0 위반)")
        # result 파일이 watched(events) 밖 processed 로 이동.
        self.assertFalse(os.path.exists(path))
        self.assertTrue(os.path.isdir(os.path.join(root, "memory/p0b_state/processed")))

    def test_relay_path_uses_relay_fn(self):
        root = _drv_dirs(self._tmp)
        tid = "task-drv-relay"
        path = self._aged_result(root, tid, relay_hints={"gemini_finding": True})
        relay_calls, launcher_calls = [], []
        rec = DRV.process_one(
            path, root=root, pickup_fn=self._pickup_with_loader,
            verify_fn=_drv_verify_authoritative,
            launcher_fn=lambda *a, **k: launcher_calls.append((a, k)),
            relay_fn=lambda *a, **k: relay_calls.append((a, k)),
            ledger_path=os.path.join(root, "l.jsonl"),
            clock=lambda: _NOW, sleep_fn=lambda *a, **k: None,
        )
        self.assertEqual(rec.verdict, DRV.VERDICT_WAKE_BUILT)
        self.assertEqual(len(relay_calls), 1, "relay-path 에서 relay_fn 1회 호출")
        self.assertEqual(launcher_calls, [], "relay_fn 주입 시 launcher 미호출(relay 우선)")


class TestDriverCrashRecovery(unittest.TestCase):
    """HIGH-1/HIGH-2 remediation: driver 가 ledger-hit-no-marker 를 pickup_once recovery
    로 위임하고, marker 미확정 시 result 이동을 보류한다."""

    def setUp(self):
        self._tmp = tempfile.mkdtemp(prefix="t2730-drvrec-")
        self.addCleanup(shutil.rmtree, self._tmp, ignore_errors=True)

    def _aged_green(self, root, tid):
        rdir = os.path.join(root, "memory/events")
        env = {"task_id": tid, "schedule_id": "SID", "recorded_at": _fresh_ts()}
        path = _write_result(rdir, tid, envelope=env)
        old = _NOW.timestamp() - 3600
        os.utime(path, (old, old))
        return path

    def _pickup(self, *a, **k):
        return M.pickup_once(*a, sealed_key_loader=_loader, clock=_clock, **k)

    def test_ledger_hit_no_marker_delegates_to_recovery(self):
        import hashlib
        root = _drv_dirs(self._tmp)
        tid = "task-drv-rec"
        path = self._aged_green(root, tid)
        ledger = os.path.join(root, "l.jsonl")
        # crash-after-ledger 모사: ledger 에만 기록(marker/collector_result 부재).
        sha = hashlib.sha256(open(path, "rb").read()).hexdigest()
        with open(ledger, "w", encoding="utf-8") as fh:
            fh.write(json.dumps({"event": "PICKUP_CLOSEOUT_DONE",
                                 "task_id": tid, "sha256": sha}) + "\n")
        launcher_calls = []
        rec = DRV.process_one(
            path, root=root, pickup_fn=self._pickup,
            verify_fn=_drv_verify_authoritative,
            launcher_fn=lambda *a, **k: launcher_calls.append(1),
            ledger_path=ledger, clock=lambda: _NOW, sleep_fn=lambda *a, **k: None,
        )
        # 조기 단락(driver _dedupe_hit) 폐기 → pickup_once recovery → SKIP_DEDUPE.
        self.assertEqual(rec.verdict, DRV.VERDICT_PICKUP_SKIP)
        self.assertEqual(launcher_calls, [], "recovery 는 wake 재발사 0")
        # collector_result/marker 멱등 보강 후 파일 이동.
        self.assertFalse(os.path.exists(path))
        with open(ledger, encoding="utf-8") as fh:
            self.assertEqual(sum(1 for ln in fh if ln.strip()), 1, "ledger 중복 append 0")

    def test_marker_write_fail_defers_move(self):
        root = _drv_dirs(self._tmp)
        tid = "task-drv-mfail"
        path = self._aged_green(root, tid)
        ledger = os.path.join(root, "l.jsonl")
        real_replace = M.os.replace

        def selective(src, dst):
            if str(dst).endswith(".pickup.done"):
                raise OSError("marker fail")
            return real_replace(src, dst)

        M.os.replace = selective
        try:
            rec = DRV.process_one(
                path, root=root, pickup_fn=self._pickup,
                verify_fn=_drv_verify_authoritative,
                ledger_path=ledger, clock=lambda: _NOW, sleep_fn=lambda *a, **k: None,
            )
        finally:
            M.os.replace = real_replace
        # marker 미확정 → result 이동 보류(다음 cycle recovery).
        self.assertEqual(rec.verdict, DRV.VERDICT_CLOSEOUT_DONE)
        self.assertTrue(os.path.exists(path), "marker 미확정 시 이동 보류되어야 함")
        # 다음 cycle: os.replace 정상 → recovery → marker 보강 후 이동.
        rec2 = DRV.process_one(
            path, root=root, pickup_fn=self._pickup,
            verify_fn=_drv_verify_authoritative,
            ledger_path=ledger, clock=lambda: _NOW, sleep_fn=lambda *a, **k: None,
        )
        self.assertEqual(rec2.verdict, DRV.VERDICT_PICKUP_SKIP)
        self.assertFalse(os.path.exists(path))


class TestDryRunIsolatedFlagOff(unittest.TestCase):
    """⑧ green CLOSEOUT_DONE 도 driver_enabled flag OFF 시 미실행 + canonical 0 touch."""

    def setUp(self):
        self._tmp = tempfile.mkdtemp(prefix="t2730-dryrun-")
        self.addCleanup(shutil.rmtree, self._tmp, ignore_errors=True)

    def test_flag_off_no_closeout_canonical_untouched(self):
        root = _drv_dirs(self._tmp)
        tid = "task-dryrun"
        rdir = os.path.join(root, "memory/events")
        path = _write_result(rdir, tid)  # green
        before = sorted(os.listdir(rdir))
        # flag OFF (flag_reader → None = disabled), launcher/relay 미주입.
        records = DRV.scan_once(
            root, flag_reader=lambda: None, write_evidence=False,
            clock=lambda: _NOW, sleep_fn=lambda *a, **k: None,
        )
        self.assertEqual(len(records), 1)
        self.assertEqual(records[0].verdict, DRV.VERDICT_NOOP_DISABLED)
        # result.json 및 events 디렉토리 무변동 (closeout/marker/collector_result 0).
        self.assertEqual(sorted(os.listdir(rdir)), before)
        self.assertFalse(os.path.exists(os.path.join(rdir, f"{tid}.pickup.done")))
        self.assertFalse(os.path.exists(
            os.path.join(rdir, f"{tid}.collector_result.json")))


if __name__ == "__main__":
    unittest.main()
