"""Tests for CodexAppServerSession — drive turns through a mock client.

The session adapter has the most complex behavior of the three new modules:
notification draining, server-request handling (approvals), interrupt,
deadline timeouts. These tests pin all of that without spawning real codex.
"""

from __future__ import annotations

import itertools
import logging
import time
from types import SimpleNamespace
from unittest.mock import patch
from typing import Any, Optional

import pytest

import agent.transports.codex_app_server_session as session_mod
from agent.transports.codex_app_server import CodexAppServerTransportError
from agent.transports.codex_app_server_session import (
    CodexAppServerSession,
    _ServerRequestRouting,
    _approval_choice_to_codex_decision,
    _build_turn_input,
)


class FakeClient:
    """Stand-in for CodexAppServerClient that records calls and lets the test
    drive the notification / server-request streams synchronously."""

    def __init__(self, *, codex_bin: str = "codex", codex_home=None) -> None:
        self.codex_bin = codex_bin
        self.codex_home = codex_home
        self.requests: list[tuple[str, dict]] = []
        self.notifications_responses: list[dict] = []
        self.responses: list[tuple[Any, dict]] = []
        self.error_responses: list[tuple[Any, int, str]] = []
        self._initialized = False
        self._closed = False
        self._notifications: list[dict] = []
        self._server_requests: list[dict] = []
        self._request_handler = None  # Optional[Callable[[str, dict], dict]]

    # API matching CodexAppServerClient
    def initialize(self, **kwargs):
        self._initialized = True
        return {"userAgent": "fake/0.0.0", "codexHome": "/tmp",
                "platformOs": "linux", "platformFamily": "unix"}

    def request(self, method: str, params: Optional[dict] = None, timeout: float = 30.0):
        self.requests.append((method, params or {}))
        if self._request_handler is not None:
            return self._request_handler(method, params or {})
        # Sensible defaults for protocol methods used by the session
        if method == "thread/start":
            return {"thread": {"id": "thread-fake-001"},
                    "activePermissionProfile": {"id": "workspace-write"}}
        if method == "turn/start":
            return {"turn": {"id": "turn-fake-001"}}
        if method == "turn/interrupt":
            return {}
        if method == "turn/steer":
            return {"turnId": (params or {}).get("expectedTurnId")}
        return {}

    def notify(self, method: str, params=None):
        pass

    def respond(self, request_id, result):
        self.responses.append((request_id, result))

    def respond_error(self, request_id, code, message, data=None):
        self.error_responses.append((request_id, code, message))

    def take_notification(self, timeout: float = 0.0):
        if self._notifications:
            return self._notifications.pop(0)
        # Honor a tiny sleep so the loop doesn't hot-spin; the real client
        # blocks on a queue. For tests we want determinism.
        if timeout > 0:
            time.sleep(min(timeout, 0.001))
        return None

    def take_server_request(self, timeout: float = 0.0):
        if self._server_requests:
            return self._server_requests.pop(0)
        return None

    def close(self):
        self._closed = True

    def is_alive(self) -> bool:
        # Fake is "alive" until close() is called; tests that want a dead
        # subprocess can patch this attribute or call close() directly.
        return not self._closed

    def stderr_tail(self, n: int = 20):
        return list(getattr(self, "_stderr_tail", []))[-n:]

    # Test helpers
    def queue_notification(self, method: str, **params):
        # Keep legacy fixture shorthand aligned with the IDs returned by the
        # fake thread/start and turn/start responses.
        if params.get("threadId") in {"t", "th"}:
            params["threadId"] = "thread-fake-001"
        if params.get("turnId") == "tu1":
            params["turnId"] = "turn-fake-001"
        turn = params.get("turn")
        if isinstance(turn, dict) and turn.get("id") == "tu1":
            turn = dict(turn)
            turn["id"] = "turn-fake-001"
            params["turn"] = turn
        self._notifications.append({"method": method, "params": params})

    def queue_server_request(self, method: str, request_id: Any = "srv-1", **params):
        self._server_requests.append({"id": request_id, "method": method, "params": params})

    def set_stderr_tail(self, lines):
        """Test helper: seed stderr_tail() output for OAuth-refresh classifier tests."""
        self._stderr_tail = list(lines)


def make_session(client: FakeClient, **kwargs) -> CodexAppServerSession:
    return CodexAppServerSession(
        cwd="/tmp",
        client_factory=lambda **kw: client,
        **kwargs,
    )


# ---- choice mapping ----

class TestApprovalChoiceMapping:
    @pytest.mark.parametrize("choice,expected", [
        ("once", "accept"),
        ("session", "acceptForSession"),
        ("always", "acceptForSession"),
        ("deny", "decline"),
        ("anything-else", "decline"),
    ])
    def test_mapping(self, choice, expected):
        assert _approval_choice_to_codex_decision(choice) == expected


class TestTurnInputCoercion:
    def test_image_parts_ride_natively_in_turn_start(self):
        """#51053: image attachments must reach the model as app-server image inputs, not a text marker."""
        items, text = _build_turn_input([
            {"type": "text", "text": "caption"},
            {"type": "image_url", "image_url": {"url": "data:image/png;base64,abc"}},
            {"type": "image_url", "image_url": {"url": "/tmp/shot.png"}},
        ])
        assert items == [
            {"type": "text", "text": "caption"},
            {"type": "image", "url": "data:image/png;base64,abc"},
            {"type": "localImage", "path": "/tmp/shot.png"},
        ]
        assert text == "caption"

    def test_image_only_turn_gets_default_prompt_and_plain_text_is_unchanged(self):
        items, text = _build_turn_input([{"type": "image_url", "image_url": {"url": "https://x/a.png"}}])
        assert text.strip()
        assert items == [{"type": "text", "text": text}, {"type": "image", "url": "https://x/a.png"}]
        assert _build_turn_input("hi") == ([{"type": "text", "text": "hi"}], "hi")


# ---- lifecycle ----

class TestLifecycle:
    def test_ensure_started_is_idempotent(self):
        client = FakeClient()
        s = make_session(client)
        tid_a = s.ensure_started()
        tid_b = s.ensure_started()
        assert tid_a == tid_b == "thread-fake-001"
        # thread/start should be called exactly once
        method_calls = [m for (m, _) in client.requests if m == "thread/start"]
        assert len(method_calls) == 1

    def test_thread_start_carries_hermes_prompt_and_disables_codex_personality(self):
        """thread/start carries cwd, Hermes' composed prompt as developerInstructions and
        personality "none" (#74712, #72104, #26035). We intentionally do NOT pass `permissions`
        (experimentalApi-gated + requires a matching config.toml [permissions] table)."""
        client = FakeClient()
        s = make_session(client, permission_profile="workspace-write", developer_instructions="SOUL: be terse")
        s.ensure_started()
        method, params = next(r for r in client.requests if r[0] == "thread/start")
        assert params == {"cwd": "/tmp", "developerInstructions": "SOUL: be terse", "personality": "none"}

    def test_thread_start_omits_developer_instructions_when_prompt_empty(self):
        """No prompt (or a blank one) never sends an empty developerInstructions field."""
        client = FakeClient()
        make_session(client, developer_instructions="   ").ensure_started()
        method, params = next(r for r in client.requests if r[0] == "thread/start")
        assert "developerInstructions" not in params
        assert params["personality"] == "none"

    def test_named_custom_provider_selects_codex_model_provider(self, monkeypatch):
        """#75186: for ``provider=custom`` + a configured ``providers.<name>`` entry, the session built by
        ``_ensure_codex_session`` sends ``model`` + ``modelProvider=<name>`` on thread/start and never the
        API key; openai/openai-codex agents send the selected model with codex's own provider."""
        import hermes_cli.runtime_provider as rp
        from agent.codex_runtime import _ensure_codex_session
        from agent.transports import codex_app_server_session as sess_mod
        monkeypatch.setattr(rp, "load_config", lambda: {
            "providers": {"my-gateway": {"api": "https://gateway.example.com/v1", "api_key": "sk-secret"}}})
        clients: list[FakeClient] = []

        def build(**kw):
            clients.append(FakeClient())
            return CodexAppServerSession(**{**kw, "client_factory": lambda **_: clients[-1]})
        monkeypatch.setattr(sess_mod, "CodexAppServerSession", build)

        def thread_start_params(**agent_attrs):
            agent = SimpleNamespace(_codex_session=None, session_cwd="/tmp", api_key="sk-secret", **agent_attrs)
            _ensure_codex_session(agent)
            agent._codex_session.ensure_started()
            return next(p for (m, p) in clients[-1].requests if m == "thread/start")

        named = thread_start_params(provider="custom", requested_provider="custom:my-gateway", model="gpt-5.4")
        # ``personality: "none"`` rides on every thread/start (#72104); only the provider selection varies.
        base = {"cwd": "/tmp", "personality": "none"}
        assert named == {**base, "modelProvider": "my-gateway", "model": "gpt-5.4"}
        assert "sk-secret" not in repr(named)
        assert thread_start_params(provider="openai-codex", requested_provider="openai-codex", model="gpt-5.4") == {
            **base, "model": "gpt-5.4"}
        assert thread_start_params(provider="custom", requested_provider="custom", model="gpt-5.4") == {**base, "model": "gpt-5.4"}
        # ``-900k`` is a Hermes-side alias the backend rejects; codex gets the base slug.
        assert thread_start_params(provider="openai-codex", requested_provider="openai-codex",
                                   model="gpt-5.6-sol-900k")["model"] == "gpt-5.6-sol"
        # The OpenAI API-key rung arrives as provider=custom with no codex model_providers id: codex's own
        # provider takes the bare slug, as openai-codex does.
        assert thread_start_params(provider="custom", requested_provider="openai", model="openai/gpt-5.4")["model"] == "gpt-5.4"

    def test_stored_thread_is_resumed_and_an_unresumable_one_falls_back_to_a_fresh_start(self):
        """#100531: a stored id goes out as ``thread/resume`` (same params as thread/start, never a
        second ``thread/start``); when codex cannot hand it back the failure is typed and the NEXT
        ensure_started() starts a fresh thread on the same handshaken client."""
        from agent.transports.codex_app_server import CodexAppServerError
        from agent.transports.codex_app_server_session import CodexThreadResumeError

        client = FakeClient()
        client._request_handler = lambda method, params: (
            {"thread": {"id": params["threadId"]}} if method == "thread/resume" else {"thread": {"id": "fresh-1"}})
        s = make_session(client, resume_thread_id="stored-1", developer_instructions="SOUL")
        assert s.ensure_started() == s.ensure_started() == "stored-1"
        assert [m for m, _ in client.requests] == ["thread/resume"]
        assert client.requests[0][1] == {"threadId": "stored-1", "cwd": "/tmp", "personality": "none", "developerInstructions": "SOUL"}

        def refuse(method, params):
            if method == "thread/resume":
                raise CodexAppServerError(code=-32600, message=f"no rollout found for thread id {params['threadId']}")
            return {"thread": {"id": "fresh-2"}}
        client = FakeClient()
        client._request_handler = refuse
        s = make_session(client, resume_thread_id="gone-1")
        with pytest.raises(CodexThreadResumeError) as exc_info:
            s.ensure_started()
        assert exc_info.value.thread_id == "gone-1"
        assert s.ensure_started() == "fresh-2"
        assert [m for m, _ in client.requests] == ["thread/resume", "thread/start"]

    def test_close_idempotent(self):
        client = FakeClient()
        s = make_session(client)
        s.ensure_started()
        s.close()
        s.close()
        assert client._closed is True


# ---- turn loop ----

class TestRunTurn:
    def test_simple_text_turn_returns_final_message(self):
        client = FakeClient()
        client.queue_notification("turn/started", threadId="t", turn={"id": "tu1"})
        client.queue_notification(
            "item/completed",
            item={"type": "agentMessage", "id": "m1", "text": "hello world"},
            threadId="t", turnId="tu1",
        )
        client.queue_notification(
            "turn/completed",
            threadId="t",
            turn={"id": "tu1", "status": "completed", "error": None},
        )
        s = make_session(client)
        r = s.run_turn("hi", turn_timeout=2.0)
        assert r.final_text == "hello world"
        assert r.interrupted is False
        assert r.error is None
        assert any(m["role"] == "assistant" and m.get("content") == "hello world"
                   for m in r.projected_messages)
        # turn_id propagated for downstream session-DB linkage
        assert r.turn_id == "turn-fake-001"



    def test_result_records_the_exact_submitted_input(self):
        client = FakeClient()
        client.queue_notification(
            "turn/completed", threadId="t",
            turn={"id": "turn-fake-001", "status": "completed", "error": None},
        )
        rich_input = [
            {"type": "text", "text": "caption"},
            {"type": "image_url", "image_url": {"url": "data:image/png;base64,abc"}},
        ]
        result = make_session(client).run_turn(rich_input, turn_timeout=2.0)
        _, params = next(request for request in client.requests if request[0] == "turn/start")
        assert result.submitted_user_text == params["input"][0]["text"]
        assert result.submitted_user_text != rich_input

    def test_turn_start_carries_the_turn_model_only_when_given(self):
        """An in-place ``/model`` switch keeps the session; the model rides on turn/start, which codex
        applies to this and later turns of the existing thread."""
        sent = []
        for model in ("gpt-5.5", None):
            client = FakeClient()
            client.queue_notification("turn/completed", threadId="t",
                                      turn={"id": "turn-fake-001", "status": "completed", "error": None})
            make_session(client).run_turn("hi", model=model, turn_timeout=2.0)
            sent.append(next(p for (m, p) in client.requests if m == "turn/start").get("model"))
        assert sent == ["gpt-5.5", None]

    def test_foreign_completion_in_server_request_drain_is_ignored(self):
        """Approval draining must not project a child result into the parent."""
        client = FakeClient()
        client.queue_server_request(
            "item/commandExecution/requestApproval",
            request_id="approval-1",
            command="pwd",
            cwd="/tmp",
        )
        client.queue_notification(
            "item/completed",
            threadId="thread-child-001",
            turnId="turn-child-001",
            item={
                "type": "agentMessage",
                "id": "child-message",
                "text": "child drain summary",
            },
        )
        client.queue_notification(
            "turn/completed",
            threadId="thread-child-001",
            turn={
                "id": "turn-child-001",
                "status": "completed",
                "error": None,
            },
        )

        original_respond = client.respond

        def respond_and_release_parent(request_id, response):
            original_respond(request_id, response)
            client.queue_notification(
                "item/completed",
                threadId="thread-fake-001",
                turnId="turn-fake-001",
                item={
                    "type": "agentMessage",
                    "id": "parent-message",
                    "text": "parent after approval",
                },
            )
            client.queue_notification(
                "turn/completed",
                threadId="thread-fake-001",
                turn={
                    "id": "turn-fake-001",
                    "status": "completed",
                    "error": None,
                },
            )

        client.respond = respond_and_release_parent
        session = make_session(
            client,
            request_routing=_ServerRequestRouting(auto_approve_exec=True),
        )

        result = session.run_turn("delegate then continue", turn_timeout=2.0)

        assert client.responses == [("approval-1", {"decision": "accept"})]
        assert result.final_text == "parent after approval"
        assert result.projected_messages == [
            {"role": "assistant", "content": "parent after approval"}
        ]



    def test_tool_iteration_counter_ticks(self):
        client = FakeClient()
        # Two completed exec items + one final agent message
        for i, item_id in enumerate(("ex1", "ex2"), start=1):
            client.queue_notification(
                "item/completed",
                item={
                    "type": "commandExecution", "id": item_id,
                    "command": f"cmd{i}", "cwd": "/tmp",
                    "status": "completed", "aggregatedOutput": "ok",
                    "exitCode": 0, "commandActions": [],
                },
                threadId="t", turnId="tu1",
            )
        client.queue_notification(
            "item/completed",
            item={"type": "agentMessage", "id": "m1", "text": "done"},
            threadId="t", turnId="tu1",
        )
        client.queue_notification(
            "turn/completed", threadId="t",
            turn={"id": "tu1", "status": "completed", "error": None},
        )
        s = make_session(client)
        r = s.run_turn("do stuff", turn_timeout=2.0)
        assert r.tool_iterations == 2
        # Each tool item produces (assistant, tool) — 2*2 + final assistant = 5 msgs
        assert len(r.projected_messages) == 5


    def test_turn_start_failure_attaches_redacted_stderr_tail(self):
        """When codex stderr has content (non-OAuth), the tail gets attached
        to the user-facing error so config/provider problems are debuggable
        instead of just 'Internal error'. Credential-shaped values in stderr
        are redacted via agent.redact(force=True); web-URL query params pass
        through (see fix(redact): pass web URLs through unchanged)."""
        client = FakeClient()
        client.set_stderr_tail([
            "ERROR: provider auth failed",
            "Authorization: Bearer sk-live-deadbeefdeadbeef",
            "url=https://api.example.com/v1?token=querysecret12345",
        ])
        from agent.transports.codex_app_server import CodexAppServerError

        def boom(method, params):
            if method == "turn/start":
                raise CodexAppServerError(code=-32603, message="Internal error")
            return {"thread": {"id": "t"}, "activePermissionProfile": {"id": "x"}}

        client._request_handler = boom
        s = make_session(client)
        r = s.run_turn("hi", turn_timeout=2.0)
        assert r.error is not None
        assert "Internal error" in r.error
        # Stderr tail attached
        assert "provider auth failed" in r.error
        # Credential-shaped values still redacted (sk- prefix + Bearer header)
        assert "sk-live-deadbeefdeadbeef" not in r.error
        # Non-OAuth → should NOT retire (subprocess JSON-RPC is still healthy).
        assert r.should_retire is False

    def test_turn_start_timeout_attaches_redacted_stderr_tail(self):
        """A non-OAuth TimeoutError on turn/start surfaces with codex stderr
        context attached and marks the session for retirement."""
        client = FakeClient()
        client.set_stderr_tail([
            "WARN: provider request stalled",
            "Authorization: Bearer sk-stalled-secret-abc123",
        ])

        def stall(method, params):
            if method == "turn/start":
                raise TimeoutError("codex method 'turn/start' timed out after 10s")
            return {"thread": {"id": "t"}, "activePermissionProfile": {"id": "x"}}

        client._request_handler = stall
        s = make_session(client)
        r = s.run_turn("hi", turn_timeout=2.0)
        assert r.error is not None
        assert "turn/start timed out" in r.error
        assert "provider request stalled" in r.error
        assert "sk-stalled-secret-abc123" not in r.error
        assert r.should_retire is True

    @pytest.mark.parametrize("op", ["run_turn", "compact_thread"])
    def test_plugin_401_stderr_keeps_primary_rpc_error_visible(self, op):
        """#75167: an ambient ChatGPT plugin prewarm 401 in codex stderr must not rewrite an
        unrelated RPC error into the `codex login` hint (turn and compaction paths alike)."""
        from agent.transports.codex_app_server import CodexAppServerError

        client = FakeClient()
        client.set_stderr_tail(["WARN ChatGPT plugin prewarm failed: HTTP 401 Unauthorized"])

        def boom(method, params):
            if method in ("turn/start", "thread/compact/start"):
                raise CodexAppServerError(code=-32603, message="internal error: workspace initialization failed")
            return {"thread": {"id": "t"}, "activePermissionProfile": {"id": "x"}}

        client._request_handler = boom
        s = make_session(client)
        r = getattr(s, op)(*(["hi"] if op == "run_turn" else []), turn_timeout=2.0)
        assert r.error is not None
        assert "workspace initialization failed" in r.error
        assert "plugin prewarm failed" in r.error  # stderr tail still attached for debugging
        assert "codex login" not in r.error
        assert r.should_retire is False


    def test_steer_appends_input_to_active_turn(self):
        client = FakeClient()
        s = make_session(client)
        s.ensure_started()
        with s._active_turn_lock:
            s._active_turn_id = "turn-live-123"

        assert s.request_steer("Use Postgres instead") is True
        method, params = client.requests[-1]
        assert method == "turn/steer"
        assert params == {
            "threadId": "thread-fake-001",
            "input": [{"type": "text", "text": "Use Postgres instead"}],
            "expectedTurnId": "turn-live-123",
        }







class TestCompactThread:
    def test_compact_thread_sends_rpc_and_waits_for_completion(self):
        client = FakeClient()
        client.queue_notification(
            "turn/started",
            threadId="thread-fake-001",
            turn={"id": "compact-turn-1"},
        )
        client.queue_notification(
            "item/completed",
            threadId="thread-fake-001",
            turnId="compact-turn-1",
            item={"type": "contextCompaction", "id": "compact-item-1"},
        )
        client.queue_notification(
            "item/completed",
            threadId="thread-fake-001",
            turnId="compact-turn-1",
            item={"type": "agentMessage", "id": "m1", "text": "compacted"},
        )
        client.queue_notification(
            "thread/tokenUsage/updated",
            threadId="thread-fake-001",
            turnId="compact-turn-1",
            tokenUsage={
                "last": {"inputTokens": 10, "outputTokens": 2, "totalTokens": 12},
                "total": {"inputTokens": 100, "outputTokens": 20, "totalTokens": 120},
                "modelContextWindow": 200000,
            },
        )
        client.queue_notification(
            "turn/completed",
            threadId="thread-fake-001",
            turn={"id": "compact-turn-1", "status": "completed", "error": None},
        )

        r = make_session(client).compact_thread(turn_timeout=2.0)

        assert ("thread/compact/start", {"threadId": "thread-fake-001"}) in client.requests
        assert r.error is None
        assert r.thread_id == "thread-fake-001"
        assert r.turn_id == "compact-turn-1"
        assert r.compacted is True
        assert r.final_text == "compacted"
        assert r.token_usage_last["totalTokens"] == 12
        assert r.model_context_window == 200000

    def test_compact_thread_ignores_foreign_child_completion(self):
        client = FakeClient()
        client.queue_notification(
            "turn/started",
            threadId="thread-child-001",
            turn={"id": "child-compact-turn"},
        )
        client.queue_notification(
            "item/completed",
            threadId="thread-child-001",
            turnId="child-compact-turn",
            item={
                "type": "agentMessage",
                "id": "child-compact-message",
                "text": "child compact summary",
            },
        )
        client.queue_notification(
            "turn/completed",
            threadId="thread-child-001",
            turn={
                "id": "child-compact-turn",
                "status": "completed",
                "error": None,
            },
        )
        client.queue_notification(
            "turn/started",
            threadId="thread-fake-001",
            turn={"id": "compact-turn-1"},
        )
        client.queue_notification(
            "item/completed",
            threadId="thread-fake-001",
            turnId="compact-turn-1",
            item={
                "type": "agentMessage",
                "id": "parent-compact-message",
                "text": "parent compacted",
            },
        )
        client.queue_notification(
            "turn/completed",
            threadId="thread-fake-001",
            turn={
                "id": "compact-turn-1",
                "status": "completed",
                "error": None,
            },
        )

        result = make_session(client).compact_thread(turn_timeout=2.0)

        assert result.error is None
        assert result.turn_id == "compact-turn-1"
        assert result.final_text == "parent compacted"
        assert result.projected_messages == [
            {"role": "assistant", "content": "parent compacted"}
        ]





# ---- approval bridge ----

class TestServerRequestRouting:



    def test_unknown_server_request_replied_with_error(self):
        client = FakeClient()
        client.queue_server_request("totally/unknown", request_id="req-3")
        client.queue_notification(
            "turn/completed", threadId="t",
            turn={"id": "tu1", "status": "completed", "error": None},
        )
        s = make_session(client)
        s.run_turn("hi", turn_timeout=0.2)
        assert any(
            rid == "req-3" and code == -32601
            for (rid, code, _msg) in client.error_responses
        )

    def test_on_event_fires_during_approval_drain(self):
        """When a server-initiated approval request arrives, the session
        drains up to 8 pending notifications first so per-turn state
        (e.g. _pending_file_changes for fileChange approvals) is current.
        Those drained notifications must also reach the on_event display
        hook — otherwise tool bubbles around approvals silently disappear.

        Regression for the issue where item/started events that landed
        in the queue alongside (or just before) an approval request got
        projected into messages but never displayed.
        """
        client = FakeClient()
        # An item/started notification is queued first, then a server
        # request — the session sees both during a single drain loop.
        client.queue_notification(
            "item/started",
            item={
                "type": "commandExecution",
                "id": "exec-1",
                "command": "echo drained",
                "cwd": "/tmp",
            },
        )
        client.queue_server_request(
            "item/commandExecution/requestApproval", request_id="req-d",
            command="echo drained",
            cwd="/tmp",
        )
        client.queue_notification(
            "turn/completed", threadId="t",
            turn={"id": "tu1", "status": "completed", "error": None},
        )

        events: list[dict] = []

        def cb(command, description, *, allow_permanent=True):
            return "once"

        s = make_session(
            client,
            approval_callback=cb,
            on_event=events.append,
        )
        s.run_turn("hi", turn_timeout=0.2)

        # The on_event hook must have seen the item/started even though
        # it was drained as part of the approval roundtrip — not just
        # events that arrive on the main notification path.
        item_started_events = [
            e for e in events
            if e.get("method") == "item/started"
        ]
        assert item_started_events, (
            "item/started drained alongside the approval was not "
            "forwarded to on_event — display will miss tool bubbles "
            "around approvals"
        )



    def test_routing_auto_approve_bypass(self):
        client = FakeClient()
        client.queue_server_request("item/commandExecution/requestApproval", request_id="r1",
                                    command="ls", cwd="/")
        client.queue_notification(
            "turn/completed", threadId="t",
            turn={"id": "tu1", "status": "completed", "error": None},
        )
        # No callback, but routing says auto-approve. Should approve.
        s = make_session(client, request_routing=_ServerRequestRouting(
            auto_approve_exec=True))
        s.run_turn("hi", turn_timeout=0.2)
        assert ("r1", {"decision": "accept"}) in client.responses



# ---- enriched approval prompts ----

class TestApprovalPromptEnrichment:
    """Quirk #4: apply_patch prompt should show what's changing.
    Quirk #10: exec prompt should never show empty cwd."""

    def test_exec_falls_back_to_session_cwd(self):
        """When codex omits cwd from the approval params, the prompt shows
        the session cwd, not an empty string."""
        client = FakeClient()
        client.queue_server_request(
            "item/commandExecution/requestApproval", request_id="r1",
            command="ls",  # no cwd
        )
        client.queue_notification(
            "turn/completed", threadId="t",
            turn={"id": "tu1", "status": "completed", "error": None},
        )
        captured = {}
        def cb(command, description, *, allow_permanent=True):
            captured["description"] = description
            return "once"
        s = make_session(client, approval_callback=cb)
        s.run_turn("hi", turn_timeout=0.2)
        # Session cwd is /tmp by default in make_session()
        assert "/tmp" in captured["description"]

    def test_apply_patch_prompt_summarizes_pending_changes(self):
        """When the projector has cached the fileChange item from item/started,
        the approval prompt surfaces the change summary."""
        client = FakeClient()
        # item/started fires first (carries the changes), then approval request
        client.queue_notification(
            "item/started",
            item={"type": "fileChange", "id": "fc-1",
                  "changes": [
                      {"kind": {"type": "add"}, "path": "/tmp/new.py"},
                      {"kind": {"type": "update"}, "path": "/tmp/old.py"},
                  ]},
            threadId="t", turnId="tu1",
        )
        client.queue_server_request(
            "item/fileChange/requestApproval", request_id="req-2",
            itemId="fc-1", turnId="tu1", threadId="t",
            startedAtMs=1234567890,
            reason="add and update files",
        )
        client.queue_notification(
            "turn/completed", threadId="t",
            turn={"id": "tu1", "status": "completed", "error": None},
        )
        captured = {}
        def cb(command, description, *, allow_permanent=True):
            captured["command"] = command
            captured["description"] = description
            return "once"
        s = make_session(client, approval_callback=cb)
        s.run_turn("hi", turn_timeout=0.2)
        # Both add and update kinds should be in the summary
        assert "1 add" in captured["command"] or "1 add" in captured["description"]
        assert "1 update" in captured["command"] or "1 update" in captured["description"]
        # And at least one of the paths
        joined = captured["command"] + " " + captured["description"]
        assert "/tmp/new.py" in joined or "/tmp/old.py" in joined

    def test_apply_patch_prompt_works_without_cached_summary(self):
        """When approval arrives before item/started (or without changes
        info), prompt falls back to whatever codex provided."""
        client = FakeClient()
        client.queue_server_request(
            "item/fileChange/requestApproval", request_id="req-2",
            itemId="fc-orphan", turnId="tu1", threadId="t",
            startedAtMs=1234567890,
            reason="apply some changes",
        )
        client.queue_notification(
            "turn/completed", threadId="t",
            turn={"id": "tu1", "status": "completed", "error": None},
        )
        captured = {}
        def cb(command, description, *, allow_permanent=True):
            captured["command"] = command
            return "once"
        s = make_session(client, approval_callback=cb)
        s.run_turn("hi", turn_timeout=0.2)
        # Falls back to the reason
        assert "apply some changes" in captured["command"]


# ---- openclaw beta.8 parity: retire/wedge/oauth/abort marker ----

class TestSessionRetirement:
    """Mirrors openclaw beta.8's resilience fixes:
      - retire timed-out app-server clients (should_retire on deadline)
      - post-tool silence is a warning, never a retirement: only a dead
        subprocess or the turn deadline retires (#112928)
      - <turn_aborted> raw marker as terminal (don't wait for turn/completed
        that never comes)
      - OAuth refresh failure classification (suggest `codex login` instead
        of raw RPC error strings)
      - dead subprocess detection between iterations
    """



    def test_final_agent_message_without_turn_completed_is_recovered(self):
        """A completed assistant item is still a usable terminal response when
        codex omits turn/completed and then goes quiet.
        """
        client = FakeClient()
        client.queue_notification(
            "item/completed",
            item={"type": "agentMessage", "id": "m1", "text": "done"},
            threadId="t",
            turnId="tu1",
        )
        s = make_session(client)
        r = s.run_turn(
            "hi",
            turn_timeout=0.05,
            notification_poll_timeout=0.01,
        )
        assert r.final_text == "done"
        assert r.interrupted is False
        assert r.error is None
        assert r.should_retire is False
        assert any(
            msg["role"] == "assistant" and msg.get("content") == "done"
            for msg in r.projected_messages
        )
        assert not any(method == "turn/interrupt" for method, _ in client.requests)


    def test_post_tool_silence_warns_but_does_not_retire_a_healthy_turn(self, caplog):
        """#112928: codex can reason for minutes after a large tool result without
        emitting a single wire event while the process stays alive. Silence past
        the quiet threshold must only warn; a later turn/completed ends the turn
        normally with no turn/interrupt and no retirement."""
        client = FakeClient()
        client.queue_notification(
            "item/completed",
            item={
                "type": "commandExecution", "id": "ex1",
                "command": "echo hi", "cwd": "/tmp",
                "status": "completed", "aggregatedOutput": "hi",
                "exitCode": 0, "commandActions": [],
            },
            threadId="t", turnId="tu1",
        )
        client.queue_notification(
            "turn/completed", threadId="t",
            turn={"id": "tu1", "status": "completed", "error": None},
        )
        s = make_session(client)
        # Clock: deadline arm, tool item at 999.0, then every later poll sits
        # well past the 0.15s quiet threshold until turn/completed is drained.
        monotonic_values = itertools.chain([1000.0, 999.0, 999.0, 999.0], itertools.repeat(1000.2))
        with caplog.at_level(logging.WARNING, logger=session_mod.logger.name), patch.object(
            session_mod.time,
            "monotonic",
            side_effect=lambda: next(monotonic_values),
        ):
            r = s.run_turn(
                "tool then silence",
                turn_timeout=5.0,
                notification_poll_timeout=0.0,
                post_tool_quiet_timeout=0.15,
            )
        assert r.interrupted is False
        assert r.should_retire is False
        assert r.error is None
        assert r.tool_iterations == 1
        assert not any(method == "turn/interrupt" for method, _ in client.requests)
        assert any("no events for" in rec.getMessage() for rec in caplog.records)

    def test_post_tool_activity_clears_the_quiet_timer_and_never_retires(self):
        """A tool completion followed by an agent message completes normally: further activity clears
        the post-tool quiet timer, and even when it expires it only warns, never retires."""
        client = FakeClient()
        client.queue_notification(
            "item/completed",
            item={
                "type": "commandExecution", "id": "ex1",
                "command": "echo hi", "cwd": "/tmp",
                "status": "completed", "aggregatedOutput": "hi",
                "exitCode": 0, "commandActions": [],
            },
            threadId="t", turnId="tu1",
        )
        # Non-tool activity immediately after — clears the quiet timer.
        client.queue_notification(
            "item/completed",
            item={"type": "agentMessage", "id": "m1", "text": "tool finished"},
            threadId="t", turnId="tu1",
        )
        client.queue_notification(
            "turn/completed", threadId="t",
            turn={"id": "tu1", "status": "completed", "error": None},
        )
        s = make_session(client)
        r = s.run_turn(
            "tool then talk", turn_timeout=2.0,
            notification_poll_timeout=0.01,
            post_tool_quiet_timeout=0.05,
        )
        # Tool ran, then text cleared the quiet timer, then turn/completed.
        # Should NOT be a retirement case.
        assert r.tool_iterations == 1
        assert r.final_text == "tool finished"
        assert r.should_retire is False
        assert r.interrupted is False







    def test_dead_subprocess_detected_between_iterations(self):
        """If codex dies (segfault, OOM, killed by its auth refresh
        thread), the inter-iteration is_alive check breaks the loop
        instead of waiting on a queue that will never fill."""
        client = FakeClient()
        s = make_session(client)
        s.ensure_started()
        # Simulate subprocess death by setting _closed (FakeClient's
        # is_alive returns False when closed).
        client._closed = True
        client.set_stderr_tail([
            "thread 'tokio-runtime-worker' panicked at 'oauth: invalid_grant'",
        ])
        r = s.run_turn("x", turn_timeout=2.0,
                       notification_poll_timeout=0.01)
        assert r.should_retire is True
        # Stderr-derived auth hint takes precedence over generic message
        assert r.error and "codex login" in r.error


# ---- thread/start cross-fill ----

class TestThreadStartCrossFill:
    """Mirrors openclaw beta.8's tolerance for thread.id/sessionId aliasing."""




    def test_missing_thread_id_raises(self):
        from agent.transports.codex_app_server import CodexAppServerError

        client = FakeClient()
        client._request_handler = lambda method, params: (
            {"thread": {}, "activePermissionProfile": {"id": "x"}}
            if method == "thread/start" else
            {"turn": {"id": "tu1"}}
        )
        s = make_session(client)
        with pytest.raises(CodexAppServerError, match="no thread id"):
            s.ensure_started()


class TestHasTurnAbortedMarker:
    """Unit coverage for the marker matcher itself."""

    def test_empty_string(self):
        from agent.transports.codex_app_server_session import (
            _has_turn_aborted_marker,
        )
        assert _has_turn_aborted_marker("") is False
        assert _has_turn_aborted_marker(None) is False  # type: ignore[arg-type]

    def test_plain_text_no_marker(self):
        from agent.transports.codex_app_server_session import (
            _has_turn_aborted_marker,
        )
        assert _has_turn_aborted_marker("normal response with no markers") is False

    def test_open_marker(self):
        from agent.transports.codex_app_server_session import (
            _has_turn_aborted_marker,
        )
        assert _has_turn_aborted_marker("blah <turn_aborted> blah") is True



class TestClassifyOAuthFailure:
    """Unit coverage for the OAuth classifier; conservative on purpose."""





    def test_empty_inputs(self):
        from agent.transports.codex_app_server_session import (
            _classify_oauth_failure,
        )
        assert _classify_oauth_failure() is None
        assert _classify_oauth_failure("") is None
        assert _classify_oauth_failure("", stderr=None) is None  # type: ignore[arg-type]

    @pytest.mark.parametrize(
        "primary, stderr, expected",
        [
            ("internal error: workspace initialization failed", "ChatGPT plugin prewarm failed: HTTP 401 Unauthorized", False),
            ("", "plugin discovery: oauth handshake returned 401 unauthorized", False),
            ("HTTP 401 Unauthorized", "", True),
            ("request body exceeded limit by 401 bytes", "", False),
            ("internal error", "token refresh failed: invalid_grant", True),
        ],
        ids=["plugin-401-stderr-keeps-rpc-error", "plugin-401-stderr-keeps-timeout", "primary-401", "bare-401-token-is-not-auth", "strong-stderr-signal"],
    )
    def test_generic_auth_words_count_only_in_primary_error(self, primary, stderr, expected):
        """#75167: ambient plugin 401/oauth stderr must not become the re-login hint; the
        primary error's own 401 Unauthorized and strong stderr credential signals still do,
        while a bare `401` token in an unrelated primary error does not."""
        from agent.transports.codex_app_server_session import _classify_oauth_failure

        assert (_classify_oauth_failure(primary, stderr=stderr) is not None) is expected


# ---- transport loss / close() racing a turn (#87422, #83127) ----

class TransportLossClient(FakeClient):
    """FakeClient whose writes fail the way the real client fails on a dead pipe."""

    def __init__(self, fail_on: str, **kwargs) -> None:
        super().__init__(**kwargs)
        self.fail_on = fail_on

    def _lost(self):
        raise CodexAppServerTransportError(code=-32000, message="codex app-server stdin closed unexpectedly: Broken pipe")

    def request(self, method, params=None, timeout=30.0):
        if method == self.fail_on:
            self._lost()
        return super().request(method, params, timeout)

    def respond(self, request_id, result):
        if self.fail_on == "respond":
            self._lost()
        super().respond(request_id, result)


class TestTransportLoss:
    def test_close_racing_turn_loop_ends_turn_as_interrupted(self):
        """#87422: close() landing between poll iterations must end run_turn()/compact_thread()
        with interrupted+should_retire, never an AttributeError on the nulled client."""
        import threading

        for run in ("run_turn", "compact_thread"):
            ready = threading.Event()
            client = FakeClient()
            client._closed = False
            polled = client.take_notification

            def take_notification(timeout=0.0, _polled=polled, _ready=ready):
                _ready.set()
                time.sleep(0.02)  # window for close() to land mid-iteration
                return _polled(timeout)

            client.take_notification = take_notification
            client.is_alive = lambda: True  # the fake stays "alive": only the session's own close() ends the turn
            session = make_session(client)
            out: dict = {}

            def worker():
                if run == "run_turn":
                    out["r"] = session.run_turn("hi", turn_timeout=5, notification_poll_timeout=0.001)
                else:
                    out["r"] = session.compact_thread(turn_timeout=5, notification_poll_timeout=0.001)

            t = threading.Thread(target=worker)
            t.start()
            assert ready.wait(timeout=5)
            session.close()
            t.join(timeout=5)
            assert not t.is_alive(), run
            result = out["r"]
            assert result.interrupted and result.should_retire, run
            assert "closed" in (result.error or ""), run

    def test_close_before_turn_loop_reports_session_closed_not_timeout(self):
        """close() landing after turn/start was accepted but before _drive_turn snapshots the
        client must retire as 'session closed', not be mislabelled a turn timeout."""
        client = FakeClient()
        client._closed = False
        session = make_session(client)
        started = session._run_started_turn

        def close_then_run(result, ts, *args):
            session.close()
            return started(result, ts, *args)

        session._run_started_turn = close_then_run
        result = session.run_turn("hi", turn_timeout=3, notification_poll_timeout=0.001)
        assert result.interrupted and result.should_retire
        assert "session closed" in (result.error or "")
        assert "timed out" not in (result.error or "")

    @pytest.mark.parametrize("run,fail_on", [
        ("run_turn", "turn/start"),
        ("compact_thread", "thread/compact/start"),
        ("run_turn", "respond"),
    ])
    def test_write_failure_returns_retiring_result_and_control_paths_stay_non_fatal(self, run, fail_on):
        """#83127: a transport write failure during turn/start, compaction start or an approval
        response returns a retiring TurnResult instead of escaping; steer/interrupt stay non-fatal."""
        client = TransportLossClient(fail_on)
        if fail_on == "respond":
            client.queue_server_request("item/commandExecution/requestApproval", command="ls")
        session = make_session(client, approval_callback=lambda *a, **k: "once")
        if run == "run_turn":
            result = session.run_turn("hi", turn_timeout=2, notification_poll_timeout=0.001)
        else:
            result = session.compact_thread(turn_timeout=2, notification_poll_timeout=0.001)
        assert result.should_retire
        assert "stdin closed unexpectedly" in result.error

        control = TransportLossClient("turn/steer")
        steer_session = make_session(control)
        steer_session.ensure_started()
        steer_session._active_turn_id = "turn-fake-001"
        assert steer_session.request_steer("more") is False
        control.fail_on = "turn/interrupt"
        steer_session._issue_interrupt("turn-fake-001")  # must not raise


def test_only_current_turn_progress_reaches_hermes_activity_clock():
    from agent.activity_tracking import ActivityTrackingMixin
    from agent.codex_runtime import make_codex_app_server_event_bridge

    agent = ActivityTrackingMixin()
    agent._touch_activity("starting new turn")
    generation = agent._turn_liveness_activity_generation
    client = FakeClient()
    for thread_id, turn_id in (
        ("thread-child-001", "turn-child-001"),
        ("thread-fake-001", "previous-turn"),
    ):
        client.queue_notification("item/agentMessage/delta", threadId=thread_id, turnId=turn_id, delta="foreign")
        client.queue_notification("item/completed", threadId=thread_id, turnId=turn_id,
                                  item={"id": "foreign-tool", "type": "commandExecution", "command": "true"})
    client.queue_notification("item/agentMessage/delta", threadId="t", turnId="tu1", delta="current")
    client.queue_notification("turn/completed", threadId="t", turn={"id": "tu1", "status": "completed"})
    result = make_session(client, on_event=make_codex_app_server_event_bridge(agent)).run_turn("work", turn_timeout=2)
    assert result.error is None and not result.interrupted
    assert agent._turn_liveness_activity_generation == generation + 1
    assert agent._last_activity_desc == "codex app-server: item/agentMessage/delta"
