diff --git a/contributors/emails/jakobbjelver@gmail.com b/contributors/emails/jakobbjelver@gmail.com new file mode 100644 index 000000000000..a72d1b297aca --- /dev/null +++ b/contributors/emails/jakobbjelver@gmail.com @@ -0,0 +1 @@ +jakobbjelver diff --git a/tests/tui_gateway/test_compute_host_borrowed_lease.py b/tests/tui_gateway/test_compute_host_borrowed_lease.py new file mode 100644 index 000000000000..7b643fa6c0e8 --- /dev/null +++ b/tests/tui_gateway/test_compute_host_borrowed_lease.py @@ -0,0 +1,231 @@ +"""Regression tests for #101416: isolated (compute-host) turns refused their own session. + +With ``dashboard.turn_isolation: true`` every lazy (agent-not-yet-built, i.e. every NEW) desktop +session's turn is routed to the compute-host CHILD process. The parent claims the session's +active-session lease in ``prompt.submit`` before routing; the child's freshly built session record +carried NO lease, so ``_admit_prompt_turn`` re-claimed from the child's pid and was fenced out by the +parent's own registry entry (``_is_same_writer`` requires the same pid AND the same live_session_id). + +The fix: the parent vouches on the turn frame (``active_session_lease`` = {lease_id, session_id}) and +the child installs an INERT borrow (``ActiveSessionLease(enabled=False)``) before the turn pipeline +runs. The REAL lease never leaves the parent: it is re-anchored there on a child-side compression +rotation and held past ``session.close`` until the child's turn settles. +""" + +from __future__ import annotations + +import io +import json +import os +import threading +import time +import types + +import pytest + +from tui_gateway import server +from tui_gateway.compute_host import ComputeHost + + +def _frames(out: io.StringIO) -> list[dict]: + return [json.loads(line) for line in out.getvalue().splitlines() if line.strip()] + + +def _wait(out: io.StringIO, predicate, timeout: float = 5.0) -> dict: + deadline = time.monotonic() + timeout + while time.monotonic() < deadline: + for frame in _frames(out): + if predicate(frame): + return frame + time.sleep(0.01) + raise AssertionError(f"timed out; saw={_frames(out)}") + + +def _stub_agent(deltas: list[str]) -> types.SimpleNamespace: + def run_conversation(prompt, *, conversation_history=None, stream_callback=None, **_kw): + final = "".join(deltas) + if stream_callback is not None: + for chunk in deltas: + stream_callback(chunk) + messages = [*(conversation_history or []), {"role": "user", "content": prompt}, + {"role": "assistant", "content": final}] + return {"final_response": final, "messages": messages} + + return types.SimpleNamespace( + session_id="s1-key", run_conversation=run_conversation, + clear_interrupt=lambda: None, hard_interrupt=lambda *a, **k: None) + + +def _make_frame(sid: str, **overrides) -> dict: + frame = {"type": "turn.start", "sid": sid, "request_id": "turn", "text": "hello", + "session_key": "s1-key", "source": "desktop", "cols": 80, "history": []} + frame.update(overrides) + return frame + + +def _seed_parent_lease(key: str, live_session_id: str = "parent-sid"): + """Claim the lease exactly as the parent dashboard would (same pid, its own live id).""" + from hermes_cli.active_sessions import try_acquire_active_session + + lease, message = try_acquire_active_session( + session_id=key, surface="desktop", config={}, + metadata={"live_session_id": live_session_id}, track_liveness=True) + assert message is None and lease is not None + return lease + + +def _foreign_acquire(key: str): + """A DISTINCT writer (same pid, another live id — the exact _is_same_writer fence).""" + from hermes_cli.active_sessions import try_acquire_active_session + + return try_acquire_active_session( + session_id=key, surface="cli", config={}, metadata={"live_session_id": "other-writer"}) + + +def _registry() -> list[dict]: + path = os.path.join(os.environ["HERMES_HOME"], "runtime", "active_sessions.json") + with open(path, "r", encoding="utf-8") as fh: + return json.load(fh).get("entries", []) + + +def _parent_session(sid: str, key: str, lease) -> dict: + return dict(agent=None, agent_ready=threading.Event(), session_key=key, history=[], history_version=0, + history_lock=threading.Lock(), running=True, transport=server._detached_ws_transport, + attached_images=[], cols=80, source="desktop", inflight_turn=None, created_at=time.time(), + last_active=time.time(), active_session_lease=lease, _compute_host_active=True, _sid=sid) + + +@pytest.fixture() +def isolated_env(monkeypatch, tmp_path): + """Real _build_server_session → _init_session → _run_prompt_submit → _admit_prompt_turn pipeline, + with the environment-heavy side paths neutralized. The turn BODY is cut right after admission + (``_prepare_turn_input`` → None), so the lease path under test runs REAL against the + conftest-sandboxed HERMES_HOME registry while nothing calls a provider.""" + agent = _stub_agent(["a ", "b "]) + monkeypatch.setattr(server, "_make_agent", lambda *a, **kw: agent) + monkeypatch.setattr(server, "_wire_callbacks", lambda sid: None) + monkeypatch.setattr(server, "_sync_agent_model_with_config", lambda sid, session: None) + monkeypatch.setattr(server, "_session_cwd", lambda session: str(tmp_path)) + monkeypatch.setattr(server, "_register_session_cwd", lambda session: None) + monkeypatch.setattr(server, "_tts_stream_begin", lambda: None) + monkeypatch.setattr(server, "_get_usage", lambda agent_: {}) + monkeypatch.setattr(server, "_hydrate_session_cwd", lambda *a, **k: None) + monkeypatch.setattr(server, "_wire_session_agent", lambda *a, **k: None) + monkeypatch.setattr(server, "_start_session_services", lambda *a, **k: None) + monkeypatch.setattr(server, "_schedule_mcp_late_refresh", lambda *a, **k: None) + import tui_gateway.prompt_turn as prompt_turn + + for mod in (server, prompt_turn): + if hasattr(mod, "_prepare_turn_input"): + monkeypatch.setattr(mod, "_prepare_turn_input", lambda *a, **k: None) + yield agent + for sid in [s for s in list(server._sessions) if s.startswith("s1")]: + server._sessions.pop(sid, None) + + +def _run_turn(frame: dict, timeout: float = 5.0) -> tuple[list[dict], dict | None]: + """Run one turn.start through the real child path; return (all frames, turn.end frame).""" + out = io.StringIO() + host = ComputeHost(stdout=out, heartbeat_secs=0) + try: + host.handle_frame(frame) + end = _wait(out, lambda f: f["type"] == "turn.end", timeout=timeout) + finally: + host.close() + return _frames(out), end + + +# ── Fix 1: the child borrows instead of re-claiming ───────────────────────── + + +def test_isolated_turn_runs_against_parent_leased_session(isolated_env): + """THE #101416 repro, fixed: parent holds the lease, child runs the turn end-to-end and the + registry still holds exactly the parent's entry.""" + parent = _seed_parent_lease("s1-key") + try: + frames, end = _run_turn(_make_frame( + "s1", active_session_lease={"lease_id": parent.lease_id, "session_id": "s1-key"})) + kinds = [f["type"] for f in frames] + assert "turn.started" in kinds and kinds[-1] == "turn.end" and end["session_key"] == "s1-key" + events = [(f["message"].get("params") or {}).get("type") for f in frames if f["type"] == "rpc"] + assert "error" not in events and "message.start" in events # admitted and ran, no refusal + entries = _registry() + assert [e["lease_id"] for e in entries] == [parent.lease_id] + assert entries[0]["metadata"]["live_session_id"] == "parent-sid" + finally: + parent.release() + + +def test_isolated_turn_without_matching_vouch_still_fails_closed(isolated_env): + """Negative control: a vouch for another stored id (a lease still keyed on the pre-rotation id) + keeps the legacy self-claim, which the parent's entry still fences — nothing about + _is_same_writer is relaxed, and a stale lease never authorizes the continuation.""" + parent = _seed_parent_lease("s1-key") + try: + frames, _ = _run_turn(_make_frame( + "s1", active_session_lease={"lease_id": parent.lease_id, "session_id": "stale-A"})) + errors = [f["message"]["params"]["payload"]["message"] for f in frames + if f["type"] == "rpc" and (f["message"].get("params") or {}).get("type") == "error"] + assert errors and "open in another Hermes window" in errors[0] + assert [e["lease_id"] for e in _registry()] == [parent.lease_id] + finally: + parent.release() + + +# ── Fix 2: compression rotation A->B stays owned by the parent ────────────── + + +def test_child_rotation_never_claims_and_parent_reanchors_its_real_lease(): + parent = _seed_parent_lease("A") + try: + # Child side: the borrow is retargeted locally; the registry is untouched (no child-pid lease). + child = {"session_key": "A", "history_lock": threading.Lock()} + server._install_borrowed_lease("sid", child, _make_frame( + "sid", session_key="A", active_session_lease={"lease_id": parent.lease_id, "session_id": "A"})) + assert server._transfer_active_session_slot("sid", child, new_session_id="B") is True + assert child["active_session_lease"].session_id == "B" + assert [(e["session_id"], e["pid"]) for e in _registry()] == [("A", os.getpid())] + # Parent side: a stale lease vouches for nothing; adopting the rotated key moves the REAL lease. + session = _parent_session("sid", "A", parent) + session["session_key"] = "B" + assert server._active_session_lease_vouch(session) is None + session["session_key"] = "A" + with session["history_lock"]: + server._compute_host_adopt_frame_meta(session, {"sid": "sid", "session_key": "B"}) + assert session["session_key"] == "B" and parent.session_id == "B" + assert [(e["session_id"], e["lease_id"]) for e in _registry()] == [("B", parent.lease_id)] + assert server._active_session_lease_vouch(session) == {"lease_id": parent.lease_id, "session_id": "B"} + finally: + parent.release() + + +# ── Fix 3: close keeps the lease until the isolated turn settles ──────────── + + +def test_close_holds_lease_until_isolated_turn_settles(monkeypatch): + interrupts: list[str] = [] + monkeypatch.setattr(server, "_get_compute_host_supervisor", + lambda *a, **k: types.SimpleNamespace(interrupt=lambda sid, **k: interrupts.append(sid))) + monkeypatch.setattr(server, "_load_dashboard_process_isolation_config", lambda *a: {"turn_isolation": True}) + monkeypatch.setattr(server, "_TURN_SETTLE_BEFORE_CLOSE_SECONDS", 0.2) + monkeypatch.setattr(server, "_emit", lambda *a, **k: None) + parent = _seed_parent_lease("A") + session = _parent_session("sid", "A", parent) + session["_compute_host_turn_id"] = "turn-1" # the child is still running this turn + session["_closing"] = True + try: + assert server._teardown_popped_session(session, end_reason="tui_close") is True + assert interrupts == ["sid"] + # Close returned, the child is live: ownership is still ours and still refuses a distinct writer. + assert [e["lease_id"] for e in _registry()] == [parent.lease_id] + assert parent.lease_id in server._own_live_lease_ids() + lease, refusal = _foreign_acquire("A") + assert lease is None and getattr(refusal, "reason", "") == "SESSION_NOT_OWNED" + # Child settlement (turn.end, or turn.error from _fail_pending_turns on child death) releases it. + server._on_compute_host_turn_done("rid", "sid", session, {"type": "turn.end", "sid": "sid", "session_key": "A"}) + assert _registry() == [] and parent.lease_id not in server._own_live_lease_ids() + lease, refusal = _foreign_acquire("A") + assert refusal is None and lease is not None + lease.release() + finally: + parent.release() diff --git a/tui_gateway/compute_host.py b/tui_gateway/compute_host.py index 5d4abdb2ca95..3bed5a8163ca 100644 --- a/tui_gateway/compute_host.py +++ b/tui_gateway/compute_host.py @@ -219,6 +219,12 @@ def _run_real_turn(self, frame: dict[str, Any]) -> None: try: from tui_gateway import server session = self._ensure_server_session(server, frame) + # #101416: the parent already holds this session's active-session lease (claimed in + # prompt.submit before routing here). Install the inert borrow BEFORE the turn runs, or + # _admit_prompt_turn re-claims from this child pid and is fenced out by the parent's own + # registry entry ("already has a live owner"). Unknown flag (parent predates the field): + # no borrow, legacy self-claim path, unchanged behaviour. + server._install_borrowed_lease(sid, session, frame) text = frame["text"] if "text" in frame else frame.get("prompt", "") inflight = frame["text"] if "text" in frame else frame.get("prompt") with session["history_lock"]: diff --git a/tui_gateway/compute_host_bridge.py b/tui_gateway/compute_host_bridge.py index 74b968706f20..3449a4863eb0 100644 --- a/tui_gateway/compute_host_bridge.py +++ b/tui_gateway/compute_host_bridge.py @@ -67,7 +67,24 @@ def _compute_host_turn_frame( "service_tier_override": session.get("create_service_tier_override"), "source": _session_source(session), "attached_images": attached_images, "auth_user_id": _session_auth_user_id(session), - "queued_prompt_generation": queued_prompt_generation} + "queued_prompt_generation": queued_prompt_generation, + # #101416: vouch that this process already holds the registry lease for this session, so + # the child adopts it as an inert token instead of re-claiming and being fenced out by + # our own entry ("Session ... already has a live owner"). No lease held = no vouch, and + # the child keeps its legacy self-claim path (fail-closed refusal on conflict). + "active_session_lease": _active_session_lease_vouch(session)} + + +def _active_session_lease_vouch(session: dict) -> dict | None: + """``{lease_id, session_id}`` of the REAL registry lease this process holds for the session's + current stored id, else None. Qualified, not a bare bool: after a compression rotation a lease + still keyed on the old id must not let a (replacement) child borrow the continuation.""" + lease = session.get("active_session_lease") + if lease is None or getattr(lease, "released", False) or not getattr(lease, "enabled", False): + return None + if str(lease.session_id) != str(session.get("session_key") or ""): + return None + return {"lease_id": str(lease.lease_id), "session_id": str(lease.session_id)} def _metadata_mirror(session: dict | None) -> dict: @@ -80,9 +97,17 @@ def _compute_host_session_info(session: dict) -> dict: def _compute_host_adopt_frame_meta(session: dict, frame: dict) -> None: - """Adopt a host frame's session_key / history_version. Caller holds history_lock.""" - if frame.get("session_key"): - session["session_key"] = str(frame.get("session_key")) + """Adopt a host frame's session_key / history_version. Caller holds history_lock. + + A rotated ``session_key`` means the child compressed A->B. The child only holds an inert + borrowed token (``_install_borrowed_lease``), so the REAL registry lease is re-anchored here, + by its owner — never claimed from the child pid (authority stays singular, #103737 review).""" + new_key = str(frame.get("session_key") or "") + if new_key and new_key != str(session.get("session_key") or ""): + if not _transfer_active_session_slot(str(frame.get("sid") or ""), session, new_session_id=new_key): + logger.warning("Compression session lease did not re-anchor: sid=%s old_session_id=%s new_session_id=%s", + frame.get("sid"), session.get("session_key"), new_key) + session["session_key"] = new_key if frame.get("history_version") is not None: with contextlib.suppress(Exception): session["history_version"] = max(int(session.get("history_version", 0)), @@ -212,6 +237,8 @@ def _on_compute_host_turn_done(rid: str, sid: str, session: dict, frame: dict) - message = str(frame.get("message") or "compute host turn failed") _emit("message.complete", sid, {"text": f"Error: {message}", "status": "error"}) _apply_compute_host_metadata_mirror(session, frame) + # Settlement of a turn whose session was closed mid-flight: the real lease was held for it. + _release_deferred_active_session_lease(session) info = _compute_host_session_info(session) if not frame.get("session_info_emitted"): _emit("session.info", sid, info) diff --git a/tui_gateway/session_lifecycle.py b/tui_gateway/session_lifecycle.py index 92e790c5dcc3..2beee22a250a 100644 --- a/tui_gateway/session_lifecycle.py +++ b/tui_gateway/session_lifecycle.py @@ -79,10 +79,44 @@ def _claim_active_session_slot( return (None, _SESSION_OWNERSHIP_UNAVAILABLE) +def _install_borrowed_lease(sid: str, session: dict, frame: dict) -> None: + """Adopt the parent's registry slot as an INERT token on a compute-host child. + + An isolated turn runs in the compute-host CHILD process. The parent dashboard already + claimed the session's active-session lease before routing the turn here, but the child's + freshly built session record never carried it — so ``_admit_prompt_turn`` re-claimed from + the child's pid and was fenced out by the parent's own entry (``_is_same_writer`` requires + same pid AND same live_session_id): every isolated turn failed with "Session ... already + has a live owner" (#101416). The slot is real and owned upstream, so the child must + neither claim a second one nor be able to release/transfer the parent's: the token is + ``enabled=False`` — ``release()`` is a no-op and ``transfer_active_session`` only retargets + the token locally, so a compression rotation A->B inside the child never reaches the + registry (the parent re-anchors its real lease from the reported ``session_key``, see + ``_compute_host_adopt_frame_meta``). NOT ``released=True``: a released token makes the + transfer fall through to a real registry claim under the child pid. + + Only installed when the frame's ``active_session_lease`` vouch names THIS stored session + id — a parent lease still keyed on a pre-rotation id must not authorize its continuation. + Anything else keeps the legacy behaviour — the child claims for itself and any ownership + conflict fails CLOSED with the visible refusal. + """ + vouch = frame.get("active_session_lease") + if not isinstance(vouch, dict) or session.get("active_session_lease") is not None: + return + key = str(session.get("session_key") or "") + if not key or str(vouch.get("session_id") or "") != key: + return + from hermes_cli.active_sessions import ActiveSessionLease + session["active_session_lease"] = ActiveSessionLease( + lease_id=f"borrowed:{vouch.get('lease_id') or sid}", session_id=key, + surface=str(frame.get("source") or "desktop"), enabled=False) + + def _ensure_active_session_slot(sid: str, session: dict) -> str | None: """Claim this session's cap slot on its first real turn; None when ok. session.create/resume deliberately do NOT claim: tile paints, reconnect-resumes and abandoned drafts would hold invisible slots (no DB row) - that starve the messaging gateway sharing the cap. Anything holding a slot must be user-visible.""" + that starve the messaging gateway sharing the cap. Anything holding a slot must be user-visible. An + inert borrowed token (see _install_borrowed_lease) also lands here: present = slot held upstream.""" if session.get("active_session_lease") is not None: return None lease, limit_message = _claim_active_session_slot( @@ -139,10 +173,12 @@ def _release_hosted_room_turn_slot(session: dict) -> None: def _own_live_lease_ids(*, exclude=None) -> set[str]: - """Snapshot leases still backed by this process's live session records.""" + """Snapshot leases still backed by this process's live session records (plus leases deferred past a + close for an unsettled isolated turn — still ours until the child settles).""" with _sessions_lock: return {str(lease.lease_id) for session in _sessions.values() - if (lease := session.get("active_session_lease")) is not None and lease is not exclude} + if (lease := session.get("active_session_lease")) is not None and lease is not exclude + } | set(_deferred_active_session_leases) @contextlib.contextmanager @@ -423,10 +459,49 @@ def _teardown_popped_session(session: dict | None, *, end_reason: str = "tui_clo "session turn thread still alive after %.1fs teardown grace", _TURN_SETTLE_BEFORE_CLOSE_SECONDS) except Exception: logger.debug("failed waiting for session turn thread", exc_info=True) + if end_reason != "tui_shutdown": + _settle_isolated_turn_before_close(session) _teardown_session(session, end_reason=end_reason) return True +# lease_id -> REAL lease of a closed session whose isolated child turn has not settled yet. Still live +# authority for the orphan sweep (``_own_live_lease_ids``); released by ``_release_deferred_active_session_lease``. +_deferred_active_session_leases: dict[str, Any] = {} + + +def _settle_isolated_turn_before_close(session: dict) -> None: + """An isolated turn runs in the compute-host child, not on ``_run_thread``: interrupt it and give it the + same close grace, and if it still has not settled keep the REAL lease out of finalize's release — the + completion callback releases it on the correlated turn.end/turn.error (child death fails pending turns + the same way). The RPC close is bounded; ownership ends with the child's last write, never with the + grace timer, else a second backend acquires the stored session while the child is still writing.""" + if not session.get("_compute_host_turn_id") or not _session_uses_compute_host(session): + return + with contextlib.suppress(Exception): + _interrupt_session_turn(_lifecycle_own_sid(session), session) + deadline = time.monotonic() + _TURN_SETTLE_BEFORE_CLOSE_SECONDS + while session.get("_compute_host_turn_id") and time.monotonic() < deadline: + time.sleep(0.05) + with session["history_lock"]: + if not session.get("_compute_host_turn_id") or (lease := session.pop("active_session_lease", None)) is None: + return + session["_deferred_active_session_lease"] = lease + _deferred_active_session_leases[str(lease.lease_id)] = lease + logger.warning("isolated turn still live after %.1fs close grace; holding lease for %s until the child settles", + _TURN_SETTLE_BEFORE_CLOSE_SECONDS, session.get("session_key")) + + +def _release_deferred_active_session_lease(session: dict) -> None: + """Settlement half of ``_settle_isolated_turn_before_close``; a no-op for sessions that never deferred.""" + lease = session.pop("_deferred_active_session_lease", None) + if lease is None: + return + _deferred_active_session_leases.pop(str(lease.lease_id), None) + if (err := _lease_retry(3, lease.release)) is not None: + logger.warning("Failed to release deferred active session slot", exc_info=err) + + def _close_session_by_id( sid: str, *, end_reason: str = "tui_close", predicate: Callable[[dict], bool] | None = None) -> bool: """Idempotent teardown funnel for callers with no resume race (resume-sensitive callers pop under