fix(tui): extend deferred context-engine finalize to the compute-host compress routes
The salvaged deferred-finalize wiring (#65670) covered cli.py, the gateway slash command, and the three in-process tui_gateway/server.py compress sites — but the dashboard compute-host isolated session.compress / slash.compress routes run the compress mirror inside the host child, where a control-handler exception after compress_context() queued the boundary notification would leak it: never fired on success paths, or worse, fired by a LATER compress against a boundary the host had rejected. - tui_gateway/compute_host.py: _handle_control discards any pending context-engine compression notification (committed=False) when a session.compress / slash.compress control frame errors. finalize is exactly-once, so this is a no-op when the mirror already emitted or discarded it. - tui_gateway/server.py: _mirror_slash_side_effects normalizes the /compact alias onto the compress branch so isolated-session compacts actually compress and hit the same finalize/discard wiring instead of silently no-oping. - tests/tui_gateway/test_compute_host_phase1.py: compute-host route tests — commit-then-notify ordering, discard on host commit failure, /compact alias routing.
This commit is contained in:
@@ -334,3 +334,232 @@ for raw in sys.stdin:
|
||||
assert supervisor.is_running()
|
||||
finally:
|
||||
supervisor.shutdown()
|
||||
|
||||
|
||||
def _make_compress_host_session(events: list) -> dict:
|
||||
class _Agent:
|
||||
model = "host-model"
|
||||
provider = "host-provider"
|
||||
tools = []
|
||||
_cached_system_prompt = ""
|
||||
session_input_tokens = 1
|
||||
session_output_tokens = 1
|
||||
session_prompt_tokens = 1
|
||||
session_completion_tokens = 1
|
||||
session_total_tokens = 2
|
||||
session_api_calls = 1
|
||||
session_id = "rotated-id"
|
||||
|
||||
agent = _Agent()
|
||||
agent.context_compressor = type("ContextEngineStub", (), {})()
|
||||
agent.context_compressor.on_session_start = (
|
||||
lambda *_args, **_kwargs: events.append("notify")
|
||||
)
|
||||
return {
|
||||
"agent": agent,
|
||||
"session_key": "before-key",
|
||||
"history": [
|
||||
{"role": "user", "content": "before"},
|
||||
{"role": "assistant", "content": "before"},
|
||||
],
|
||||
"history_lock": threading.Lock(),
|
||||
"history_version": 2,
|
||||
"running": False,
|
||||
"manual_compression_lock": threading.Lock(),
|
||||
}
|
||||
|
||||
|
||||
def test_compute_host_compress_control_notifies_engine_after_commit(monkeypatch):
|
||||
"""The compute-host slash.compress route must fire the context-engine
|
||||
boundary hook exactly once, and only AFTER the host commits the compressed
|
||||
history + session-key sync (salvaged #65670, extended to this route)."""
|
||||
from agent.conversation_compression import (
|
||||
_queue_context_engine_compression_notification,
|
||||
finalize_context_engine_compression_notification,
|
||||
)
|
||||
from tui_gateway import server
|
||||
|
||||
out = io.StringIO()
|
||||
host = ComputeHost(stdout=out, max_workers=1, heartbeat_secs=0)
|
||||
events: list[str] = []
|
||||
session = _make_compress_host_session(events)
|
||||
|
||||
def _compress(sess, focus_topic=None, **_kwargs):
|
||||
# Simulate agent._compress_context(defer_context_engine_notification=True)
|
||||
_queue_context_engine_compression_notification(
|
||||
sess["agent"],
|
||||
new_session_id="rotated-id",
|
||||
old_session_id="before-key",
|
||||
)
|
||||
with sess["history_lock"]:
|
||||
sess["history"] = [{"role": "summary", "content": "compressed"}]
|
||||
sess["history_version"] = 3
|
||||
|
||||
def _sync(sid, sess):
|
||||
events.append("sync")
|
||||
sess["session_key"] = "after-key"
|
||||
|
||||
server._sessions["sid"] = session
|
||||
monkeypatch.setenv("HERMES_COMPUTE_HOST_CHILD", "1")
|
||||
monkeypatch.setattr(server, "_compress_session_history", _compress)
|
||||
monkeypatch.setattr(server, "_sync_session_key_after_compress", _sync)
|
||||
monkeypatch.setattr(server, "_emit", lambda *_args, **_kwargs: None)
|
||||
monkeypatch.setattr(
|
||||
server,
|
||||
"_session_info",
|
||||
lambda _agent, _session=None: {"model": "host-model", "usage": {"total": 2}},
|
||||
)
|
||||
|
||||
try:
|
||||
host.handle_frame(
|
||||
{
|
||||
"type": "control",
|
||||
"sid": "sid",
|
||||
"request_id": "compress-1",
|
||||
"route_name": "slash.compress",
|
||||
"command": "/compress",
|
||||
}
|
||||
)
|
||||
ack = _wait_for_frame(
|
||||
out,
|
||||
lambda f: f.get("type") == "control.ack" and f.get("request_id") == "compress-1",
|
||||
)
|
||||
finally:
|
||||
server._sessions.pop("sid", None)
|
||||
host.close()
|
||||
|
||||
# Exactly one notification, after the session-key commit.
|
||||
assert events == ["sync", "notify"]
|
||||
assert ack["session_key"] == "after-key"
|
||||
# Nothing pending leaks onto the agent for a later compress to misfire.
|
||||
assert (
|
||||
finalize_context_engine_compression_notification(
|
||||
session["agent"], committed=True
|
||||
)
|
||||
is False
|
||||
)
|
||||
|
||||
|
||||
def test_compute_host_compress_control_failure_discards_notification(monkeypatch):
|
||||
"""When the host-side compress mirror fails after compression queued the
|
||||
boundary notification, the pending hook must be discarded — never left to
|
||||
fire against a boundary the host rejected."""
|
||||
from agent.conversation_compression import (
|
||||
_queue_context_engine_compression_notification,
|
||||
finalize_context_engine_compression_notification,
|
||||
)
|
||||
from tui_gateway import server
|
||||
|
||||
out = io.StringIO()
|
||||
host = ComputeHost(stdout=out, max_workers=1, heartbeat_secs=0)
|
||||
events: list[str] = []
|
||||
session = _make_compress_host_session(events)
|
||||
|
||||
def _compress(sess, focus_topic=None, **_kwargs):
|
||||
_queue_context_engine_compression_notification(
|
||||
sess["agent"],
|
||||
new_session_id="rotated-id",
|
||||
old_session_id="before-key",
|
||||
)
|
||||
|
||||
def _boom(*_args, **_kwargs):
|
||||
raise RuntimeError("synthetic host commit failure")
|
||||
|
||||
server._sessions["sid"] = session
|
||||
monkeypatch.setenv("HERMES_COMPUTE_HOST_CHILD", "1")
|
||||
monkeypatch.setattr(server, "_compress_session_history", _compress)
|
||||
monkeypatch.setattr(server, "_sync_session_key_after_compress", _boom)
|
||||
monkeypatch.setattr(server, "_emit", lambda *_args, **_kwargs: None)
|
||||
monkeypatch.setattr(
|
||||
server,
|
||||
"_session_info",
|
||||
lambda _agent, _session=None: {"model": "host-model", "usage": {"total": 2}},
|
||||
)
|
||||
|
||||
try:
|
||||
host.handle_frame(
|
||||
{
|
||||
"type": "control",
|
||||
"sid": "sid",
|
||||
"request_id": "compress-2",
|
||||
"route_name": "slash.compress",
|
||||
"command": "/compress",
|
||||
}
|
||||
)
|
||||
ack = _wait_for_frame(
|
||||
out,
|
||||
lambda f: f.get("type") == "control.ack" and f.get("request_id") == "compress-2",
|
||||
)
|
||||
finally:
|
||||
server._sessions.pop("sid", None)
|
||||
host.close()
|
||||
|
||||
assert events == []
|
||||
assert "live session sync failed" in str(ack.get("output") or "")
|
||||
# The pending notification was discarded, not left on the agent.
|
||||
assert (
|
||||
finalize_context_engine_compression_notification(
|
||||
session["agent"], committed=True
|
||||
)
|
||||
is False
|
||||
)
|
||||
assert events == []
|
||||
|
||||
|
||||
def test_compute_host_compact_alias_routes_to_compress_mirror(monkeypatch):
|
||||
"""slash.compress control frames forward the user's raw alias verbatim;
|
||||
/compact must reach the compress mirror (and its deferred-notification
|
||||
finalize wiring), not silently no-op."""
|
||||
from agent.conversation_compression import (
|
||||
_queue_context_engine_compression_notification,
|
||||
)
|
||||
from tui_gateway import server
|
||||
|
||||
out = io.StringIO()
|
||||
host = ComputeHost(stdout=out, max_workers=1, heartbeat_secs=0)
|
||||
events: list[str] = []
|
||||
session = _make_compress_host_session(events)
|
||||
calls: dict[str, object] = {}
|
||||
|
||||
def _compress(sess, focus_topic=None, **_kwargs):
|
||||
calls["focus"] = focus_topic
|
||||
_queue_context_engine_compression_notification(
|
||||
sess["agent"],
|
||||
new_session_id="rotated-id",
|
||||
old_session_id="before-key",
|
||||
)
|
||||
|
||||
server._sessions["sid"] = session
|
||||
monkeypatch.setenv("HERMES_COMPUTE_HOST_CHILD", "1")
|
||||
monkeypatch.setattr(server, "_compress_session_history", _compress)
|
||||
monkeypatch.setattr(
|
||||
server, "_sync_session_key_after_compress", lambda *_a: events.append("sync")
|
||||
)
|
||||
monkeypatch.setattr(server, "_emit", lambda *_args, **_kwargs: None)
|
||||
monkeypatch.setattr(
|
||||
server,
|
||||
"_session_info",
|
||||
lambda _agent, _session=None: {"model": "host-model", "usage": {"total": 2}},
|
||||
)
|
||||
|
||||
try:
|
||||
host.handle_frame(
|
||||
{
|
||||
"type": "control",
|
||||
"sid": "sid",
|
||||
"request_id": "compact-1",
|
||||
"route_name": "slash.compress",
|
||||
"command": "/compact focus topic",
|
||||
}
|
||||
)
|
||||
ack = _wait_for_frame(
|
||||
out,
|
||||
lambda f: f.get("type") == "control.ack" and f.get("request_id") == "compact-1",
|
||||
)
|
||||
finally:
|
||||
server._sessions.pop("sid", None)
|
||||
host.close()
|
||||
|
||||
assert calls == {"focus": "focus topic"}
|
||||
assert events == ["sync", "notify"]
|
||||
assert ack["route_name"] == "slash.compress"
|
||||
|
||||
@@ -567,6 +567,28 @@ class ComputeHost:
|
||||
}
|
||||
)
|
||||
except Exception as exc:
|
||||
if route_name in {"session.compress", "slash.compress"}:
|
||||
# The compress mirror defers the context-engine boundary
|
||||
# notification until the host commits. If anything raises
|
||||
# between queueing and finalize (e.g. building the ack's
|
||||
# session_info), discard the pending notification so it can't
|
||||
# leak onto the agent and fire against a rejected boundary on
|
||||
# a later compress. finalize is exactly-once, so this is a
|
||||
# no-op when the mirror already emitted or discarded it.
|
||||
try:
|
||||
from tui_gateway import server as _server
|
||||
from agent.conversation_compression import (
|
||||
finalize_context_engine_compression_notification,
|
||||
)
|
||||
|
||||
_agent = (_server._sessions.get(sid) or {}).get("agent")
|
||||
if _agent is not None:
|
||||
finalize_context_engine_compression_notification(
|
||||
_agent,
|
||||
committed=False,
|
||||
)
|
||||
except Exception:
|
||||
pass
|
||||
self.emit({"type": "control.error", "sid": sid, "request_id": request_id, "message": str(exc)})
|
||||
|
||||
def _bump_progress(self) -> None:
|
||||
|
||||
@@ -15052,6 +15052,13 @@ def _mirror_slash_side_effects(sid: str, session: dict, command: str) -> str:
|
||||
(parts[1].strip() if len(parts) > 1 else ""),
|
||||
session.get("agent"),
|
||||
)
|
||||
if name == "compact":
|
||||
# /compact is an alias of /compress in every host. The compute-host
|
||||
# slash.compress control forwards the user's raw alias verbatim, so
|
||||
# without normalizing here the child mirror silently no-ops — the
|
||||
# session never compresses and the deferred context-engine
|
||||
# notification wiring below is never exercised for that route.
|
||||
name = "compress"
|
||||
|
||||
# Reject agent-mutating commands during an in-flight turn. These
|
||||
# all do read-then-mutate on live agent/session state that the
|
||||
|
||||
Reference in New Issue
Block a user