"""Tests for FileSink (spens/sinks/file.py)."""

from __future__ import annotations

import json

from spens.events import Emitter
from spens.sessions import read_ndjson
from spens.sinks.file import FileSink, read_state


def _base(session_id: str = "sid") -> dict:
    return {
        "state": "building",
        "session_id": session_id,
        "exit_code": None,
        "agent_container": f"spens-agent-{session_id}",
        "interceptor_container": f"spens-interceptor-{session_id}",
        "started_at": "2026-01-01T00:00:00+00:00",
        "updated_at": "2026-01-01T00:00:00+00:00",
        "prompt": "do it",
        "env": "node-20",
        "agent": "opencode",
    }


def _emit_lifecycle(session_dir, base) -> None:
    emitter = Emitter(base["session_id"], [FileSink(session_dir, base)])
    emitter.emit("started", session_dir=str(session_dir))
    emitter.emit("build_started", image="img")
    emitter.emit("build_finished", image="img")
    emitter.emit("interceptor_starting")
    emitter.emit("interceptor_ready")
    emitter.emit("agent_started", container=f"spens-agent-{base['session_id']}")
    emitter.emit("agent_output", chunk="hello")
    emitter.emit("agent_exited", exit_code=0)
    emitter.emit("finished", exit_code=0, summary_path="/s")


def test_events_appended_in_order(tmp_path) -> None:
    session_dir = tmp_path / "sessions" / "sid"
    _emit_lifecycle(session_dir, _base())
    rows = read_ndjson(session_dir / "events.jsonl")
    assert [r["event"] for r in rows] == [
        "started", "build_started", "build_finished", "interceptor_starting",
        "interceptor_ready", "agent_started", "agent_output", "agent_exited",
        "finished",
    ]
    for row in rows:
        assert row["schema_version"] == 1
        assert row["session_id"] == "sid"


def test_state_json_transitions_through_the_lifecycle(tmp_path) -> None:
    session_dir = tmp_path / "sessions" / "sid"
    sink = FileSink(session_dir, _base())
    emitter = Emitter("sid", [sink])

    emitter.emit("started")
    assert read_state(session_dir)["state"] == "building"
    emitter.emit("interceptor_starting")
    assert read_state(session_dir)["state"] == "interceptor_starting"
    emitter.emit("agent_started", container="spens-agent-sid")
    state = read_state(session_dir)
    assert state["state"] == "running"
    assert state["agent_container"] == "spens-agent-sid"
    emitter.emit("agent_exited", exit_code=3)
    state = read_state(session_dir)
    assert state["state"] == "exited"
    assert state["exit_code"] == 3
    emitter.emit("finished", exit_code=3, summary_path="/s")
    state = read_state(session_dir)
    assert state["state"] == "finished"
    assert state["exit_code"] == 3


def test_state_json_is_valid_after_every_event_and_no_tmp_left(tmp_path) -> None:
    session_dir = tmp_path / "sessions" / "sid"
    sink = FileSink(session_dir, _base())
    emitter = Emitter("sid", [sink])
    for name in ("started", "warning", "build_started", "build_finished",
                 "agent_started", "agent_exited", "finished"):
        emitter.emit(name)
        state = read_state(session_dir)
        assert state is not None, f"state.json missing after {name}"
        assert state["session_id"] == "sid"
    # atomic rewrite: no temp file is ever left behind
    assert list(session_dir.glob("*.tmp")) == []
    # every field of the spec shape is present
    state = read_state(session_dir)
    assert set(state) >= {
        "state", "session_id", "exit_code", "agent_container",
        "interceptor_container", "started_at", "updated_at", "prompt",
        "env", "agent",
    }


def test_preexisting_terminal_state_is_never_overwritten(tmp_path) -> None:
    """The cancel race guarantee: first terminal state wins."""
    session_dir = tmp_path / "sessions" / "sid"
    session_dir.mkdir(parents=True)
    (session_dir / "state.json").write_text(json.dumps(
        {**_base(), "state": "canceled"}
    ))

    # the still-running session keeps emitting through its own FileSink
    sink = FileSink(session_dir, _base())
    emitter = Emitter("sid", [sink])
    emitter.emit("agent_exited", exit_code=137)
    emitter.emit("finished", exit_code=137, summary_path="/s")

    state = read_state(session_dir)
    assert state["state"] == "canceled"           # not overwritten
    assert state["exit_code"] == 137              # but the exit code is recorded
    # events still land in the stream
    rows = read_ndjson(session_dir / "events.jsonl")
    assert [r["event"] for r in rows] == ["agent_exited", "finished"]


def test_illegal_transition_refused(tmp_path) -> None:
    session_dir = tmp_path / "sessions" / "sid"
    sink = FileSink(session_dir, _base())
    emitter = Emitter("sid", [sink])
    emitter.emit("started")
    emitter.emit("agent_started")  # building -> running is illegal
    assert read_state(session_dir)["state"] == "building"


def test_existing_state_wins_over_base_metadata(tmp_path) -> None:
    """A background child reusing the session keeps the parent's started_at."""
    session_dir = tmp_path / "sessions" / "sid"
    sink = FileSink(session_dir, _base())
    emitter = Emitter("sid", [sink])
    emitter.emit("started")
    parent_started_at = read_state(session_dir)["started_at"]

    later_base = _base()
    later_base["started_at"] = "2099-01-01T00:00:00+00:00"
    child_sink = FileSink(session_dir, later_base)
    child = Emitter("sid", [child_sink])
    child.emit("build_started", image="img")

    state = read_state(session_dir)
    assert state["started_at"] == parent_started_at
    assert state["state"] == "building"


def test_warning_before_started_defaults_to_building(tmp_path) -> None:
    session_dir = tmp_path / "sessions" / "sid"
    sink = FileSink(session_dir, _base())
    emitter = Emitter("sid", [sink])
    emitter.emit("warning", message="[spens] Warning: x")
    state = read_state(session_dir)
    assert state["state"] == "building"


def test_truncated_final_events_line_tolerated_on_read(tmp_path) -> None:
    session_dir = tmp_path / "sessions" / "sid"
    _emit_lifecycle(session_dir, _base())
    with open(session_dir / "events.jsonl", "a", encoding="utf-8") as fh:
        fh.write('{"schema_version": 1, "event": "fin')  # torn mid-write
    rows = read_ndjson(session_dir / "events.jsonl")
    assert [r["event"] for r in rows][-1] == "finished"


def test_session_lock_is_exclusive(tmp_path) -> None:
    """A second taker blocks until the holder releases (POSIX and Windows)."""
    import threading

    from spens.sinks.file import session_lock

    session_dir = tmp_path / "sessions" / "sid"
    acquired = threading.Event()
    release = threading.Event()

    def taker() -> None:
        with session_lock(session_dir):
            acquired.set()
            release.wait(timeout=10)

    with session_lock(session_dir):
        thread = threading.Thread(target=taker, daemon=True)
        thread.start()
        assert not acquired.wait(0.5)  # still blocked while we hold the lock

    release.set()
    assert acquired.wait(10)
    thread.join(timeout=10)
