fix(tui): publish session.control.update on every control mutation path; drop the Desktop-only heartbeat driver

- /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.
This commit is contained in:
Teknium
2026-09-06 03:49:14 -07:00
parent 45b7dabfb9
commit b0f87d7a7a
10 changed files with 125 additions and 415 deletions
@@ -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
@@ -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})]
+87 -123
View File
@@ -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
-86
View File
@@ -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=("_",))
+2 -2
View File
@@ -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()
+25 -4
View File
@@ -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)})
+3 -21
View File
@@ -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):
+3
View File
@@ -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:
+2 -2
View File
@@ -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
+3 -2
View File
@@ -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: