"""The prompt turn: ``_run_prompt_submit`` and the per-phase helpers it drives.

Bodies are rebound onto server.py's globals at install time (method_ctx.bind_module),
so they reference server.py globals bare.  Turn shape (one fresh daemon thread):
admit -> crash marker -> bind scopes -> resolve message -> run_conversation ->
commit history / message.complete -> goal & loop hooks -> release scopes ->
post-turn follow-ups (queued prompt, goal continuation, notifications).
"""

from __future__ import annotations

import dataclasses

from .method_ctx import HandlerRegistry, bind_module

_registry = HandlerRegistry()


def _bot_mode_delivery_text(response: Any, *, successful: bool) -> Any:
    """Return the text Bot Mode may render or relay after a completed turn.

    The gateway owns the canonical marker set.  Keep the row and completion
    event intact, but make a successful bare marker invisible at every Bot
    Mode delivery boundary.  Failed turns intentionally fail open so their
    diagnostic text is never swallowed.
    """
    from gateway.response_filters import is_intentional_silence_response
    return "" if successful and is_intentional_silence_response(response) else response


def _is_bot_mode_session(session: dict) -> bool:
    """Whether this completion belongs to the canonical Bot Chat surface.

    Same resolution as the system-prompt gate: the agent's title hint first (the DB
    title lands after turn 1 and ``pending_title`` is cleared once it does), then the
    live title from the session store.
    """
    from tools.bot_mode_probe import BOT_CHAT_TITLE
    hint = str(getattr(session.get("agent"), "_session_title_hint", "") or "").strip()
    if hint:  # any explicit hint decides; only an empty one costs a session-store read
        return hint == BOT_CHAT_TITLE
    return _session_live_title(session, _session_lookup_key(session)) == BOT_CHAT_TITLE


def _hook_failure(what: str, exc: BaseException) -> None:
    print(f"[tui_gateway] {what} failed: {type(exc).__name__}: {exc}", file=sys.stderr)


def _is_successful_goal_turn(result: Any, status: str, raw: Any) -> bool:
    """Whether a turn produced a real response the goal judge can use.

    A non-failed ``max_iterations_reached(...)`` handoff is a resumable turn boundary, not a
    failure: its summary must reach the judge so an active goal continues (#102213). Failed,
    interrupted and other ``completed is False`` turns still stay out (cf. #63180)."""
    from agent.turn_failure_copy import is_max_iteration_handoff
    if status != "complete" or not isinstance(raw, str) or not raw.strip():
        return False
    if not isinstance(result, dict):
        return True
    if result.get("failed") or result.get("interrupted"):
        return False
    return result.get("completed") is not False or is_max_iteration_handoff(result)


def _active_goal_manager(session: dict):
    """The session's GoalManager when a goal is active, else None."""
    from hermes_cli.goals import GoalManager
    try:
        max_turns = int((_load_cfg().get("goals") or {}).get("max_turns", 20) or 20)
    except Exception:
        max_turns = 20
    goal_mgr = GoalManager(
        session_id=str(session.get("session_key") or ""), default_max_turns=max_turns)
    return goal_mgr if goal_mgr.is_active() else None


def _plan_goal_compression_recovery(
    session: dict, result: Any, *, status: str, raw: Any) -> tuple[str | None, str | None]:
    """Bounded active-goal retry after compression exhaustion: ``(continuation, notice)``.
    Exhaustion is a failed turn (never judge input, never a spent goal turn); one fresh
    continuation is allowed, a second exhaustion pauses the goal instead of spinning."""
    if not (isinstance(result, dict) and result.get("compression_exhausted")):
        if _is_successful_goal_turn(result, status, raw):
            session.pop(_GOAL_COMPRESSION_RECOVERY_ATTEMPTS, None)
        return None, None
    if not str(session.get("session_key") or ""):
        return None, None
    if (goal_mgr := _active_goal_manager(session)) is None:
        session.pop(_GOAL_COMPRESSION_RECOVERY_ATTEMPTS, None)
        return None, None
    goal_created_at = float(getattr(goal_mgr.state, "created_at", 0.0) or 0.0)
    goal_text = getattr(goal_mgr.state, "goal", "")
    recovery_state = session.get(_GOAL_COMPRESSION_RECOVERY_ATTEMPTS)
    attempts = 0
    if (
        isinstance(recovery_state, dict)
        and recovery_state.get("goal_created_at") == goal_created_at
        and recovery_state.get("goal") == goal_text):
        with contextlib.suppress(TypeError, ValueError):
            attempts = int(recovery_state.get("attempts", 0) or 0)
    continuation_prompt = goal_mgr.next_continuation_prompt()
    if attempts < _GOAL_COMPRESSION_RECOVERY_LIMIT and continuation_prompt:
        session[_GOAL_COMPRESSION_RECOVERY_ATTEMPTS] = {
            "goal_created_at": goal_created_at, "goal": goal_text, "attempts": attempts + 1}
        return (
            continuation_prompt,
            "Context compression was exhausted. Retrying the active goal once.")
    goal_mgr.pause(reason="context compression exhausted twice consecutively")
    # A later explicit /goal resume gets a fresh bounded recovery cycle.
    session.pop(_GOAL_COMPRESSION_RECOVERY_ATTEMPTS, None)
    return None, (
        "Goal paused after context compression was exhausted twice. "
        "Run /compress, then /goal resume to continue.")


def _admit_prompt_turn(
    sid: str, session: dict, text: Any, image_paths: list[str] | None,
    queued_prompt_generation: int | None, display_kind: str | None,
    display_metadata: dict | None) -> tuple[list[str], Any] | None:
    """Ownership + liveness gate every turn source must cross; ``(images, agent)`` or None.
    Synthesized turns (auto-continue, wake-ups) call ``_run_prompt_submit`` directly — the
    bypass that once let a second backend run a duplicate turn."""
    held_lease = session.get("active_session_lease")
    # When the session already holds its lease this is a cheap dict check. See #94778.
    if not session.get("_closing") and (ownership_refusal := _ensure_active_session_slot(sid, session)) is not None:
        logger.info(
            "Refusing turn for session %s at _run_prompt_submit: %s",
            session.get("session_key") or sid,
            getattr(ownership_refusal, "reason", None) or "refused")
        with session["history_lock"]:
            session["running"] = False
            session.pop("_submit_user_row", None)  # no turn runs: the submit-time row stays as the send
        _emit("error", sid, {"message": str(ownership_refusal)})
        return None
    with session["history_lock"]:
        if session.get("_closing") or (
            queued_prompt_generation is not None
            and int(session.get("_queued_prompt_generation", 0)) != queued_prompt_generation):
            session["running"] = False
            session.pop("_submit_user_row", None)
            if session.get("_closing") and session.get("active_session_lease") is not held_lease:
                # Close stops waiting for this thread after a grace and then finalizes. A lease this
                # admission claimed after that finalize has no other code path that releases it; one
                # the session already held stays for close's own handoff (_settle_isolated_turn_before_close).
                _release_active_session_slot(session)
            return None
        images = list(session.get("attached_images", []) if image_paths is None else image_paths)
        if image_paths is None:
            session["attached_images"] = []
        inflight = session.get("inflight_turn")
        # A retained failed turn (see _fail_inflight_turn) is a stale leftover
        # by the time a new turn starts — replace it, never append onto it.
        if not isinstance(inflight, dict) or inflight.get("status") == "error":
            _start_inflight_turn(
                session, text, display_kind=display_kind, display_metadata=display_metadata)
        agent = session["agent"]
        if agent is None:
            session["running"] = False
        else:
            with contextlib.suppress(Exception):
                agent.clear_interrupt()
    if agent is None:
        # A deferred build can finish without attaching an agent (its record was replaced or closed
        # mid-build: ``agent_ready`` set, ``agent`` None, see ``_start_agent_build``).  Every turn source
        # crosses this gate, so refuse here with a retryable frame: the turn body used to dereference the
        # missing agent twice (in ``_invoke_agent`` and again in its ``finally``), which killed the turn
        # thread with ``running`` still True — the prompt vanished and the session stayed "busy" (#111531).
        reason = session.get("agent_error") or AGENT_MISSING_FOR_TURN
        logger.info("Refusing turn for session %s: no agent attached (%s)", session.get("session_key") or sid, reason)
        _emit_terminal_turn_error(
            sid, session, reason,
            error_surface={"layer": "runtime", "code": "agent_init_failed", "retryable": True})
        return None
    return images, agent


def _record_turn_marker(session: dict, text: Any, *, auto_continue: bool = True,
                        notification_category: str | None = None) -> str:
    """Write the durable crash marker; returns the session key it was written under (compression
    can rotate session_key mid-turn).  A surviving marker means the process died mid-turn.
    The key is published before the disk write so an interrupt racing startup can retire
    it; the post-write cancel check closes the inverse race (Stop landed first, no file)."""
    marker_home = _session_home(session)
    marker_key = str(session.get("session_key") or "")
    marker_attempt = int(session.pop("_auto_continue_attempt", 0) or 0)
    marker_text = session.pop("_auto_continue_prompt", None) or text
    if isinstance(marker_text, str) and marker_text.strip():
        with session["history_lock"]:
            session["_active_turn_marker_key"] = marker_key
        record_turn_start(marker_home, marker_key, marker_text, attempts=marker_attempt,
                          auto_continue=auto_continue,
                          **({"notification_category": "diagnostic"}
                             if notification_category == "diagnostic" else {}))
        with session["history_lock"]:
            marker_cancelled = bool(session.get("_turn_cancel_requested"))
        if marker_cancelled:
            clear_turn_marker(marker_home, marker_key)
    return marker_key


@dataclasses.dataclass(slots=True)
class _TurnScopes:
    """Reset tokens for the thread/context scopes a turn binds (filled incrementally)."""

    approval: Any = None
    session_tokens: list = dataclasses.field(default_factory=list)
    home: Any = None  # per-turn HERMES_HOME override for a resumed remote profile
    secret: Any = None
    terminal: Any = None


def _route_turn_images(agent, prompt: Any, images: list[str]) -> Any:
    """Run message for a turn with attached images: "native" content parts, or "text" path
    references the agent analyzes in-loop (never blocking submit on vision calls).
    Decision table: agent/image_routing.py."""
    try:
        from agent.image_routing import build_native_content_parts, decide_image_input_mode
        from hermes_cli.config import load_config as _tui_load_config
        _provider, _model = _active_image_routing_identity(agent)
        mode = decide_image_input_mode(
            _provider, _model, _tui_load_config(),
            requested_provider=getattr(agent, "requested_provider", ""))
        if getattr(agent, "api_mode", "") == "codex_app_server":
            mode = "text"
    except Exception as _img_exc:
        print(f"[tui_gateway] image_routing decision failed, defaulting to text: {_img_exc}",
              file=sys.stderr)
        mode = "text"
    if mode != "native":
        return _build_image_ref_message(prompt, images)
    try:
        parts, skipped = build_native_content_parts(prompt, images)
        if skipped:
            print(
                f"[tui_gateway] native image attachment skipped {len(skipped)} unreadable path(s)",
                file=sys.stderr)
        if any(p.get("type") == "image_url" for p in parts):
            return parts
    except Exception as _img_exc:
        print(f"[tui_gateway] native attach failed, falling back to text: {_img_exc}",
              file=sys.stderr)
    return _build_image_ref_message(prompt, images)


def _start_turn_voice() -> tuple[Any, bool]:
    """Arm voice-mode turn audio; ``(tts_queue, thinking_started)``.  ``_tts_stream_begin``
    goes first: cutting a still-speaking previous turn IS this turn's barge-in, so it must
    latch before the caller consumes the latch."""
    tts_queue = _tts_stream_begin()
    if not _voice_mode_enabled():
        return tts_queue, False
    if _voice_cfg_dict().get("barge_in", True):
        _arm_full_duplex_listener()
    try:
        from tools.voice_mode import is_audio_output_active, start_thinking_sound

        def _thinking_should_play() -> bool:
            if is_audio_output_active():
                return False
            try:
                from hermes_cli.voice import is_continuous_active
                return not is_continuous_active()
            except Exception:
                return True
        return tts_queue, start_thinking_sound(should_play=_thinking_should_play)
    except Exception:
        return tts_queue, False


def _commit_turn_history(
    session: dict, result: dict, history: list, history_version: int) -> str | None:
    """Write the agent's messages back to session history; returns a client warning or None.
    If history_version moved mid-turn, the only tolerated mutation is a gateway-inserted
    pivot marker (compare content, not indices: ``_append_model_switch_marker`` strips prior
    markers in place); any other desync is surfaced, never dropped."""
    with session["history_lock"]:
        current_version = int(session.get("history_version", 0))
        if current_version == history_version:
            session["history"] = result["messages"]
            session["history_version"] = history_version + 1
            return None
        # History mutated externally during the turn. Check if the only mutation was a pivot marker the
        # gateway itself inserted mid-turn (#76870). If so the agent output is still valid — merge it into
        # the current history that now contains the marker. A personality change counts here too: unlike a
        # model switch it has no pending queue, so `/personality` during a running turn lands immediately
        # and used to read as a genuine desync, dropping the finished turn (#82756).
        # _append_model_switch_marker strips prior markers in-place then appends a new one, so the delta is
        # NOT a simple tail-slice — we must compare content, not indices.
        current_history = list(session["history"])
        history_no_markers = [e for e in history if not _is_pivot_marker(e)]
        current_no_markers = [e for e in current_history if not _is_pivot_marker(e)]
        if current_no_markers == history_no_markers and any(
                _is_pivot_marker(e) for e in current_history):
            # Auto-compression can leave the result shorter than the turn-start history.
            msgs = result["messages"]
            new_messages = msgs[len(history):] if len(msgs) > len(history) else list(msgs)
            session["history"] = current_history + new_messages
            session["history_version"] = current_version + 1
            return None
        print(
            f"[tui_gateway] prompt.submit: history_version mismatch "
            f"(expected={history_version} current={current_version}) — "
            f"agent output NOT written to session history",
            file=sys.stderr)
        return (
            "History changed during this turn — the response above is visible "
            "but was not saved to session history.")


def _result_status(result: dict) -> str:
    return (
        "interrupted" if result.get("interrupted")
        else "error" if result.get("error") else "complete")


def _turn_outcome(result: Any, error_surface: dict | None = None) -> tuple[Any, str, str | None]:
    """Reduce a run_conversation result to ``(raw_text, status, last_reasoning)``."""
    if not isinstance(result, dict):
        return str(result), "complete", None
    raw = result.get("final_response", "")
    status = _result_status(result)
    # No visible response AND a real error: the assistant slot carries a plain account of the
    # failure (title from ``error_surface``, raw provider detail on a ``Details:`` line, next
    # step) rather than the bare provider body.  An empty successful turn still renders as empty.
    if (not raw) and result.get("error") and (result.get("failed") or result.get("partial")):
        raw = turn_error_text(result.get("error"), error_surface)
    # "Operation interrupted: waiting for model response (…)" is cancellation
    # metadata, not assistant prose (gateway/run.py and ACP suppress it too).
    # "Operation interrupted: waiting for model response (…)" is cancellation metadata, not assistant prose.
    # gateway/run.py and the ACP adapter already suppress this sentinel; without this the desktop paints it
    # as the agent's reply whenever a stop/steer lands mid-request (#7921).
    if status == "interrupted" and isinstance(raw, str) and raw.strip().startswith(
            INTERRUPT_WAITING_FOR_MODEL_PREFIX):
        raw = ""
    lr = result.get("last_reasoning")
    last_reasoning = lr.strip() if isinstance(lr, str) and lr.strip() else None
    return raw, status, last_reasoning


def _goal_followup_after_turn(
    sid: str, session: dict, result: Any, status: str, raw: Any) -> str | None:
    """/goal continuation (mirrors gateway/run._post_turn_goal_continuation): the prompt to
    chain once ``running`` is released, or None.  Compression failures are never judge
    input: the error text is not work toward the goal, and judging it spends a turn."""
    goal_followup = None
    compression_exhausted = bool(isinstance(result, dict) and result.get("compression_exhausted"))
    try:
        recovery_prompt, recovery_notice = _plan_goal_compression_recovery(
            session, result, status=status, raw=raw)
        if recovery_notice:
            from gateway.warning_notifications import render_notification
            render_notification(
                lambda: _emit("status.update", sid, {"kind": "goal", "text": recovery_notice}),
                platform="tui", user_config=getattr(session.get("agent"), "_notification_config", None))
        goal_followup = recovery_prompt or None
    except Exception as _goal_recovery_exc:
        _hook_failure("goal compression recovery", _goal_recovery_exc)
    if compression_exhausted or not _is_successful_goal_turn(result, status, raw):
        return goal_followup
    try:
        if session.get("session_key") and (goal_mgr := _active_goal_manager(session)) is not None:
            _active_deleg = 0
            try:
                from hermes_cli.goals import count_active_delegations, gather_background_processes as _gather_bg
                # Only THIS session's processes (TUI turns register under session_key): subagents'
                # pollers must not park the parent's goal. Same rule as the CLI and gateway loops.
                _bg_procs = _gather_bg(owner_task_id=session.get("session_key") or None)
                _active_deleg = count_active_delegations(getattr(session.get("agent"), "session_id", None))
            except Exception:
                _bg_procs = None
            decision = goal_mgr.evaluate_after_turn(
                raw, user_initiated=True, background_processes=_bg_procs, active_delegations=_active_deleg)
            if verdict_msg := decision.get("message") or "":
                _emit("status.update", sid, {"kind": "goal", "text": verdict_msg})
            if decision.get("should_continue") and (
                cont_prompt := decision.get("continuation_prompt") or ""):
                goal_followup = cont_prompt
    except Exception as _goal_exc:
        _hook_failure("goal continuation hook", _goal_exc)
    return goal_followup


def _after_complete_turn(sid: str, session: dict, st: _TurnRun, raw: Any) -> None:
    """Hooks for a ``complete`` turn: /loop tick evaluation, pending title, voice fallback."""
    try:
        from hermes_cli.loops import LoopManager
        loop_sid_key = session.get("session_key") or ""
        if loop_sid_key:
            loop_mgr = LoopManager(session_id=loop_sid_key)
            loop_state = loop_mgr.state
            if loop_state is not None and loop_state.awaiting_response:
                loop_decision = loop_mgr.complete_tick(raw if isinstance(raw, str) else "")
                if loop_msg := loop_decision.get("message") or "":
                    _emit("status.update", sid, {"kind": "loop", "text": loop_msg})
    except Exception as _loop_exc:
        _hook_failure("loop completion hook", _loop_exc)
    # Apply pending_title now that the DB row exists — in the session-owned profile store.
    if _pending := session.get("pending_title"):
        _session_key = session.get("session_key") or sid
        try:
            with _session_db(session) as _pdb:
                if _pdb and _pdb.set_session_title(_session_key, _pending):
                    session["pending_title"] = None
        except ValueError as exc:
            # Invalid/duplicate title — non-retryable, drop it; auto-title takes over.
            session["pending_title"] = None
            logger.info("Dropping pending title for session %s: %s", _session_key, exc)
        except Exception:
            pass  # transient DB failure — keep pending_title for retry
    # Voice fallback when the streaming pipeline couldn't start (tts_queue already spoke
    # everything otherwise); barge-aware.
    if st.tts_queue is None and isinstance(raw, str) and raw.strip() and _voice_tts_enabled():
        try:
            threading.Thread(target=_speak_text_with_barge, args=(raw,), daemon=True).start()
        except ImportError:
            logger.warning("voice TTS skipped: hermes_cli.voice unavailable")
        except Exception as e:
            logger.warning("voice TTS dispatch failed: %s", e)


def _dispatch_followup_turn(rid, sid: str, session: dict, prompt: Any, what: str, *,
                            on_done=None, on_error=None) -> None:
    """Chain one follow-up turn (caller set ``running``); on failure run ``on_error``, log,
    release ``running``."""
    try:
        _emit("message.start", sid)
        _run_prompt_submit(rid, sid, session, prompt)
        if on_done is not None:
            on_done()
    except Exception as exc:
        if on_error is not None:
            on_error()
        _hook_failure(what, exc)
        with session["history_lock"]:
            session["running"] = False


def _run_post_turn_followups(
    rid, sid: str, session: dict, result: Any, goal_followup: str | None) -> None:
    """Chain whatever should run after ``running`` was released.  Order: a mid-turn user
    prompt wins over every auto follow-up (drain it, skip the rest); a leftover /steer is
    requeued first so it isn't dropped; then goal continuation, then completion
    notifications.  Each nested submit re-checks ``running`` under the lock."""
    steer = result.get("pending_steer") if isinstance(result, dict) else None
    if isinstance(steer, str) and steer.strip():
        with session["history_lock"]:
            _enqueue_prompt(session, steer, session.get("transport"))
    if _drain_queued_prompt(rid, sid, session):
        return
    if goal_followup:
        with _session_turn_admission(session) as admitted:
            if not admitted or session.get("running"):
                return  # user already sent something — their turn wins
            if session.get("_turn_cancel_requested"):
                return  # the user pressed Stop; the goal resumes after their next prompt
            session["running"] = True
        _dispatch_followup_turn(rid, sid, session, goal_followup, "goal continuation dispatch")
    # Safety net for completion events that arrived mid-turn.  Ownership is positive-proof
    # and compression-chain aware (same fail-closed gate as the poller): session B must
    # not consume session A's event.  Unclaimable events are requeued for the poller.
    try:
        from tools.process_registry import process_registry
        # _finish_turn has released the worker's runtime scope. Queue ownership,
        # notification policy and nested dispatch still belong to this session.
        with _session_profile_runtime_scope(session):
            drained = process_registry.drain_notifications(
                session_key=session.get("session_key", ""),
                owns_event=lambda e: _session_owns_notification_event(sid, session, e),
                skip_poll_observed=False)
            from tools.process_registry_notifications import format_process_notification
            deferred = []
            _notif_handle_ready(
                sid, session, [event for event, _text in drained],
                session.setdefault("_notification_emitted", set()), process_registry,
                format_process_notification, deferred, owned=True)
            for event in deferred:
                process_registry.completion_queue.put(event)
    except Exception as _drain_exc:
        _hook_failure("completion queue drain", _drain_exc)


@dataclasses.dataclass(slots=True)
class _TurnRun:
    """Shared state of one turn thread.  ``agent`` is bound eagerly so except/finally always
    have one; ``error_retained`` makes the finally keep the failed inflight snapshot for
    resume replay; ``error_detail`` is the "tui turn finished" failure cause."""

    agent: Any
    one_turn_restore: Any
    terminal_callback: Any
    receipt_committed: bool
    scopes: _TurnScopes = dataclasses.field(default_factory=_TurnScopes)
    result: Any = None  # read after the finally for leftover /steer
    tts_queue: Any = None
    thinking_started: bool = False
    history: list = dataclasses.field(default_factory=list)
    history_version: int = 0
    compression_count: int | None = None
    run_kwargs: Any = None
    error_retained: bool = False
    error_detail: str = ""
    prompt_text: str = ""
    marker_key: str = ""
    receipt_attempted: bool = False


def _adopt_out_of_band_turns(session: dict) -> None:
    """Fold turns another surface appended to this session (Telegram reply, cron run) into the model-facing
    history before the turn snapshots it. The desktop already repaints them from the DB (#86588); without
    this the next prompt still ran on the in-memory history and the model never saw them (#42962).
    Foreign rows are the active rows between the highest ``_row_id`` the agent's own flushes stamped onto
    the in-memory messages (``sync_flushed_message_markers``; a local compaction re-stamps them too) and
    this turn's own user row, which ``_persist_submit_user_row`` already wrote (#111868). When a foreign
    row is a compaction summary the other surface rewrote the transcript under us, so the in-memory history
    is stale from the root and is re-hydrated from the DB the way ``session.resume`` does. Nothing stamped
    yet (seeded branch before its first turn) means nothing to adopt; a cold resume arrives stamped — but a
    record whose list was rebuilt from provider-format messages carries no stamp at all, so the row-id
    boundary cannot key on anything and the store-ahead fallback takes over (#81951)."""
    with session["history_lock"]:
        history, version = list(session.get("history") or ()), int(session.get("history_version", 0))
    seen = max((rid for m in history if isinstance(m, dict) and (rid := _message_row_id(m)) is not None),
               default=None)
    ceiling = _message_row_id(session.get("_submit_user_row") or {})

    def _below_ceiling(rid) -> bool:
        return isinstance(rid, int) and (ceiling is None or rid < ceiling)
    if seen is None:
        _adopt_store_tail_without_row_ids(session, history, version, _below_ceiling)
        return

    def _foreign(rid) -> bool:
        return _below_ceiling(rid) and rid > seen
    # Keyset probe first: the common turn has nothing to adopt and must not pay a full transcript decode.
    # Address the live session, not session_key: a rotation moves the tip mid-session, so the parent
    # holds none of the foreign rows (Telegram reply, cron) this function exists to adopt and the model
    # silently never sees them (#123545).
    with _session_db(session) as db:
        try:
            newer = db.get_messages(_submit_row_target_key(session), after_id=seen) if db is not None else []
        except Exception:
            logger.debug("out-of-band history probe failed; turn runs on the in-memory history", exc_info=True)
            return
    newer = [row for row in newer if _foreign(row.get("id"))]
    if not newer:
        return
    rewritten = any(row.get("_compressed_summary") for row in newer)
    rows = _load_durable_truncation_history(session, repair_alternation=rewritten) or []
    keep = _below_ceiling if rewritten else _foreign
    tail = canonicalize_replay_history([m for m in rows if keep(_message_row_id(m))])
    if not tail:
        return
    with session["history_lock"]:
        if int(session.get("history_version", 0)) != version:
            return  # /compress, a rewind or a pivot marker landed meanwhile; the next turn re-derives
        session["history"] = tail if rewritten else history + tail
        session["history_version"] = version + 1


def _adopt_store_tail_without_row_ids(session: dict, history: list, version: int, below_ceiling) -> None:
    """Catch a live record up with the store when its in-memory messages carry NO durable ``_row_id``.

    ``_adopt_out_of_band_turns`` keys off the highest ``_row_id`` in memory; a record whose list was
    rebuilt from provider-format messages (the unstamped shape ``_resolve_truncate_row_id`` heals for
    rewind — #82959) carries none. A writer in ANOTHER process (the messaging gateway appending to the
    same session while the desktop's serve record holds it) was then silently dropped and the next
    provider request truncated to the early snapshot while the UI showed the full transcript (#81951).

    With no row ids there is no boundary to key on, so adoption is gated on proof instead: the durable
    lineage (below this turn's own user row, which ``_persist_submit_user_row`` already wrote and the turn
    appends itself) must be STRICTLY longer and its head must equal the in-memory view entry for entry.
    A diverged view (prefix mismatch, e.g. a rewind that has not landed locally), a store that is not
    ahead (an unflushed local tail is the fresher record) and an empty in-memory history (no head to
    prove against) are left alone.

    A compaction by another surface is the one rewrite that is NOT an append: the store then holds a
    ``_compressed_summary`` row the in-memory view has never seen, so the view is stale from the root and
    is re-hydrated from the DB like the stamped path does — a positional check alone would keep it because
    the compacted store is shorter.
    """
    rows = [m for m in _load_durable_truncation_history(session, repair_alternation=True) or []
            if below_ceiling(_message_row_id(m))]
    if not history:
        return
    known = {m.get("content") for m in history if isinstance(m, dict) and m.get("_compressed_summary")}
    if any(m.get("_compressed_summary") and m.get("content") not in known for m in rows):
        tail = canonicalize_replay_history(rows)
        with session["history_lock"]:
            if tail and int(session.get("history_version", 0)) == version:
                session["history"] = tail
                session["history_version"] = version + 1
        return
    if len(rows) <= len(history):
        return
    if not all(isinstance(mem, dict) and mem.get("role") == stored.get("role")
               and mem.get("content") == stored.get("content")
               for mem, stored in zip(history, rows)):
        return
    tail = canonicalize_replay_history(rows[len(history):])
    if not tail:
        return
    with session["history_lock"]:
        if int(session.get("history_version", 0)) != version:
            return  # a local rewrite landed meanwhile; the next turn re-derives
        session["history"] = history + tail
        session["history_version"] = version + 1


def _install_has_prior_sessions(session: dict) -> bool:
    """True when this install already has session rows beyond the current one.

    Mirrors ``gateway.session.SessionStore.has_any_sessions`` (the messaging
    first-contact gate): ``_run_prompt_submit`` persists the session's own row
    before the turn runs (``_ensure_session_db_row``), so a fresh install on
    its first-ever message holds exactly one row.
    """
    try:
        with _session_db(session) as db:
            return db is not None and db.session_count_ge(2)
    except Exception:
        logger.debug("session count probe failed for first-contact check", exc_info=True)
        return False


def _stage_first_contact_onboarding_note(session: dict, agent, history_empty: bool) -> None:
    """Stage the install's first-message onboarding note for THIS turn (#82750).

    The messaging gateway appends the consent-gated profile-build directive to
    the very first message ever (``_hmwa_first_contact_notes``); the
    TUI/Desktop surface never did, so a fresh install's first Desktop chat
    skipped the opt-in profile flow entirely. Stage the same note through
    ``agent._gateway_turn_context_notes`` — consumed by
    ``agent.turn_context`` on the user message — never the ephemeral system
    prompt, which must stay byte-stable for the conversation (prompt-cache
    invariant). Fires at most once per install: the directive path persists
    ``onboarding.seen.profile_build_offered`` before the turn runs.
    """
    try:
        from agent.onboarding import first_contact_turn_note
        from hermes_cli.config import load_config as _load_onboarding_config
        from hermes_constants import get_hermes_home

        note = first_contact_turn_note(
            _load_onboarding_config() or {},
            get_hermes_home() / "config.yaml",
            session_history_empty=history_empty,
            install_has_prior_sessions=_install_has_prior_sessions(session),
        )
        if not note:
            return
        prior = getattr(agent, "_gateway_turn_context_notes", "") or ""
        agent._gateway_turn_context_notes = f"{prior}\n\n{note}" if prior else note
    except Exception:
        logger.debug("first-contact onboarding note failed", exc_info=True)


def _prepare_turn_input(sid: str, session: dict, st: _TurnRun, text: Any, images: list[str]):
    """Bind scopes, sync the agent, snapshot history, build the run message; returns
    ``(prompt, run_message, cols, streamer)`` or None when @-expansion was refused.
    Scopes fill field by field so a failure midway still leaves every bound token for the
    finally; the profile's terminal policy is bound too (a failed install leaves a
    fail-closed refusal scope).  The config-model sync is skipped under a /model --once
    override (not pinned as model_override, the sync would clobber it); a model picked
    mid-turn is applied first so the explicit pick wins over a config change."""
    from tools.approval_context import set_current_session_key
    scopes = st.scopes
    scopes.approval = set_current_session_key(session["session_key"])
    scopes.session_tokens = _set_session_context(session["session_key"], ui_session_id=sid)
    # Profile turn: that profile's home + secrets + terminal policy. Launch-profile turn: unscoped in a
    # single-profile process; once multiplexing is active (#68559 / #107422 residual) its OWN scope,
    # built from the env frozen at activation — get_secret() fails closed then, so an unscoped default
    # member's hosted-room turn otherwise died with UnscopedSecretError, and ambient TERMINAL_* a
    # secondary context poisoned must never become the launch turn's authority.
    bound = _profile_runtime_scope_tokens(session.get("profile_home"))
    if bound is not None:
        scopes.home, scopes.secret, scopes.terminal = bound.home, bound.secret, bound.terminal
    # The sudo password callback is thread-local: without re-wiring here, sudo prompts
    # fall through to /dev/tty and hang the headless gateway (re-run is a no-op).
    _wire_callbacks(sid)
    if not st.one_turn_restore:
        # Skip the config-model sync while a /model --once override is active: the once-model is
        # intentionally not pinned as a session model_override (it must not persist), so without this guard
        # the sync would see "agent model != config model" and clobber the once-override back to the config
        # model before the turn runs (#29923 review defect). Any config.yaml change is adopted on the NEXT
        # turn, after the finally-restore below.
        _apply_pending_model_switch(sid, session)
        _sync_agent_model_with_config(sid, session)
        _sync_agent_compression_with_config(sid, session)
    _sync_agent_fallback_with_config(sid, session)  # chain added after the chat opened reaches this turn
    _sync_bot_capabilities(sid, session)  # Bot Chat: adopt Settings->Capabilities edits
    _adopt_out_of_band_turns(session)
    st.agent = agent = session["agent"]
    # Snapshot after the model sync: a deferred switch's history mutation belongs to this turn.
    with session["history_lock"]:
        st.history = list(session["history"])
        st.history_version = int(session.get("history_version", 0))
    # Install-first-message onboarding (#82750): gateway parity for the TUI/Desktop
    # surface — no-op unless this is the install's very first message ever.
    _stage_first_contact_onboarding_note(session, agent, not st.history)
    cwd = _session_cwd(session)
    _register_session_cwd(session)
    cols = session.get("cols", 80)
    streamer = make_stream_renderer(cols)
    prompt = text
    if isinstance(prompt, str) and "@" in prompt:
        from agent.context_references import preprocess_context_references
        from agent.model_metadata import get_model_context_length
        ctx_len = get_model_context_length(
            getattr(agent, "model", "") or _resolve_model(),
            base_url=getattr(agent, "base_url", "") or "",
            api_key=getattr(agent, "api_key", "") or "",
            provider=getattr(agent, "provider", "") or "",
            config_context_length=getattr(agent, "_config_context_length", None))
        ctx = preprocess_context_references(
            prompt, cwd=cwd, allowed_root=cwd, context_length=ctx_len)
        if ctx.blocked:
            _emit(
                "error", sid, {"message": "\n".join(ctx.warnings) or "Context injection refused."})
            return None
        prompt = ctx.message
    st.prompt_text = prompt if isinstance(prompt, str) else ""
    run_message: Any = _route_turn_images(agent, prompt, images) if images else prompt
    from agent.notification_presentation import event_presentation_muted
    if not event_presentation_muted("message.delta", sid):
        st.tts_queue, st.thinking_started = _start_turn_voice()
    # Per-turn API-message notes: barge mid-speech, reactions, HUD surface (per-turn state
    # that must not touch the byte-stable system prompt).
    from tools.tts_streaming import SPEECH_INTERRUPTED_NOTE, take_speech_interrupted
    if take_speech_interrupted():
        run_message = _prepend_note(run_message, SPEECH_INTERRUPTED_NOTE)
    run_message = _prepend_note(run_message, _pending_reaction_notes(session))
    return prompt, _prepend_note(run_message, _hud_surface_note(session)), cols, streamer


def _invoke_agent(
    sid: str, session: dict, st: _TurnRun, prompt: Any, run_message: Any, streamer,
    images: list[str], display_kind: str | None, display_metadata: dict | None,
    turn_author: dict | None = None, text: Any = None) -> None:
    """Wire the streaming callbacks and run the conversation into ``st.result``.
    ``text`` is the turn's raw submit, matched against the row staged by prompt.submit."""
    agent = st.agent
    st.compression_count = getattr(getattr(agent, "context_compressor", None), "compression_count", None)
    # Bot Chat mirrors gateway.stream_consumer: deltas are withheld while the streamed buffer
    # could still resolve to a silence marker ("NO"->"NO_REPLY"), so a bare marker is never
    # shown and then retracted (the client keeps streamed text when message.complete is "").
    hold = {"buf": "", "held": ""} if _is_bot_mode_session(session) else None

    def _stream(delta):
        if getattr(agent, "_mute_notification_reply", False):
            return
        if hold is not None and isinstance(delta, str):
            from gateway.response_filters import is_partial_silence_marker
            hold["buf"] += delta
            if is_partial_silence_marker(hold["buf"]):
                hold["held"] += delta
                return
            delta, hold["held"] = hold["held"] + delta, ""
        with session["history_lock"]:
            _append_inflight_delta(session, delta)
        payload = {"text": delta}
        if streamer and (r := streamer.feed(delta)) is not None:
            payload["rendered"] = r
        if st.tts_queue is not None and isinstance(delta, str):
            st.tts_queue.put(delta)
        _emit("message.delta", sid, payload)

    # Interim assistant text (commentary beside tool calls, pre-nudge final answer) is sealed
    # by the desktop as its own segment instead of being lost to message.complete.
    def _interim_assistant_cb(text: str, *, already_streamed: bool = False) -> None:
        if getattr(agent, "_mute_notification_reply", False):
            return
        _emit("message.interim", sid, {"text": text, "already_streamed": already_streamed})
    agent.interim_assistant_callback = (
        _interim_assistant_cb if _load_interim_assistant_messages() else None)
    # A synthesized turn is typed at turn START so a crash persist writes a timeline event,
    # not a raw user bubble; the post-turn stamp is the fallback for an older agent.
    st.run_kwargs = run_kwargs = {
        "conversation_history": list(st.history),
        "stream_callback": _stream,
        "persist_user_message": (
            _build_persist_user_message(prompt, images, run_message) if images else prompt)}
    try:
        run_params = inspect.signature(agent.run_conversation).parameters
    except (TypeError, ValueError):
        run_params = {}
    if "task_id" in run_params:
        run_kwargs["task_id"] = session["session_key"]
    if display_kind and "persist_user_display_kind" in run_params:
        run_kwargs["persist_user_display_kind"] = display_kind
    if display_metadata and "persist_user_display_metadata" in run_params:
        run_kwargs["persist_user_display_metadata"] = display_metadata
    if turn_author and "turn_author" in run_params:
        run_kwargs["turn_author"] = turn_author
    _adopt_submit_user_row(session, agent, run_kwargs["persist_user_message"], text)
    # Live-rename hook: auto-titling fires inside the turn prologue.
    _title_key = session.get("session_key") or sid
    agent._on_session_title = lambda t, _src, _k=_title_key: _emit(
        "session.title", sid, {"session_id": _k, "title": t})
    _usage_stop, _usage_thread = _start_usage_ticker(sid, agent)
    try:
        from agent.notification_presentation import notification_turn, event_presentation_muted
        with notification_turn(agent, muted=event_presentation_muted("message.delta", sid), session_id=sid):
            st.result = agent.run_conversation(run_message, **st.run_kwargs)
    finally:
        # Stop AND join before anything emits: a tick surviving past message.complete would
        # roll the client's usage back to a stale snapshot (unbounded join: same worst case).
        _usage_stop.set()
        _usage_thread.join()


def _absorb_turn_result(
    sid: str, session: dict, st: _TurnRun, text: Any, display_kind: str | None, display_metadata
) -> str | None:
    """Stamp, restore /moa, commit history, re-sync the session key; returns the history warning."""
    result, agent = st.result, st.agent
    if display_kind and isinstance(text, str):
        # Post-turn fallback stamp of a synthesized turn's display kind (DB row + result).
        db = getattr(agent, "_session_db", None)
        current_session_id = getattr(agent, "session_id", None) or session.get("session_key")
        if db is not None:
            try:
                db.set_latest_matching_message_display_kind(
                    current_session_id, role="user", content=text, display_kind=display_kind,
                    display_metadata=display_metadata)
            except Exception:
                logger.debug("failed to stamp synthetic display kind", exc_info=True)
        if isinstance(result, dict) and isinstance(result.get("messages"), list):
            for message in reversed(result["messages"]):
                if message.get("role") == "user" and message.get("content") == text:
                    message["display_kind"] = display_kind
                    if display_metadata:
                        message["display_metadata"] = display_metadata
                    break
    if "moa_one_shot_restore" in session:
        # Undo a /moa one-shot through the switch path: resetting model_override alone
        # would leave the live client pinned to MoA after the in-place switch_model().
        _restore = session.pop("moa_one_shot_restore", None)
        # Restore the model the user was on before the /moa one-shot. See #53444.
        if isinstance(_restore, dict):
            _prev_override = _restore.get("override")
            _prev_model = _restore.get("model")
            _prev_provider = _restore.get("provider")
            if _prev_override is None:
                session.pop("model_override", None)
            else:
                session["model_override"] = _prev_override
            if _prev_model:
                _raw = (
                    f"{_prev_model} --provider {_prev_provider}" if _prev_provider else _prev_model)
                try:
                    _apply_model_switch(
                        sid, session, _raw, confirm_expensive_model=False,
                        pin_session_override=bool(_prev_override),
                        persist_override=False, count_switch=False)  # session-internal restore, never config.yaml
                except Exception as _moa_restore_exc:
                    logger.warning("MoA one-shot model restore failed: %s", _moa_restore_exc)
        elif _restore is None:
            session.pop("model_override", None)
        else:
            session["model_override"] = _restore
    status_note = None
    if isinstance(result, dict):
        if isinstance(result.get("messages"), list):
            status_note = _commit_turn_history(session, result, st.history, st.history_version)
        # Auto-compression may have rotated agent.session_id: sync session_key before
        # title/goal/finalize use it, keep pending_title (user intent), restart the slash
        # worker so worker-backed commands target the live session.
        # Fix for #20001.
        _sync_session_key_after_compress(
            sid, session, clear_pending_title=False, restart_slash_worker=True)
    return status_note


def _persisted_turn_receipt(st: _TurnRun, raw: Any, status: str) -> dict | None:
    """Report committed row addresses, never a text/timestamp search for a matching turn.

    The agent's persistence cursor is re-anchored during compaction. A partial receipt can
    address its surviving rows, but cannot retire a client's entire streamed turn. Full coverage
    additionally requires the unchanged pre-turn prefix and no redirected user boundary.
    """
    from agent.context_compressor import _DB_PERSISTED_MARKER

    messages = st.result.get("messages")
    start = getattr(st.agent, "_persist_user_message_idx", None)
    if (not isinstance(messages, list) or type(start) is not int or not 0 <= start < len(messages)
            or messages[start].get("role") != "user"):
        return None

    def committed_id(message):
        row_id = message.get("_row_id")
        return row_id if message.get(_DB_PERSISTED_MARKER) and type(row_id) is int and row_id > 0 else None

    tail = messages[start:]
    anchor_id = committed_id(tail[0])
    if any(before is tail[0] or (anchor_id is not None and anchor_id == committed_id(before))
           for before in st.history):
        return None  # preflight returned the old transcript, not a new turn
    row_ids = [rid for message in tail if (rid := committed_id(message)) is not None]
    if not row_ids:
        return None
    receipt = {"row_ids": row_ids, "complete": False}
    if (user_row_id := committed_id(tail[0])) is not None:
        receipt["user_row_id"] = user_row_id
    last = tail[-1]
    # Equality only verifies the structurally selected final row's body: it never selects an identity.
    if (status == "complete" and last.get("role") == "assistant" and not last.get("tool_calls")
            and last.get("content") == raw and (final_id := committed_id(last)) is not None):
        receipt["final_assistant_row_id"] = final_id
    prefix_unchanged = start == len(st.history) and all(
        committed_id(before) is not None and committed_id(before) == committed_id(after)
        for before, after in zip(st.history, messages[:start]))
    receipt["complete"] = bool(
        prefix_unchanged and type(st.compression_count) is int
        and st.compression_count == getattr(getattr(st.agent, "context_compressor", None), "compression_count", None)
        and len(row_ids) == len(tail)
        and sum(message.get("role") == "user" for message in tail) == 1
        and "final_assistant_row_id" in receipt)
    return receipt


def _complete_turn_payload(session: dict, st: _TurnRun, status_note: str | None, cols: int):
    """``(payload, raw, status)`` for message.complete; retains/clears the inflight turn and
    settles the hosted-room terminal receipt."""
    result, agent = st.result, st.agent
    # Advisory {layer, code, retryable} descriptor; computed before the retain so resume
    # replay carries the same one, and before the text so the fallback copy can use it.
    _error_surface = None
    if _result_status(result) == "error":
        try:
            from agent.error_surface import build_error_surface_from_result
            _error_surface = build_error_surface_from_result(
                result, provider=str(getattr(agent, "provider", "") or ""),
                model=str(getattr(agent, "model", "") or ""))
        except Exception:
            _error_surface = None
    raw, status, last_reasoning = _turn_outcome(result, _error_surface)
    if _is_bot_mode_session(session):
        raw = _bot_mode_delivery_text(raw, successful=status == "complete")
    payload = {"text": raw, "usage": _get_usage(agent), "status": status}
    if receipt := _persisted_turn_receipt(st, raw, status):
        payload["persisted_turn"] = receipt
    if last_reasoning:
        payload["reasoning"] = last_reasoning
    if status_note:
        payload["warning"] = status_note
    # A runtime that delivers its final message as an interim (the Codex app-server bridge routes every
    # completed agentMessage there) never sets response_previewed; the client would render it twice (#125951).
    was_delivered = getattr(agent, "_interim_text_was_delivered", None)
    if result.get("response_previewed") or (callable(was_delivered) and was_delivered(raw) is True):
        payload["response_previewed"] = True
    # transform_llm_output may rewrite the final after streaming: the renderer must treat
    # this payload as the authoritative replacement even without a prefix relationship.
    if result.get("response_transformed"):
        payload["response_transformed"] = True
    # Structured billing-wall descriptor: the client renders recovery without re-parsing text.
    if _billing_block := result.get("billing_block"):
        payload["billing"] = _billing_block
        payload["failure_reason"] = result.get("failure_reason")
    if rendered := render_message(raw, cols):
        payload["rendered"] = rendered
    error_value = result.get("error")
    final_text = result.get("final_response")
    has_partial_text = bool(
        result.get("partial") and isinstance(final_text, str)
        and final_text.strip() and final_text.strip() != str(error_value or "").strip())
    with session["history_lock"]:
        if status == "error":
            # Retain the failed turn: resume's inflight payload is the only carrier of the
            # failure if this frame is lost to a disconnect.
            if has_partial_text and not (session.get("inflight_turn") or {}).get("assistant"):
                # Non-streaming results need a replay body too; keep existing streamed segments intact.
                _append_inflight_delta(session, raw)
            _fail_inflight_turn(session, error_value, error_surface=_error_surface)
            st.error_retained = True
            st.error_detail = _turn_failure_detail(
                error_value, result.get("failure_reason"), st.prompt_text)
        else:
            _clear_inflight_turn(session)
    if status == "error":
        payload["error"] = str(error_value or raw)
        payload["recoverable"] = True
        # Desktop distinguishes retained answer text from error copy using this flag.
        if has_partial_text:
            payload["partial"] = True
        if _error_surface:
            payload["error_surface"] = _error_surface
    if st.terminal_callback is not None:
        st.receipt_attempted = True
        st.terminal_callback({
            "status": {"interrupted": "cancelled", "error": "failed"}.get(status, "settled"),
            "text": raw if isinstance(raw, str) else str(raw),
            **({"error": str(error_value or raw)} if status == "error" else {})})
        st.receipt_committed = True
    if st.receipt_committed:
        _retire_turn_marker(session, st.marker_key)
    return payload, raw, status


def _recover_turn_exception(sid: str, session: dict, st: _TurnRun, e: BaseException) -> None:
    """Except-path of the turn: crash log, history restore, terminal error frame."""
    import traceback
    with contextlib.suppress(Exception):
        os.makedirs(os.path.dirname(_CRASH_LOG), exist_ok=True)
        with open(_CRASH_LOG, "a", encoding="utf-8") as f:
            f.write(
                f"\n=== turn-dispatcher exception · "
                f"{time.strftime('%Y-%m-%d %H:%M:%S')} · sid={sid} ===\n")
            f.write(traceback.format_exc())
    print(f"[gateway-turn] {type(e).__name__}: {e}", file=sys.stderr, flush=True)
    # A finalizer exception can leave in-memory history at the turn-start snapshot.
    _restore_agent_history_after_turn_error(session, st.agent)
    if st.terminal_callback is not None and not st.receipt_attempted:
        st.receipt_attempted = True
        try:
            st.terminal_callback({"status": "failed", "text": "", "error": str(e)})
            st.receipt_committed = True
        except Exception:
            logger.exception("hosted room terminal receipt commit failed")
    try:
        # Same terminal error frame shape as the returned-error path.
        _emit_terminal_turn_error(sid, session, e, retire_marker=st.receipt_committed)
        st.error_retained = True
        st.error_detail = _turn_failure_detail(e, type(e).__name__, st.prompt_text)
    except Exception as emit_exc:
        print(
            f"[gateway-turn] terminal error emit failed: {type(emit_exc).__name__}: {emit_exc}",
            file=sys.stderr, flush=True)
        _emit("error", sid, {"message": str(e)})


def _finish_turn(sid: str, session: dict, st: _TurnRun) -> None:
    """Finally-path of the turn: release everything, then the "tui turn finished" bookend."""
    # Drop both pre-turn history snapshots before asking glibc to return pages (a test
    # inspects these two locals by name).
    history, run_kwargs = st.history, st.run_kwargs
    history.clear()
    if isinstance(run_kwargs, dict):
        run_kwargs.clear()
    try:  # while the profile HERMES_HOME override is still active (session's own config)
        from hermes_cli.mem_trim import trim_memory
        # The finishing session is still marked running here; every OTHER session must be idle (#58576).
        if _sessions_quiescent(exclude=sid):
            trim_memory(reason="tui turn completion")
    except Exception:
        logger.debug("post-turn memory trim failed", exc_info=True)
    if st.thinking_started:
        with contextlib.suppress(Exception):
            from tools.voice_mode import stop_thinking_sound
            stop_thinking_sound()
    if st.tts_queue is not None:
        st.tts_queue.put(None)  # end-of-text sentinel — flush + finish speaking
    if st.one_turn_restore:
        try:
            _restore_agent_model_runtime(st.agent, st.one_turn_restore)
            _restart_slash_worker(sid, session)
            _persist_live_session_runtime(session)
            _persist_live_session_system_prompt(session)
        except Exception:
            logger.debug("TUI one-turn model restore failed", exc_info=True)
    scopes = st.scopes
    with contextlib.suppress(Exception):
        if scopes.approval is not None:
            from tools.approval_context import reset_current_session_key
            reset_current_session_key(scopes.approval)
    if scopes.home is not None:
        reset_hermes_home_override(scopes.home)
    if scopes.secret is not None:
        reset_secret_scope(scopes.secret)
    if scopes.terminal is not None:
        from tools.terminal_scope import reset_terminal_scope
        reset_terminal_scope(scopes.terminal)
    _clear_session_context(scopes.session_tokens)


# Bounded so a contended state.db cannot hold ``_sessions_lock``; a skipped heal is retried on the next prompt.
_ROUTING_REOPEN_PATIENCE_S = 0.5


def _routing_provenance_db(session: dict):
    """The session's own SessionDB for :func:`_reopen_routed_session_row`, or a ``None`` context."""
    try:
        return _session_db(session)
    except Exception:
        logger.debug("could not resolve the session db for routing provenance", exc_info=True)
        return contextlib.nullcontext(None)


def _reopen_routed_session_row(db, sid: str, session: dict) -> None:
    """Routing provenance for #106459, called under ``_sessions_lock`` while this backend still has the
    session registered and is about to start a turn for it: an explicit-close stamp on its row missed a
    live conversation, so it is cleared (the gateway's #54878 stale-route self-heal, which this backend
    lacked). ``_pop_session_by_id`` claims teardown under the same lock and sets ``_closing`` before
    ``_finalize_session`` stamps, so a close of THIS session is either not written yet or already stops
    the turn -- it is never cleared here. Best-effort: a failure never blocks the turn."""
    session_id = getattr(session.get("agent"), "session_id", None)  # the compression tip this turn writes
    if db is None or not session_id:
        return
    try:
        db.reopen_if_explicitly_closed(
            session_id, provenance=f"TUI session {sid} is still registered and accepting a turn",
            patience_s=_ROUTING_REOPEN_PATIENCE_S)
    except Exception:
        logger.debug("routing-provenance reopen failed for %s", session_id, exc_info=True)


def _run_prompt_submit(
    rid, sid: str, session: dict, text: Any, *, display_kind: str | None = None,
    display_metadata: dict | None = None, image_paths: list[str] | None = None,
    queued_prompt_generation: int | None = None,
    terminal_callback: Callable[[dict[str, Any]], None] | None = None,
    turn_author: dict | None = None) -> bool:
    # Every dispatch binds the session's own row (session_key, real source) before the turn writes:
    # the synthesized turns that enter here directly (crash auto-continue, queued-prompt drain,
    # wake-ups) bypass prompt.submit's persist, and a row-less turn is otherwise materialized by
    # the token-accounting guard as an anonymous session (#111999).
    if _ensure_session_db_row(session) is False:
        logger.warning(
            "prompt dispatch: session store unavailable for %s — this turn may not persist",
            session.get("session_key") or sid)
    admitted = _admit_prompt_turn(
        sid, session, text, image_paths, queued_prompt_generation, display_kind, display_metadata)
    if admitted is None:
        return False
    images, agent = admitted
    from gateway.warning_notifications import diagnostic_turn_muted
    from agent.notification_presentation import notification_config_snapshot
    with _session_profile_runtime_scope(session):
        notification_config = notification_config_snapshot()
        muted = diagnostic_turn_muted(display_metadata, "tui", notification_config)
    if muted:
        display_kind = "hidden"
    # The ONE INFO record proving a prompt was accepted by THIS process; ties ui sid,
    # session_key and the agent's live session_id together.  No prompt content is logged.
    _turn_started_monotonic = time.monotonic()
    logger.info(
        # Desktop/TUI observability (#86647): this is the ONE INFO record proving a Desktop/TUI prompt was
        # accepted by THIS process, and it ties together every id a rotation-mute trace needs — the UI
        # session id, the gateway session_key, and the agent's live session_id (which compression rotates
        # independently of the other two). Before this line a Desktop request left no trace in agent.log at
        # all ("0 platform=desktop" — see #86647), so a muted window was structurally indistinguishable from
        # a request that never arrived.
        "tui prompt accepted: ui_session=%s session_key=%s agent_session_id=%s "
        "kind=%s chars=%s images=%d",
        sid, session.get("session_key") or "", getattr(agent, "session_id", "") or "",
        display_kind or "user", len(text) if isinstance(text, str) else "-", len(images))
    if not muted:
        _emit("message.start", sid)

    def run_body():
        # RPC-dispatcher ContextVars do not follow onto this thread: rebind the transport
        # before any tool can commission a child (delegate_task captures it as authority).
        transport_token = bind_transport(session.get("transport"))
        runtime_session_token = _current_runtime_session_record.set(session)
        st = _TurnRun(
            session["agent"], session.pop("one_turn_model_restore", None), terminal_callback,
            receipt_committed=terminal_callback is None)
        st.marker_key = _record_turn_marker(session, text, auto_continue=terminal_callback is None,
            notification_category=(display_metadata or {}).get("notification_category"))
        goal_followup = None
        try:
            prepared = _prepare_turn_input(sid, session, st, text, images)
            if prepared is None:
                if st.terminal_callback is not None and not st.receipt_attempted:
                    st.receipt_attempted = True
                    st.terminal_callback({
                        "status": "failed", "text": "", "error": "Context injection refused."})
                    st.receipt_committed = True
                return
            prompt, run_message, cols, streamer = prepared
            _invoke_agent(
                sid, session, st, prompt, run_message, streamer, images, display_kind,
                display_metadata, turn_author, text)
            status_note = _absorb_turn_result(
                sid, session, st, text, display_kind, display_metadata)
            payload, raw, status = _complete_turn_payload(session, st, status_note, cols)
            _emit("message.complete", sid, payload)
            goal_followup = _goal_followup_after_turn(sid, session, st.result, status, raw)
            if status == "complete":
                _after_complete_turn(sid, session, st, raw)
            # Goal judge + loop tick evaluation mutate persisted state AFTER message.complete: publish the
            # structured control snapshot now so the Desktop card never paints the pre-judge turn count.
            _publish_session_control_snapshot(sid, session, only_if_present=True)
        except Exception as e:
            _recover_turn_exception(sid, session, st, e)
        finally:
            _finish_turn(sid, session, st)
            _current_runtime_session_record.reset(runtime_session_token)
            reset_transport(transport_token)
            # A stale interim closure must not fire during a later turn.
            st.agent.interim_assistant_callback = None
            with session["history_lock"]:
                session["running"] = False
                session["last_active"] = time.time()
                if not st.error_retained:
                    _clear_inflight_turn(session)
                _release_hosted_room_turn_slot(session)
            # Closing bookend of "tui prompt accepted" — exactly one per accepted prompt.
            # agent.session_id is re-read because compression may have rotated it (an
            # accepted/finished pair whose id changed IS a rotation trace).
            if isinstance(st.result, dict):
                status = _result_status(st.result)
            else:
                status = "error" if st.error_retained else "complete"
            logger.info(
                "tui turn finished: ui_session=%s session_key=%s agent_session_id=%s status=%s "
                "error_retained=%s duration=%.1fs%s",
                sid, session.get("session_key") or "", getattr(st.agent, "session_id", "") or "",
                status, st.error_retained, time.monotonic() - _turn_started_monotonic,
                st.error_detail)
            # Backstop for turns that never reached a terminal frame.
            if st.receipt_committed:
                _retire_turn_marker(session, st.marker_key)
                with session["history_lock"]:
                    if session.get("_active_turn_marker_key") == st.marker_key:
                        session.pop("_active_turn_marker_key", None)
                    session.pop("_hosted_room_task", None)
            session.pop("_auto_continue_scheduled", None)
            _emit_settled_session_info(sid, session, st.agent)
        return st.result, goal_followup
    def run():
        from agent.notification_presentation import notification_turn
        # _prepare_turn_input owns profile binding for the worker. The context
        # here only gates presentation; do not introduce a second runtime scope.
        from agent.notification_presentation import notification_policy_snapshot
        with notification_policy_snapshot(agent, "tui", notification_config), notification_turn(agent, muted=muted, session_id=sid):
            followup = run_body()
        if followup is not None:
            _run_post_turn_followups(rid, sid, session, *followup)
    # The handle is resolved BEFORE _sessions_lock: a profile session opens its own SessionDB through the
    # state registry, and _sessions_lock gates every create/close/prompt on this backend.
    with _routing_provenance_db(session) as routing_db, _sessions_lock:
        registered = _sessions.get(sid)
        can_start = not session.get("_closing") and (registered is None or registered is session)
        if can_start:
            # Only a registered session is proof the conversation is routed here; an unregistered one may
            # still run its turn, but its stamp stays (#106459).
            if registered is session:
                _reopen_routed_session_row(routing_db, sid, session)
            can_start = _start_session_work(run, name=f"prompt-turn-{sid}", session=session) is not None
    if not can_start:
        with session["history_lock"]:
            session["running"] = False
    return can_start


def register(server) -> None:
    """Publish this module's helpers onto ``server``, rebound to its globals."""
    bind_module(globals(), server, skip=("_",))
