From b0f87d7a7ad63d87e9777d5849c9aa5f676f2141 Mon Sep 17 00:00:00 2001 From: Teknium <127238744+teknium1@users.noreply.github.com> Date: Sun, 6 Sep 2026 03:49:14 -0700 Subject: [PATCH] fix(tui): publish session.control.update on every control mutation path; drop the Desktop-only heartbeat driver MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - /goal and /loop are command.dispatch built-ins (never reach the slash worker), so the worker-branch publish never fired for the most common typed commands and the structured card stayed stale until the next turn. The publish helper now lives beside _snapshot_control in methods_session_control.py and is called from both slash paths plus the end of the turn tail — after the goal judge / loop tick evaluation, which mutate state AFTER message.complete. - session.control skips the duplicate emit for dispatcher-backed actions (still exactly one update per mutation). - Removed tui_gateway/desktop_heartbeat_driver.py, the entry/ws hooks and the hermes_cli lock rework: #104224 drives /heartbeat from the existing per-session poller for every TUI client. - test_session_control.py: dropped presentation-shaped cases, added the publication invariants (both slash paths, unknown dispatch publishes nothing, real _run_prompt_submit turn publishes post-judge). 3/3 new tests fail on the contributor backend. --- .../test_desktop_heartbeat_driver.py | 118 ---------- .../test_desktop_heartbeat_slash_sync.py | 57 ----- tests/tui_gateway/test_session_control.py | 210 ++++++++---------- tui_gateway/desktop_heartbeat_driver.py | 86 ------- tui_gateway/entry.py | 4 +- tui_gateway/methods_session_control.py | 29 ++- tui_gateway/methods_tools.py | 24 +- tui_gateway/prompt_turn.py | 3 + tui_gateway/server.py | 4 +- tui_gateway/ws.py | 5 +- 10 files changed, 125 insertions(+), 415 deletions(-) delete mode 100644 tests/tui_gateway/test_desktop_heartbeat_driver.py delete mode 100644 tests/tui_gateway/test_desktop_heartbeat_slash_sync.py delete mode 100644 tui_gateway/desktop_heartbeat_driver.py diff --git a/tests/tui_gateway/test_desktop_heartbeat_driver.py b/tests/tui_gateway/test_desktop_heartbeat_driver.py deleted file mode 100644 index b056934955..0000000000 --- a/tests/tui_gateway/test_desktop_heartbeat_driver.py +++ /dev/null @@ -1,118 +0,0 @@ -"""Regression coverage for Desktop-owned Heartbeat scheduling.""" - -from __future__ import annotations - -import io -import threading -import time - - -def test_due_heartbeat_dispatches_a_desktop_turn_and_publishes_the_new_count(tmp_path, monkeypatch): - """A due idle Desktop session must fire through the normal turn path exactly once.""" - home = tmp_path / ".hermes" - home.mkdir() - monkeypatch.setattr("pathlib.Path.home", lambda: tmp_path) - monkeypatch.setenv("HERMES_HOME", str(home)) - - from hermes_cli import goals - from hermes_cli.heartbeat import HeartbeatManager, HeartbeatState, save_heartbeat - from tui_gateway import server - - goals._DB_CACHE.clear() - sid, key = "desktop-heartbeat-sid", "desktop-heartbeat-key" - session = { - "session_key": key, - "history_lock": threading.RLock(), - "running": False, - "queued_prompt": None, - "queued_prompts": [], - "_closing": False, - "agent": object(), - "attached_images": [], - } - server._sessions[sid] = session - save_heartbeat( - key, - HeartbeatState( - prompt="Check the desktop Heartbeat path.", - interval_seconds=60, - created_at=time.time() - 61, - ), - ) - dispatched, events = [], [] - monkeypatch.setattr( - server, - "_run_prompt_submit", - lambda rid, got_sid, got_session, text, **kwargs: dispatched.append((rid, got_sid, got_session, text)), - ) - monkeypatch.setattr(server, "_emit", lambda event, got_sid, payload: events.append((event, got_sid, payload))) - - try: - assert server._poll_desktop_heartbeats_once() == 1 - assert len(dispatched) == 1 - assert dispatched[0][1:3] == (sid, session) - assert "[Heartbeat — recurring instruction" in dispatched[0][3] - assert session["running"] is True - assert HeartbeatManager(key).state.fire_count == 1 - assert any( - event == "session.control.update" - and got_sid == sid - and payload["control"]["heartbeat"]["fire_count"] == 1 - for event, got_sid, payload in events - ) - finally: - server._sessions.pop(sid, None) - goals._DB_CACHE.clear() - - -def test_entry_starts_desktop_heartbeat_driver(monkeypatch): - from hermes_cli import model_switch_providers - from tui_gateway import entry, server - - started = {"count": 0} - monkeypatch.setattr(server, "_start_desktop_heartbeat_driver", lambda: started.__setitem__("count", started["count"] + 1)) - monkeypatch.setattr(server, "_start_backend_heartbeat_refresher", lambda: None) - monkeypatch.setattr(server, "_schedule_startup_orphan_sweep", lambda: None) - monkeypatch.setattr(entry, "_install_sidecar_publisher", lambda: None) - monkeypatch.setattr(entry, "ensure_mcp_discovery_started", lambda: None) - monkeypatch.setattr(entry, "resolve_skin", lambda: "default") - monkeypatch.setattr(entry.server, "_ensure_skin_watcher", lambda: None) - monkeypatch.setattr(entry, "_log_exit", lambda reason: None) - monkeypatch.setattr(entry, "handle_spurious_eof", lambda *args: False) - monkeypatch.setattr(entry, "write_json", lambda payload: True) - monkeypatch.setattr(entry.sys, "stdin", io.StringIO("")) - monkeypatch.setattr(model_switch_providers, "prewarm_picker_cache_async", lambda: None) - - entry.main() - assert started["count"] == 1 - - -def test_websocket_starts_desktop_heartbeat_driver(monkeypatch): - import asyncio - - from tui_gateway import server, ws - - started = {"count": 0} - monkeypatch.setattr(server, "_start_desktop_heartbeat_driver", lambda: started.__setitem__("count", started["count"] + 1)) - monkeypatch.setattr(server, "_start_backend_heartbeat_refresher", lambda: None) - monkeypatch.setattr(server, "_schedule_startup_orphan_sweep", lambda: None) - monkeypatch.setattr(server, "resolve_skin", lambda: "default") - monkeypatch.setattr(server, "_ensure_skin_watcher", lambda: None) - monkeypatch.setattr(server, "register_live_transport", lambda *args, **kwargs: None) - monkeypatch.setattr(server, "_WS_ORPHAN_REAP_GRACE_S", 0) - - class FakeWS: - async def accept(self): - pass - - async def send_text(self, line): - pass - - async def receive_text(self): - raise ws._WebSocketDisconnect() - - async def close(self): - pass - - asyncio.run(ws.handle_ws(FakeWS())) - assert started["count"] == 1 diff --git a/tests/tui_gateway/test_desktop_heartbeat_slash_sync.py b/tests/tui_gateway/test_desktop_heartbeat_slash_sync.py deleted file mode 100644 index 8acef7f4bd..0000000000 --- a/tests/tui_gateway/test_desktop_heartbeat_slash_sync.py +++ /dev/null @@ -1,57 +0,0 @@ -"""Desktop must receive Heartbeat state immediately after slash setup.""" - -from __future__ import annotations - - -class _Worker: - def __init__(self, output: str) -> None: - self.output = output - self.commands: list[str] = [] - - def run(self, command: str) -> str: - self.commands.append(command) - return self.output - - -def test_worker_backed_heartbeat_slash_publishes_control_snapshot(monkeypatch): - """A slash-created Heartbeat appears without leaving and re-entering its chat.""" - from tui_gateway import server - - sid = "heartbeat-session" - session_key = "persistent-heartbeat-session" - worker = _Worker("♥ Heartbeat set (every 1m): Reply HEARTBEAT_OK") - snapshot = { - "goal": None, - "heartbeat": { - "created_at": 1_700_000_000, - "fire_count": 0, - "interval_seconds": 60, - "last_fired_at": 1_700_000_000, - "prompt": "Reply HEARTBEAT_OK", - "status": "active", - }, - "loop": None, - "revision": "heartbeat-created", - "updated_at": 1_700_000_000, - } - emitted: list[tuple[str, str, dict]] = [] - session = { - "agent": None, - "running": False, - "session_key": session_key, - "slash_worker": worker, - } - server._sessions[sid] = session - monkeypatch.setattr(server, "_snapshot_control", lambda key: snapshot if key == session_key else None) - monkeypatch.setattr(server, "_emit", lambda event, event_sid, payload=None: emitted.append((event, event_sid, payload))) - - try: - response = server._methods["slash.exec"]( - "request-1", {"command": "/heartbeat every 1m Reply HEARTBEAT_OK", "session_id": sid} - ) - finally: - server._sessions.pop(sid, None) - - assert response["result"]["output"] == worker.output - assert worker.commands == ["/heartbeat every 1m Reply HEARTBEAT_OK"] - assert emitted == [("session.control.update", sid, {"control": snapshot})] diff --git a/tests/tui_gateway/test_session_control.py b/tests/tui_gateway/test_session_control.py index 2b05cd28e5..b7872f8816 100644 --- a/tests/tui_gateway/test_session_control.py +++ b/tests/tui_gateway/test_session_control.py @@ -210,22 +210,6 @@ class TestStructuredRead: ): assert forbidden not in serialized - def test_wait_barrier_is_absolute_and_unchanged_reads_do_not_count_down(self, server, session, monkeypatch): - sid, key, _ = session - _save_goal(key, waiting_until=2_000.0, waiting_reason="rate limit") - clock = SimpleNamespace(now=1_000.0) - monkeypatch.setattr(server, "time", SimpleNamespace(time=lambda: clock.now)) - - first = _control(server, sid) - clock.now = 1_001.0 - second = _control(server, sid) - - assert first == second - assert first["goal"]["wait_barrier"] == { - "type": "until", "until_at": 2_000.0, "reason": "rate limit", - } - assert "remaining_seconds" not in first["goal"]["wait_barrier"] - @pytest.mark.parametrize( ("kind", "setup"), [ @@ -251,45 +235,6 @@ class TestStructuredRead: assert _control(server, sid)["revision"] != before["revision"] - def test_session_and_pid_wait_barriers_keep_absolute_targets(self, server, session): - sid, key, _ = session - _save_goal(key, waiting_on_session="bg-session", waiting_reason="CI") - assert _control(server, sid)["goal"]["wait_barrier"] == { - "type": "session", "target": "bg-session", "reason": "CI", - } - _save_goal(key, waiting_on_pid=4242, waiting_reason="build") - assert _control(server, sid)["goal"]["wait_barrier"] == { - "type": "pid", "target": 4242, "reason": "build", - } - - -@pytest.mark.parametrize( - ("action", "name", "arg", "status"), - [ - ("goal.pause", "goal", "pause", "active"), - ("goal.resume", "goal", "resume", "paused"), - ("goal.clear", "goal", "clear", "active"), - ("loop.pause", "loop", "pause", "active"), - ("loop.resume", "loop", "resume", "paused"), - ("loop.stop", "loop", "stop", "active"), - ], -) -def test_dispatched_actions_use_exact_mapping_and_caller_rpc_id( - server, session, monkeypatch, action, name, arg, status, -): - sid, key, _ = session - if name == "goal": - _save_goal(key, status=status) - else: - _save_loop(key, status=status) - calls = _observe_dispatch(server, monkeypatch) - - response = _call(server, "session.control", rid=743, session_id=sid, action=action) - - assert "result" in response - assert calls == [(743, {"session_id": sid, "name": name, "arg": arg})] - - class TestDispatcherBackedMutations: def test_goal_pause_resume_clear_mutate_real_persisted_state(self, server, session): sid, key, _ = session @@ -435,77 +380,96 @@ class TestErrorsAndEvents: assert _error(response)["code"] == 4004 assert emitted == [] - def test_success_response_snapshot_is_exactly_the_one_update_event(self, server, session, monkeypatch): - sid, key, _ = session - _save_goal(key) + +class TestUpdatePublication: + """The Desktop card repaints from ``session.control.update``; every path that mutates control state must + publish it exactly once, and the read after a turn must reflect the post-turn hooks.""" + + def _capture(self, server, monkeypatch): emitted = [] monkeypatch.setattr(server, "_emit", lambda event, event_sid, payload=None: emitted.append((event, event_sid, payload))) + return emitted - response = _call(server, "session.control", session_id=sid, action="goal.pause") - assert emitted == [("session.control.update", sid, {"control": response["result"]["control"]})] - - def test_dispatch_envelope_preserves_all_user_visible_fields(self, server, session, monkeypatch): - sid, key, _ = session - _save_goal(key, status="paused") - expected = { - "type": "send", - "output": "already-rendered output", - "notice": "Goal resumed", - "message": "Continue toward the goal.", - "display": "/goal resume", - } - dispatch_result = { - **expected, - "model_context": {"private_trace": "never expose this nested value"}, - "private_reasoning": "never expose this model-facing field", - } - emitted = [] - monkeypatch.setattr(server, "_emit", lambda event, event_sid, payload=None: emitted.append((event, event_sid, payload))) - monkeypatch.setitem( - server._methods, - "command.dispatch", - lambda rid, params: {"jsonrpc": "2.0", "id": rid, "result": dispatch_result}, - ) - - response = _call(server, "session.control", session_id=sid, action="goal.resume") - assert response["result"]["dispatch"] == expected - assert "model_context" not in response["result"]["dispatch"] - assert "private_reasoning" not in response["result"]["dispatch"] - assert emitted == [("session.control.update", sid, {"control": response["result"]["control"]})] - assert "model_context" not in emitted[0][2] - assert "private_reasoning" not in emitted[0][2] - - def test_emit_failure_must_not_turn_mutation_into_rpc_error(self, server, session, monkeypatch): - """Event delivery is best-effort: a failing _emit must not swallow an - already-applied mutation. The persisted goal must be paused, the RPC - response must carry the correct snapshot and dispatch envelope, and the - emission error must be debug-logged rather than propagated.""" - from hermes_cli.goals import GoalManager - - sid, key, _ = session - _save_goal(key, status="active") - monkeypatch.setattr( - server, "_emit", lambda *a, **kw: (_ for _ in ()).throw(RuntimeError("emit broken")), - ) - - response = _call(server, "session.control", session_id=sid, action="goal.pause") - - assert "result" in response, f"_emit failure leaked as RPC error: {response}" - assert response["result"]["control"]["goal"]["status"] == "paused" - assert response["result"]["dispatch"]["type"] == "exec" - assert response["result"]["dispatch"]["output"] - assert GoalManager(key).state.status == "paused" - - def test_snapshot_failure_after_mutation_does_not_emit_an_all_null_fallback(self, server, session, monkeypatch): - from hermes_cli.goals import GoalManager - + def test_control_action_publishes_exactly_one_update_matching_the_response(self, server, session, monkeypatch): sid, key, _ = session _save_goal(key) - emitted = [] - monkeypatch.setattr(server, "_emit", lambda *event: emitted.append(event)) - monkeypatch.setattr(server, "_snapshot_control", lambda _key: (_ for _ in ()).throw(RuntimeError("storage failed"))) - + emitted = self._capture(server, monkeypatch) response = _call(server, "session.control", session_id=sid, action="goal.pause") - assert _error(response)["code"] == 5031 - assert GoalManager(key).state.status == "paused" - assert emitted == [] + updates = [e for e in emitted if e[0] == "session.control.update"] + assert updates == [("session.control.update", sid, {"control": response["result"]["control"]})] + + @pytest.mark.parametrize("name,arg", [("goal", "pause"), ("loop", "pause")]) + def test_typed_slash_via_command_dispatch_publishes_the_new_state(self, server, session, monkeypatch, name, arg): + """/goal and /loop never reach the slash worker (``_PENDING_INPUT_COMMANDS``) — the built-in path must + publish too, or the structured card stays stale until the next turn.""" + sid, key, _ = session + if name == "goal": + _save_goal(key) + else: + from hermes_cli.loops import LoopManager + + LoopManager(key).set("poll CI", interval_seconds=300) + emitted = self._capture(server, monkeypatch) + response = _call(server, "command.dispatch", session_id=sid, name=name, arg=arg) + assert "result" in response + updates = [e for e in emitted if e[0] == "session.control.update"] + assert len(updates) == 1 and updates[0][1] == sid + assert updates[0][2]["control"][name]["status"] == "paused" + + def test_unknown_dispatch_publishes_nothing(self, server, session, monkeypatch): + sid, _, _ = session + emitted = self._capture(server, monkeypatch) + assert "error" in _call(server, "command.dispatch", session_id=sid, name="definitely-not-a-command", arg="") + assert [e for e in emitted if e[0] == "session.control.update"] == [] + + def test_real_turn_publishes_after_the_goal_judge_ran(self, server, session, monkeypatch, tmp_path): + """``message.complete`` precedes the post-turn goal judge, so a client refresh keyed on that event reads + the pre-judge turn count. A real ``_run_prompt_submit`` turn (inline thread, stub agent) must publish the + snapshot AFTER the judge with the incremented count.""" + sid, key, entry = session + _save_goal(key, turns_used=3) + emitted = self._capture(server, monkeypatch) + + class _InlineThread: + def __init__(self, target=None, daemon=None, args=(), kwargs=None): + self._t, self._a, self._k = target, args, kwargs or {} + + def start(self): + self._t(*self._a, **self._k) + + def is_alive(self): + return False + + def join(self, timeout=None): + return None + + for name, value in { + "_wire_callbacks": lambda sid_: None, "_sync_agent_model_with_config": lambda sid_, s: None, + "_session_cwd": lambda s: str(tmp_path), "_register_session_cwd": lambda s: None, + "_tts_stream_begin": lambda: None, "_sync_session_key_after_compress": lambda *a, **k: None, + "_get_usage": lambda agent: {}, "_hermes_home": tmp_path, + }.items(): + monkeypatch.setattr(server, name, value) + monkeypatch.setattr(server.threading, "Thread", _InlineThread) + # Deterministic judge: the model-graded verdict is not under test, the ordering is. + from hermes_cli.goals import GoalManager + + def judge(self, raw, **kwargs): + from hermes_cli.goals import save_goal + + self.state.turns_used += 1 + save_goal(self.session_id, self.state) + return {"should_continue": False, "message": ""} + + monkeypatch.setattr(GoalManager, "evaluate_after_turn", judge) + entry.update({"image_counter": 0, "slash_worker": None, "show_reasoning": False, "tool_progress_mode": "all", + "inflight_turn": None, "running": True, + "agent": SimpleNamespace(session_id=key, clear_interrupt=lambda: None, + run_conversation=lambda message, **kw: {"final_response": "done"})}) + + assert server._run_prompt_submit("rid", sid, entry, "work on it") + + order = [e[0] for e in emitted] + assert "message.complete" in order and "session.control.update" in order + assert order.index("message.complete") < order.index("session.control.update") + assert emitted[order.index("session.control.update")][2]["control"]["goal"]["turns_used"] == 4 diff --git a/tui_gateway/desktop_heartbeat_driver.py b/tui_gateway/desktop_heartbeat_driver.py deleted file mode 100644 index eedfb1d6f4..0000000000 --- a/tui_gateway/desktop_heartbeat_driver.py +++ /dev/null @@ -1,86 +0,0 @@ -"""Desktop-owned Heartbeat scheduling through the normal TUI turn path.""" - -from __future__ import annotations - -import threading -import time -import uuid - -from .method_ctx import bind_module - - -_HEARTBEAT_POLL_SECONDS = 5.0 -_desktop_heartbeat_driver_started = False -_desktop_heartbeat_driver_lock = threading.Lock() - - -def _poll_desktop_heartbeats_once() -> int: - """Start one due Heartbeat turn per idle live Desktop session. - - The session lock covers the persisted fire claim and the in-memory ``running`` - claim. That makes a user prompt win whenever it arrived first, and prevents two - poll passes from incrementing the same Heartbeat more than once. - """ - fired = 0 - for sid, session in list(_sessions.items()): - if not isinstance(session, dict): - continue - session_key = str(session.get("session_key") or "") - if not session_key: - continue - lock = session.get("history_lock") - if lock is None: - continue - try: - with lock: - if ( - session.get("_closing") - or session.get("running") - or session.get("queued_prompt") is not None - or session.get("queued_prompts") - ): - continue - with _session_profile_runtime_scope(session): - from hermes_cli.heartbeat import HeartbeatManager - - prompt = HeartbeatManager(session_key).due_prompt() - if not prompt: - continue - control = _snapshot_control(session_key) - session["running"] = True - except Exception as exc: - logger.debug("desktop heartbeat check failed for %s: %s", sid, exc) - continue - - try: - _emit("session.control.update", sid, {"control": control}) - _run_prompt_submit(f"heartbeat-{uuid.uuid4().hex}", sid, session, prompt) - fired += 1 - except Exception as exc: - with lock: - session["running"] = False - logger.debug("desktop heartbeat dispatch failed for %s: %s", sid, exc) - return fired - - -def _start_desktop_heartbeat_driver() -> None: - """Start the one daemon poller that drives active Desktop session Heartbeats.""" - global _desktop_heartbeat_driver_started - with _desktop_heartbeat_driver_lock: - if _desktop_heartbeat_driver_started: - return - _desktop_heartbeat_driver_started = True - - def _loop() -> None: - while True: - try: - _poll_desktop_heartbeats_once() - except Exception: - logger.debug("desktop heartbeat driver tick failed", exc_info=True) - time.sleep(_HEARTBEAT_POLL_SECONDS) - - threading.Thread(target=_loop, name="desktop-heartbeat-driver", daemon=True).start() - - -def register(server) -> None: - bind_module(globals(), server, skip=("_",)) diff --git a/tui_gateway/entry.py b/tui_gateway/entry.py index 95d759ad6f..8697e01f0d 100644 --- a/tui_gateway/entry.py +++ b/tui_gateway/entry.py @@ -241,10 +241,10 @@ def _write_or_exit(payload: dict, reason: str) -> None: def main(): _install_sidecar_publisher() - # The process liveness row and the Desktop session heartbeat driver must start before the sweep. + # The heartbeat row lets the orphan sweep tell "live but idle" from "truly orphaned", + # so it must start BEFORE the sweep. for start, what in ( (server._start_backend_heartbeat_refresher, "backend heartbeat refresher start"), - (server._start_desktop_heartbeat_driver, "desktop heartbeat driver start"), (server._schedule_startup_orphan_sweep, "startup orphan sweep scheduling")): try: start() diff --git a/tui_gateway/methods_session_control.py b/tui_gateway/methods_session_control.py index b1a7b7e12c..51996ffc0c 100644 --- a/tui_gateway/methods_session_control.py +++ b/tui_gateway/methods_session_control.py @@ -209,6 +209,25 @@ def _goal_blocks_loop_tick(session_key: str) -> bool: return goal_blocks_loop_tick(session_key) +# Slash commands whose success changes the snapshot; ``command.dispatch`` (/goal, /loop built-ins) and the +# slash worker (/heartbeat, /subgoal) both publish after these so the Desktop card never waits for a turn. +_SESSION_CONTROL_SLASHES = frozenset({"goal", "heartbeat", "loop", "subgoal"}) + + +def _publish_session_control_snapshot(sid: str, session: dict | None) -> None: + """Best-effort ``session.control.update`` for one live session. Also called after the post-turn hooks, + because the goal judge and loop tick evaluation mutate persisted state AFTER ``message.complete`` — a + client refresh keyed on that event reads the pre-judge turn count.""" + if not session or not (session_key := str(session.get("session_key") or "")): + return + try: + with _session_profile_runtime_scope(session): + control = _snapshot_control(session_key) + _emit("session.control.update", sid, {"control": control}) + except Exception: + logger.debug("session.control.update publish failed for %s", sid, exc_info=True) + + @method("session.control.read") @_profile_scoped def _(rid, params: dict) -> dict: @@ -273,10 +292,12 @@ def _(rid, params: dict) -> dict: logger.debug("session.control snapshot after %s failed: %s", action, exc, exc_info=True) return _err(rid, 5031, f"session.control snapshot failed: {exc}") - try: - _emit("session.control.update", params.get("session_id") or "", {"control": control}) - except Exception as exc: - logger.debug("session.control.update emit failed (best-effort): %s", exc, exc_info=True) + # command.dispatch already published the update for goal/loop actions; manager actions publish here. + if action not in _ACTION_COMMAND_MAP: + try: + _emit("session.control.update", params.get("session_id") or "", {"control": control}) + except Exception as exc: + logger.debug("session.control.update emit failed (best-effort): %s", exc, exc_info=True) return _ok(rid, {"control": control, "dispatch": _dispatch_envelope(action_result)}) diff --git a/tui_gateway/methods_tools.py b/tui_gateway/methods_tools.py index be0b3f0d00..7d9582954e 100644 --- a/tui_gateway/methods_tools.py +++ b/tui_gateway/methods_tools.py @@ -810,26 +810,6 @@ _SLASH_BUILTINS = { "loop": _cmd_loop, "undo": _cmd_undo, "snapshot": _cmd_snapshot, "snap": _cmd_snapshot, "compress": _cmd_compress, "compact": _cmd_compress} -_SESSION_CONTROL_SLASHES = frozenset({"goal", "heartbeat", "loop", "subgoal"}) - - -def _publish_session_control_snapshot(sid: str, session: dict, command_name: str) -> None: - """Best-effort immediate Desktop repaint after a worker-backed automation slash command.""" - session_key = str(session.get("session_key") or "") - if not session_key: - return - try: - with _session_profile_runtime_scope(session): - control = _snapshot_control(session_key) - except Exception: - logger.debug("session.control snapshot after /%s failed", command_name, exc_info=True) - return - try: - _emit("session.control.update", sid, {"control": control}) - except Exception: - logger.debug("session.control.update emit after /%s failed", command_name, exc_info=True) - - @method("command.dispatch") def _(rid, params: dict) -> dict: name, arg = _resolve_name(params.get("name", "").lstrip("/")), params.get("arg", "") @@ -840,6 +820,8 @@ def _(rid, params: dict) -> dict: for stage in filter(None, stages): res = stage(rid, params, session, name, arg) if res is not None: + if name in _SESSION_CONTROL_SLASHES and "error" not in res: + _publish_session_control_snapshot(params.get("session_id", ""), session) return res return _err(rid, 4018, f"not a quick/plugin/bundle/skill command: {name}") @@ -896,7 +878,7 @@ def _(rid, params: dict) -> dict: if warning := _mirror_slash_side_effects(sid, session, cmd): payload["warning"] = warning if base in _SESSION_CONTROL_SLASHES: - _publish_session_control_snapshot(sid, session, base) + _publish_session_control_snapshot(sid, session) return _ok(rid, payload) except Exception as e: with contextlib.suppress(Exception): diff --git a/tui_gateway/prompt_turn.py b/tui_gateway/prompt_turn.py index d098bfc7ed..3832590448 100644 --- a/tui_gateway/prompt_turn.py +++ b/tui_gateway/prompt_turn.py @@ -800,6 +800,9 @@ def _run_prompt_submit( 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) except Exception as e: _recover_turn_exception(sid, session, st, e) finally: diff --git a/tui_gateway/server.py b/tui_gateway/server.py index 582d3cb696..316ad72ba1 100644 --- a/tui_gateway/server.py +++ b/tui_gateway/server.py @@ -3208,7 +3208,7 @@ from . import ( # noqa: E402 methods_profiles as _methods_profiles, methods_prompt as _methods_prompt, methods_session as _methods_session, methods_tools as _methods_tools, prompt_turn as _prompt_turn, billing_view as _billing_view, methods_projects as _methods_projects, methods_session_foreign as _methods_session_foreign, - methods_session_control as _methods_session_control, desktop_heartbeat_driver as _desktop_heartbeat_driver) + methods_session_control as _methods_session_control) for _m in ( _session_reaper, _session_lifecycle, _session_workdir, _compute_host_bridge, _model_switch, @@ -3218,6 +3218,6 @@ for _m in ( _methods_browser_control, _methods_session, _methods_prompt, _methods_config, _methods_config_set, _methods_complete, _methods_tools, _methods_profiles, _methods_images, _methods_bot_relay, _prompt_turn, _billing_view, _methods_projects, _methods_session_foreign, - _methods_session_control, _desktop_heartbeat_driver): + _methods_session_control): _m.register(sys.modules[__name__]) del _m diff --git a/tui_gateway/ws.py b/tui_gateway/ws.py index 8a5ec42be6..fb7e866a97 100644 --- a/tui_gateway/ws.py +++ b/tui_gateway/ws.py @@ -280,10 +280,11 @@ async def handle_ws(ws: Any, *, auth_identity: dict | None = None, subprotocol: # global broadcasts write_json can't route. server._ensure_skin_watcher() server.register_live_transport(transport) - # The process liveness row and Desktop session driver are both idempotent and once-per-process. + # Cross-backend liveness: a heartbeat row lets the startup orphan sweep tell "live but idle + # backend" from "truly orphaned". Idempotent and once-per-process, like the orphan sweep (the + # desktop app and web dashboard reach the agent via this sidecar, not entry.main()). for start, what in ( (server._start_backend_heartbeat_refresher, "backend heartbeat refresher start"), - (server._start_desktop_heartbeat_driver, "desktop heartbeat driver start"), (server._schedule_startup_orphan_sweep, "startup orphan sweep scheduling"), ): try: