"""Tests for the store-level CAS fire claim (Phase 4C).

`claim_job_for_fire` gives multi-machine at-most-once semantics when an external
scheduler (Chronos) fires a job: across N gateway replicas, exactly ONE wins the
claim for a given fire. Single-machine deployments always win (unaffected).

These exercise the real store against a temp HERMES_HOME (no mocks) per the
E2E-over-mocks discipline for file-touching code.
"""
import threading
import time

import pytest


@pytest.fixture
def temp_home(tmp_path, monkeypatch):
    """Isolated HERMES_HOME so jobs.json doesn't touch the real store."""
    monkeypatch.setenv("HERMES_HOME", str(tmp_path))
    # cron.jobs caches no home at import; get_hermes_home() reads the env live.
    yield tmp_path


def test_claim_succeeds_once_then_blocks(temp_home):
    """First claim for a fire wins; a second claim for the same fire loses, and
    next_run_at is advanced (a re-delivery for the old time can't re-fire)."""
    from cron.jobs import create_job, claim_job_for_fire, get_job

    job = create_job(prompt="x", schedule="every 5m", name="t")
    jid = job["id"]
    before = get_job(jid)["next_run_at"]

    assert claim_job_for_fire(jid) is True
    assert claim_job_for_fire(jid) is False
    assert get_job(jid)["next_run_at"] != before


def test_claim_oneshot_cannot_be_double_claimed(temp_home):
    """A one-shot can't be double-claimed (the fresh claim blocks the retry)."""
    from cron.jobs import create_job, claim_job_for_fire

    job = create_job(prompt="x", schedule="in 30m", name="o")
    assert claim_job_for_fire(job["id"]) is True
    assert claim_job_for_fire(job["id"]) is False


def test_claim_unknown_job_returns_false(temp_home):
    from cron.jobs import claim_job_for_fire

    assert claim_job_for_fire("nope-does-not-exist") is False


def test_claim_paused_job_returns_false(temp_home):
    """A paused job can't be claimed."""
    from cron.jobs import create_job, claim_job_for_fire, pause_job

    job = create_job(prompt="x", schedule="every 5m", name="p")
    pause_job(job["id"])
    assert claim_job_for_fire(job["id"]) is False


def test_forced_claim_atomically_resumes_paused_job(temp_home):
    """Explicit manual fire may resume a paused job without exposing a due
    intermediate state to the ticker."""
    from cron.jobs import create_job, claim_job_for_fire, get_job, pause_job

    job = create_job(prompt="x", schedule="every 5m", name="manual")
    pause_job(job["id"])

    assert claim_job_for_fire(job["id"], force=True) is True
    claimed = get_job(job["id"])
    assert claimed["enabled"] is True
    assert claimed["state"] == "scheduled"
    assert claimed["paused_at"] is None
    assert claimed["paused_reason"] is None
    assert claimed["fire_claim"] is not None


def test_stale_claim_is_reclaimable(temp_home, monkeypatch):
    """A claim older than the TTL is overwritten — the fire isn't stuck forever
    if the winning machine crashed before mark_job_run cleared the claim."""
    from cron.jobs import create_job, claim_job_for_fire

    job = create_job(prompt="x", schedule="every 5m", name="s")
    jid = job["id"]
    assert claim_job_for_fire(jid) is True
    # With a 0s TTL, the existing claim is always considered stale.
    assert claim_job_for_fire(jid, claim_ttl_seconds=0) is True


def test_mark_job_run_clears_claim(temp_home):
    """After a recurring job completes, its claim is cleared so the next fire
    can be claimed again."""
    from cron.jobs import create_job, claim_job_for_fire, mark_job_run, get_job

    job = create_job(prompt="x", schedule="every 5m", name="c")
    jid = job["id"]
    assert claim_job_for_fire(jid) is True
    assert get_job(jid).get("fire_claim") is not None

    mark_job_run(jid, success=True)
    assert get_job(jid).get("fire_claim") is None
    # …and the re-armed recurring job is claimable again.
    assert claim_job_for_fire(jid) is True


def test_fire_claim_heartbeat_refreshes_only_expected_owner(temp_home, monkeypatch):
    from datetime import datetime, timedelta

    import cron.jobs as jobs

    job = jobs.create_job(prompt="x", schedule="every 5m", name="heartbeat")
    assert jobs.claim_job_for_fire(job["id"]) is True
    claimed = jobs.get_job(job["id"])["fire_claim"]
    claimed_at = datetime.fromisoformat(claimed["at"])
    monkeypatch.setattr(
        jobs,
        "_hermes_now",
        lambda: claimed_at + timedelta(seconds=30),
    )

    assert jobs.heartbeat_fire_claim(
        job["id"],
        expected_owner=claimed["by"],
    ) is True
    refreshed = jobs.get_job(job["id"])["fire_claim"]
    assert refreshed["at"] != claimed["at"]
    assert refreshed["by"] == claimed["by"]
    assert jobs.heartbeat_fire_claim(
        job["id"],
        expected_owner="replacement-owner",
    ) is False


def test_reclaimed_fire_uses_new_owner_token(temp_home, monkeypatch):
    from datetime import datetime, timedelta

    import cron.jobs as jobs

    job = jobs.create_job(prompt="x", schedule="every 5m", name="reclaim")
    assert jobs.claim_job_for_fire(job["id"]) is True
    original = dict(jobs.get_job(job["id"])["fire_claim"])
    original_at = datetime.fromisoformat(original["at"])
    monkeypatch.setattr(
        jobs,
        "_hermes_now",
        lambda: original_at + timedelta(seconds=301),
    )

    assert jobs.claim_job_for_fire(job["id"]) is True
    replacement = dict(jobs.get_job(job["id"])["fire_claim"])
    assert replacement["by"] != original["by"]
    assert jobs.heartbeat_fire_claim(
        job["id"],
        expected_owner=original["by"],
    ) is False
    assert jobs.get_job(job["id"])["fire_claim"] == replacement


def test_stale_fire_owner_cannot_mark_replacement_run(temp_home):
    import cron.jobs as jobs

    job = jobs.create_job(prompt="x", schedule="every 5m", name="fenced")
    assert jobs.claim_job_for_fire(job["id"]) is True
    original = dict(jobs.get_job(job["id"])["fire_claim"])
    records = jobs.load_jobs()
    records[0]["fire_claim"] = {"at": original["at"], "by": "replacement"}
    jobs.save_jobs(records)

    assert jobs.mark_job_run(
        job["id"],
        success=True,
        expected_fire_owner=original["by"],
    ) is False
    persisted = jobs.get_job(job["id"])
    assert persisted["fire_claim"]["by"] == "replacement"
    assert persisted.get("last_run_at") is None


def test_fire_claim_fence_serializes_terminal_revocation(temp_home):
    """A side effect authorized by owner linearizes before terminal revocation."""
    from cron.jobs import (
        claim_job_for_fire,
        create_job,
        fire_claim_fence,
        mark_job_run,
    )

    job = create_job(prompt="x", schedule="every 5m", name="fenced-side-effect")
    claimed = claim_job_for_fire(job["id"], return_job=True)
    assert isinstance(claimed, dict)
    owner = claimed["fire_claim"]["by"]
    terminal_done = threading.Event()

    def finish_run():
        mark_job_run(job["id"], True, expected_fire_owner=owner)
        terminal_done.set()

    with fire_claim_fence(job["id"], expected_owner=owner) as owns_claim:
        assert owns_claim is True
        thread = threading.Thread(target=finish_run)
        thread.start()
        time.sleep(0.05)
        assert terminal_done.is_set() is False

    thread.join(timeout=1)
    assert terminal_done.is_set() is True


def test_fire_claim_fence_rejects_stale_owner(temp_home):
    from cron.jobs import claim_job_for_fire, create_job, fire_claim_fence

    job = create_job(prompt="x", schedule="every 5m", name="stale-fence")
    claim_job_for_fire(job["id"])

    with fire_claim_fence(job["id"], expected_owner="stale") as owns_claim:
        assert owns_claim is False


def test_same_process_fire_fence_refuses_second_claim_after_timeout(temp_home, monkeypatch):
    """A wedged local holder must not indefinitely block another claimant."""
    import cron.jobs as jobs

    job = jobs.create_job(prompt="x", schedule="every 5m", name="local-fence-timeout")
    monkeypatch.setattr(jobs, "_JOBS_LOCK_TIMEOUT_SECONDS", 0.1)
    completed = threading.Event()
    result = {}

    def second_claimant():
        result["claimed"] = jobs.claim_job_for_fire(job["id"])
        completed.set()

    with jobs._fire_job_lock(job["id"]) as acquired:
        assert acquired is True
        thread = threading.Thread(target=second_claimant)
        thread.start()
        assert completed.wait(timeout=2), "same-process claimant waited past the fire-fence timeout"
        assert result["claimed"] is False

    thread.join(timeout=2)
    assert thread.is_alive() is False
    assert jobs.claim_job_for_fire(job["id"]) is True


def test_same_thread_fire_fence_reentrancy_preserves_ownership(temp_home):
    """Nested same-thread callers retain the existing fire fence."""
    import cron.jobs as jobs

    job = jobs.create_job(prompt="x", schedule="every 5m", name="local-fence-reentrant")
    completed = threading.Event()
    result = {}

    def reentrant_claimant():
        with jobs._fire_job_lock(job["id"]) as outer_acquired:
            result["outer"] = outer_acquired
            with jobs._fire_job_lock(job["id"]) as inner_acquired:
                result["inner"] = inner_acquired
        completed.set()

    thread = threading.Thread(target=reentrant_claimant, daemon=True)
    thread.start()
    assert completed.wait(timeout=2), "same-thread nested fire fence did not return"
    assert result == {"outer": True, "inner": True}
    thread.join(timeout=2)
    assert thread.is_alive() is False


def test_manual_claim_does_not_stamp_a_future_occurrence(temp_home):
    """An off-tick run-now must not consume the NEXT scheduled slot.

    Outside a scheduler tick ``next_run_at`` is the occurrence that has NOT happened
    yet, so stamping it as a completed occurrence makes ``_job_is_due`` skip that slot
    when it arrives — silently, with no error and no dispatch record. ``manual=True``
    is the caller's declaration that this is an off-tick fire.
    """
    from cron.jobs import create_job, claim_job_for_fire, get_job

    job = create_job(prompt="x", schedule="every 5m", name="m")
    pending = get_job(job["id"])["next_run_at"]

    claimed = claim_job_for_fire(job["id"], manual=True, return_job=True)
    assert isinstance(claimed, dict)
    assert claimed["_scheduled_instant"] is None, (
        f"manual fire stamped the future occurrence {pending}")


def test_unclassified_off_tick_claim_does_not_stamp_a_future_occurrence(temp_home, monkeypatch):
    from datetime import datetime, timedelta

    import cron.jobs as jobs

    job = jobs.create_job(prompt="x", schedule="every 5m", name="off-tick")
    pending = jobs.get_job(job["id"])["next_run_at"]
    monkeypatch.setattr(
        jobs, "_hermes_now", lambda: datetime.fromisoformat(pending) - timedelta(minutes=1))

    claimed = jobs.claim_job_for_fire(job["id"], return_job=True)

    assert isinstance(claimed, dict)
    assert claimed["_scheduled_instant"] is None


def test_claim_seconds_before_the_slot_owns_it_once(temp_home, monkeypatch):
    """A hosted fire arriving seconds early (provider clock skew) IS the fire for the armed
    slot: it must carry the slot identity so the misfire backstop cannot run the slot again."""
    from datetime import datetime, timedelta, timezone

    import cron.executions as executions
    import cron.jobs as jobs
    from cron.occurrences import scheduled_instant

    monkeypatch.setattr(executions, "EXECUTIONS_FILE", temp_home / "cron" / "executions.db")
    job = jobs.create_job(prompt="x", schedule="0 19 * * *", name="skew")
    slot = jobs.get_job(job["id"])["next_run_at"]
    slot_dt = datetime.fromisoformat(slot)

    monkeypatch.setattr(jobs, "_hermes_now", lambda: slot_dt - timedelta(seconds=2))
    monkeypatch.setattr(executions, "_hermes_now", lambda: slot_dt - timedelta(seconds=2))
    claimed = jobs.claim_job_for_fire(job["id"], return_job=True)
    assert claimed["_scheduled_instant"] == scheduled_instant(slot)
    row = executions.create_execution(
        job["id"], source="chronos", scheduled_instant=claimed["_scheduled_instant"])
    executions.finish_execution(row["id"], success=True)
    jobs.mark_job_run(job["id"], True)

    backstop = slot_dt.astimezone(timezone.utc) + timedelta(minutes=11)
    monkeypatch.setattr(jobs, "_hermes_now", lambda: backstop)
    monkeypatch.setattr(executions, "_hermes_now", lambda: backstop)
    assert jobs.claim_job_for_fire(job["id"], return_job=True) is False, (
        "misfire backstop re-ran the slot a skewed early fire already completed")
    assert datetime.fromisoformat(jobs.get_job(job["id"])["next_run_at"]) > slot_dt


def test_manual_claim_still_refuses_a_paused_job(temp_home):
    """``manual=True`` suppresses only the occurrence stamp — unlike ``force=True`` it
    must not resume a paused job, which the run-now tool relies on to refuse it."""
    from cron.jobs import create_job, claim_job_for_fire, get_job, pause_job

    job = create_job(prompt="x", schedule="every 5m", name="mp")
    pause_job(job["id"])

    assert claim_job_for_fire(job["id"], manual=True) is False
    assert get_job(job["id"]).get("paused_at") is not None


def test_fresh_claim_from_a_dead_same_host_owner_is_reclaimable(temp_home):
    """A claim younger than the TTL whose owner pid (same host) has exited is stale at once: a
    ``hermes cron run`` killed mid-flight must not block the next manual run for the whole TTL
    with "already being fired". A live owner's fresh claim still blocks."""
    import os
    import socket
    import subprocess
    import sys

    from cron.jobs import claim_job_for_fire, create_job, load_jobs, save_jobs

    jid = create_job(prompt="x", schedule="every 5m", name="s")["id"]
    assert claim_job_for_fire(jid) is True

    # Live same-host owner (this process) → still blocked.
    jobs = load_jobs()
    job = next(j for j in jobs if j["id"] == jid)
    job["fire_claim"]["by"] = f"{socket.gethostname()}:{os.getpid()}:tok"
    save_jobs(jobs)
    assert claim_job_for_fire(jid) is False

    # Owner that has provably exited → reclaimable despite the fresh timestamp.
    child = subprocess.Popen([sys.executable, "-c", "pass"])
    child.wait()
    jobs = load_jobs()
    job = next(j for j in jobs if j["id"] == jid)
    job["fire_claim"]["by"] = f"{socket.gethostname()}:{child.pid}:tok"
    save_jobs(jobs)
    assert claim_job_for_fire(jid) is True


def test_heartbeat_does_not_wait_on_the_fence_its_own_run_holds(temp_home, monkeypatch):
    """The run thread holds the per-job fire fence across delivery; the heartbeat thread must
    refresh the claim without taking it, or every long run reads as a false ownership loss."""
    import cron.jobs as jobs

    job = jobs.create_job(prompt="x", schedule="every 5m", name="long-run")
    assert jobs.claim_job_for_fire(job["id"]) is True
    owner = jobs.get_job(job["id"])["fire_claim"]["by"]
    # Keep the pre-fix path fast: the heartbeat used to block for the full fence timeout (30s).
    monkeypatch.setattr(jobs, "_JOBS_LOCK_TIMEOUT_SECONDS", 0.2)

    fence_held, release, result = threading.Event(), threading.Event(), {}

    def hold_fence():
        with jobs.fire_claim_fence(job["id"], expected_owner=owner) as owns:
            result["owns"] = owns
            fence_held.set()
            release.wait(timeout=5)

    holder = threading.Thread(target=hold_fence, daemon=True)
    holder.start()
    try:
        assert fence_held.wait(timeout=5)
        assert jobs.heartbeat_fire_claim(job["id"], expected_owner=owner) is True
        # A genuine takeover is still detected while the fence is busy.
        assert jobs.heartbeat_fire_claim(job["id"], expected_owner="replacement-owner") is False
    finally:
        release.set()
        holder.join(timeout=5)
    assert result == {"owns": True}
    assert holder.is_alive() is False
