Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions contributors/emails/jakobbjelver@gmail.com
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
jakobbjelver
231 changes: 231 additions & 0 deletions tests/tui_gateway/test_compute_host_borrowed_lease.py
Original file line number Diff line number Diff line change
@@ -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()
6 changes: 6 additions & 0 deletions tui_gateway/compute_host.py
Original file line number Diff line number Diff line change
Expand Up @@ -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"]:
Expand Down
35 changes: 31 additions & 4 deletions tui_gateway/compute_host_bridge.py
Original file line number Diff line number Diff line change
Expand Up @@ -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:
Expand All @@ -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)),
Expand Down Expand Up @@ -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)
Expand Down
Loading
Loading