From 54d0948d38db0d2b2b4007348776c14b5ecb1d13 Mon Sep 17 00:00:00 2001 From: 2751738943 <41409874+2751738943@users.noreply.github.com> Date: Sat, 11 Jul 2026 03:18:18 +0800 Subject: [PATCH] fix(tui): route post-turn completions by owner Apply positive-proof routing to every addressed notification in the registry and TUI poller while preserving ownerless legacy behavior and TUI delivery for poll-observed completions. Remove the unused exact-key drain helper and cover ordinary success and failure, origin, compression-lineage, orphan, and poll-observed paths. Complements NousResearch/hermes-agent#54785. --- tests/test_tui_gateway_server.py | 471 ++++++++++++++++++++++----- tests/tools/test_process_registry.py | 205 ++++++++++++ tools/process_registry.py | 105 +++--- tui_gateway/server.py | 207 ++++-------- 4 files changed, 707 insertions(+), 281 deletions(-) diff --git a/tests/test_tui_gateway_server.py b/tests/test_tui_gateway_server.py index 1bd54d77f7f6..0bde25b6fd52 100644 --- a/tests/test_tui_gateway_server.py +++ b/tests/test_tui_gateway_server.py @@ -2284,34 +2284,44 @@ def test_prompt_submit_rejects_negative_truncate_ordinal(monkeypatch): server._sessions.pop("trunc-sid", None) -def test_notification_poller_drops_orphaned_events(monkeypatch): - """Completion events whose owner is gone are dropped, not hijacked.""" +@pytest.mark.parametrize( + "routing", + [ + {"session_key": "session-b"}, + {"origin_ui_session_id": "sid_gone"}, + ], +) +def test_notification_poller_drops_orphaned_events(monkeypatch, routing): + """Addressed completions whose owner is gone are dropped, not hijacked.""" import queue as _queue_mod from tools.process_registry import process_registry emitted = [] - + delivered = [] sess = _session(session_key="session-a") server._sessions["sid_a"] = sess monkeypatch.setattr(server, "_emit", lambda *a, **kw: emitted.append(a)) - monkeypatch.setattr(server, "make_stream_renderer", lambda cols: None) - monkeypatch.setattr(server, "render_message", lambda raw, cols: None) + monkeypatch.setattr( + server, + "_run_prompt_submit", + lambda _rid, _sid, _session, text: delivered.append(text), + ) + monkeypatch.setattr(server, "_get_db", lambda: None) - # Isolate the completion queue so no other test/poller can interfere. isolated_queue: _queue_mod.Queue = _queue_mod.Queue() monkeypatch.setattr(process_registry, "completion_queue", isolated_queue) process_registry._completion_consumed.discard("proc_ghost") - - # An event owned by session-b — but session-b is NOT in _sessions (gone). - isolated_queue.put({ - "type": "completion", - "session_id": "proc_ghost", - "session_key": "session-b", - "command": "echo from ghost", - "exit_code": 0, - "output": "ghost output", - }) + isolated_queue.put( + { + "type": "completion", + "session_id": "proc_ghost", + "command": "echo from ghost", + "exit_code": 0, + "output": "ghost output", + **routing, + } + ) stop = threading.Event() stop.set() @@ -2319,57 +2329,65 @@ def test_notification_poller_drops_orphaned_events(monkeypatch): try: server._notification_poller_loop(stop, "sid_a", sess) - # No status.update emitted — the orphaned event was dropped. - status_calls = [a for a in emitted if a[0] == "status.update"] - assert len(status_calls) == 0, ( - f"orphaned event should be dropped, got {len(status_calls)} status.update calls" - ) + assert [a for a in emitted if a[0] == "status.update"] == [] + assert delivered == [] finally: server._sessions.pop("sid_a", None) while not process_registry.completion_queue.empty(): process_registry.completion_queue.get_nowait() -def test_notification_poller_delivers_owned_events(monkeypatch): - """Events owned by *this* session are delivered normally (regression guard).""" +@pytest.mark.parametrize( + ("routing", "resolved_key"), + [ + ({"session_key": "session-a"}, None), + ( + { + "session_key": "stale-durable-key", + "origin_ui_session_id": "sid_a", + }, + None, + ), + ({"session_key": "old-parent-key"}, "session-a"), + ], +) +def test_notification_poller_delivers_owned_events( + monkeypatch, routing, resolved_key +): + """Direct, UI-origin, and compression-lineage owners are delivered.""" import queue as _queue_mod from tools.process_registry import process_registry - turns = [] + class _CompressionDB: + def resolve_resume_session_id(self, key): + return resolved_key if key == "old-parent-key" and resolved_key else key + + delivered = [] emitted = [] - - class _Agent: - def run_conversation(self, prompt, conversation_history=None, stream_callback=None): - turns.append(prompt) - return {"final_response": "ok", "messages": [{"role": "assistant", "content": "ok"}]} - - class _ImmediateThread: - def __init__(self, target=None, daemon=None): - self._target = target - def start(self): - self._target() - - sess = _session(session_key="session-a", agent=_Agent()) + sess = _session(session_key="session-a") server._sessions["sid_a"] = sess - monkeypatch.setattr(server.threading, "Thread", _ImmediateThread) monkeypatch.setattr(server, "_emit", lambda *a, **kw: emitted.append(a)) - monkeypatch.setattr(server, "make_stream_renderer", lambda cols: None) - monkeypatch.setattr(server, "render_message", lambda raw, cols: None) + monkeypatch.setattr( + server, + "_run_prompt_submit", + lambda _rid, _sid, _session, text: delivered.append(text), + ) + monkeypatch.setattr(server, "_get_db", lambda: _CompressionDB()) isolated_queue: _queue_mod.Queue = _queue_mod.Queue() monkeypatch.setattr(process_registry, "completion_queue", isolated_queue) process_registry._completion_consumed.discard("proc_mine") - - # An event owned by THIS session. - isolated_queue.put({ - "type": "completion", - "session_id": "proc_mine", - "session_key": "session-a", - "command": "echo mine", - "exit_code": 0, - "output": "mine", - }) + isolated_queue.put( + { + "type": "completion", + "session_id": "proc_mine", + "command": "echo mine", + "exit_code": 0, + "output": "mine", + **routing, + } + ) stop = threading.Event() stop.set() @@ -2380,58 +2398,333 @@ def test_notification_poller_delivers_owned_events(monkeypatch): status_calls = [a for a in emitted if a[0] == "status.update"] assert len(status_calls) == 1 assert status_calls[0][2]["kind"] == "process" - assert len(turns) == 1 - assert "proc_mine" in turns[0] + assert len(delivered) == 1 + assert "proc_mine" in delivered[0] finally: server._sessions.pop("sid_a", None) while not process_registry.completion_queue.empty(): process_registry.completion_queue.get_nowait() -def test_drain_owned_notifications_routes_by_session_key(monkeypatch): - """_drain_owned_notifications filters events by session ownership.""" +def _configure_immediate_prompt_run( + monkeypatch, tmp_path, *, immediate_threads=True +): + class _ImmediateThread: + def __init__(self, target=None, daemon=None, **_kwargs): + self._target = target + + def start(self): + if self._target is not None: + self._target() + + def is_alive(self): + return False + + if immediate_threads: + monkeypatch.setattr(server.threading, "Thread", _ImmediateThread) + monkeypatch.setattr(server, "_emit", lambda *_args, **_kwargs: None) + monkeypatch.setattr(server, "make_stream_renderer", lambda _cols: None) + monkeypatch.setattr(server, "render_message", lambda _raw, _cols: None) + monkeypatch.setattr(server, "_wire_callbacks", lambda _sid: None) + monkeypatch.setattr(server, "_sync_agent_model_with_config", lambda *_args: None) + monkeypatch.setattr(server, "_session_cwd", lambda _session: str(tmp_path)) + monkeypatch.setattr(server, "_register_session_cwd", lambda _session: None) + monkeypatch.setattr(server, "_set_session_context", lambda *_args, **_kwargs: []) + monkeypatch.setattr(server, "_clear_session_context", lambda _tokens: None) + monkeypatch.setattr(server, "_session_info", lambda *_args: {}) + monkeypatch.setattr(server, "_get_usage", lambda _agent: {}) + monkeypatch.setattr( + server, "_sync_session_key_after_compress", lambda *_args, **_kwargs: None + ) + monkeypatch.setattr(server, "_drain_queued_prompt", lambda *_args: False) + monkeypatch.setattr(server, "_voice_tts_enabled", lambda: False) + monkeypatch.setattr(server, "_get_db", lambda: None) + + +class _RecordingAgent: + model = "test-model" + provider = "test-provider" + + def __init__(self, turns): + self._turns = turns + + def clear_interrupt(self): + return None + + def run_conversation( + self, prompt, conversation_history=None, stream_callback=None + ): + self._turns.append(prompt) + return {"final_response": "", "messages": []} + + +@pytest.mark.parametrize("exit_code", [0, 7]) +def test_run_prompt_submit_requeues_foreign_completion( + monkeypatch, tmp_path, exit_code +): import queue as _queue_mod - sess_a = _session(session_key="session-a") - sess_b = _session(session_key="session-b") - server._sessions["sid_a"] = sess_a - server._sessions["sid_b"] = sess_b + from tools.process_registry import process_registry - q: _queue_mod.Queue = _queue_mod.Queue() - # Mix of events: ours, foreign (live), orphan, global - q.put({"type": "completion", "session_id": "1", "session_key": "session-a", - "command": "a", "exit_code": 0, "output": "a"}) - q.put({"type": "completion", "session_id": "2", "session_key": "session-b", - "command": "b", "exit_code": 0, "output": "b"}) - q.put({"type": "completion", "session_id": "3", "session_key": "ghost", - "command": "c", "exit_code": 0, "output": "c"}) - q.put({"type": "completion", "session_id": "4", - "command": "d", "exit_code": 0, "output": "d"}) + _configure_immediate_prompt_run(monkeypatch, tmp_path) + turns = [] + session_a = _session(session_key="session-a") + session_b = _session( + session_key="session-b", + agent=_RecordingAgent(turns), + running=True, + ) + event = { + "type": "completion", + "session_id": f"proc_foreign_{exit_code}", + "session_key": "session-a", + "command": "safe-test-command", + "exit_code": exit_code, + "output": "foreign", + } + isolated_queue: _queue_mod.Queue = _queue_mod.Queue() + isolated_queue.put(event) + monkeypatch.setattr(process_registry, "completion_queue", isolated_queue) + server._sessions["sid_a"] = session_a + server._sessions["sid_b"] = session_b - results = server._drain_owned_notifications(q, sess_a) + try: + server._run_prompt_submit("rid-b", "sid_b", session_b, "session-b-turn") + + assert turns == ["session-b-turn"] + assert isolated_queue.get_nowait() == event + assert isolated_queue.empty() + finally: + server._sessions.pop("sid_a", None) + server._sessions.pop("sid_b", None) + process_registry._completion_consumed.discard(event["session_id"]) + + +def test_run_prompt_submit_delivers_completion_observed_by_poll(monkeypatch, tmp_path): + import queue as _queue_mod + + from tools.process_registry import process_registry + + _configure_immediate_prompt_run(monkeypatch, tmp_path) + turns = [] + session = _session( + session_key="session-a", + agent=_RecordingAgent(turns), + running=True, + ) + event = { + "type": "completion", + "session_id": "proc_polled", + "session_key": "session-a", + "command": "safe-test-command", + "exit_code": 0, + "output": "observed but not consumed", + } + isolated_queue: _queue_mod.Queue = _queue_mod.Queue() + isolated_queue.put(event) + monkeypatch.setattr(process_registry, "completion_queue", isolated_queue) + process_registry._completion_consumed.discard(event["session_id"]) + process_registry._poll_observed.add(event["session_id"]) + server._sessions["sid_a"] = session + + try: + server._run_prompt_submit("rid-a", "sid_a", session, "session-a-turn") + + assert turns[0] == "session-a-turn" + assert len(turns) == 2 + assert "proc_polled" in turns[1] + assert isolated_queue.empty() + finally: + server._sessions.pop("sid_a", None) + process_registry._completion_consumed.discard(event["session_id"]) + process_registry._poll_observed.discard(event["session_id"]) + + +def test_run_prompt_submit_requeues_all_unstarted_notifications_with_real_threading( + monkeypatch, tmp_path +): + import queue as _queue_mod + + from tools.process_registry import process_registry + + _configure_immediate_prompt_run( + monkeypatch, tmp_path, immediate_threads=False + ) + real_thread_class = threading.Thread + threads = [] + nested_started = threading.Event() + release_nested = threading.Event() + turns = [] + + def _recording_thread(*args, **kwargs): + thread = real_thread_class(*args, **kwargs) + threads.append(thread) + return thread + + class _BlockingNotificationAgent(_RecordingAgent): + def run_conversation( + self, prompt, conversation_history=None, stream_callback=None + ): + turns.append(prompt) + if "proc_batch_1" in prompt: + nested_started.set() + if not release_nested.wait(timeout=5): + raise TimeoutError("notification turn was not released") + return {"final_response": "", "messages": []} + + monkeypatch.setattr(server.threading, "Thread", _recording_thread) + session = _session( + session_key="session-a", + agent=_BlockingNotificationAgent(turns), + running=True, + ) + events = [ + { + "type": "completion", + "session_id": f"proc_batch_{index}", + "session_key": "session-a", + "command": "safe-test-command", + "exit_code": 0, + "output": f"owned-{index}", + } + for index in range(1, 4) + ] + isolated_queue: _queue_mod.Queue = _queue_mod.Queue() + for event in events: + isolated_queue.put(event) + monkeypatch.setattr(process_registry, "completion_queue", isolated_queue) + server._sessions["sid_a"] = session + + try: + server._run_prompt_submit("rid-a", "sid_a", session, "session-a-turn") + + assert nested_started.wait(timeout=5) + threads[0].join(timeout=5) + assert not threads[0].is_alive() + queued = [] + while not isolated_queue.empty(): + queued.append(isolated_queue.get_nowait()) + assert [event["session_id"] for event in queued] == [ + "proc_batch_2", + "proc_batch_3", + ] + finally: + release_nested.set() + for thread in threads: + thread.join(timeout=5) + server._sessions.pop("sid_a", None) + while not isolated_queue.empty(): + isolated_queue.get_nowait() + for event in events: + process_registry._completion_consumed.discard(event["session_id"]) + process_registry._poll_observed.discard(event["session_id"]) + + +def test_run_prompt_submit_delivers_completion_owned_through_compression_lineage( + monkeypatch, tmp_path +): + import queue as _queue_mod + + from tools.process_registry import process_registry + + class _CompressionDB: + def resolve_resume_session_id(self, key): + return "new-child-key" if key == "old-parent-key" else key + + _configure_immediate_prompt_run(monkeypatch, tmp_path) + monkeypatch.setattr(server, "_get_db", lambda: _CompressionDB()) + ownership_checks = [] + original_owns_event = server._session_owns_notification_event + + def _record_ownership_check(sid, checked_session, checked_event): + ownership_checks.append(checked_event["session_id"]) + return original_owns_event(sid, checked_session, checked_event) + + monkeypatch.setattr( + server, "_session_owns_notification_event", _record_ownership_check + ) + turns = [] + session = _session( + session_key="new-child-key", + agent=_RecordingAgent(turns), + running=True, + ) + event = { + "type": "completion", + "session_id": "proc_precompression", + "session_key": "old-parent-key", + "command": "safe-test-command", + "exit_code": 0, + "output": "owned", + } + isolated_queue: _queue_mod.Queue = _queue_mod.Queue() + isolated_queue.put(event) + monkeypatch.setattr(process_registry, "completion_queue", isolated_queue) + server._sessions["sid_b"] = session + + try: + server._run_prompt_submit("rid-b", "sid_b", session, "session-b-turn") + + assert turns[0] == "session-b-turn" + assert len(turns) == 2 + assert "proc_precompression" in turns[1] + assert ownership_checks == ["proc_precompression"] + assert isolated_queue.empty() + finally: + server._sessions.pop("sid_b", None) + process_registry._completion_consumed.discard(event["session_id"]) + + +def test_run_prompt_submit_prefers_origin_ui_session_id(monkeypatch, tmp_path): + import queue as _queue_mod + + from tools.process_registry import process_registry + + _configure_immediate_prompt_run(monkeypatch, tmp_path) + ownership_checks = [] + original_owns_event = server._session_owns_notification_event + + def _record_ownership_check(sid, checked_session, checked_event): + ownership_checks.append(checked_event["session_id"]) + return original_owns_event(sid, checked_session, checked_event) + + monkeypatch.setattr( + server, "_session_owns_notification_event", _record_ownership_check + ) + turns = [] + session = _session( + session_key="current-key", + agent=_RecordingAgent(turns), + running=True, + ) + event = { + "type": "completion", + "session_id": "proc_origin_owned", + "session_key": "stale-durable-key", + "origin_ui_session_id": "sid_b", + "command": "safe-test-command", + "exit_code": 0, + "output": "owned", + } + isolated_queue: _queue_mod.Queue = _queue_mod.Queue() + isolated_queue.put(event) + monkeypatch.setattr(process_registry, "completion_queue", isolated_queue) + server._sessions["sid_b"] = session + + try: + server._run_prompt_submit("rid-b", "sid_b", session, "session-b-turn") + + assert turns[0] == "session-b-turn" + assert len(turns) == 2 + assert "proc_origin_owned" in turns[1] + assert ownership_checks == ["proc_origin_owned"] + assert isolated_queue.empty() + finally: + server._sessions.pop("sid_b", None) + process_registry._completion_consumed.discard(event["session_id"]) - # Only our event (session-a) + global (no session_key) should be returned. - assert len(results) == 2, f"expected 2 events (ours + global), got {len(results)}" - owned_sids = {evt.get("session_id") for evt, _ in results} - assert "1" in owned_sids # ours - assert "4" in owned_sids # global - - # session-b's event should be re-queued for its poller. - # Orphan should be dropped (not in results, not re-queued). - assert q.qsize() == 1, f"1 event (session-b) should be re-queued, got {q.qsize()}" - - # Now drain as session-b — should get its re-queued event. - results_b = server._drain_owned_notifications(q, sess_b) - assert len(results_b) == 1, f"session-b should pick up its re-queued event, got {len(results_b)}" - assert results_b[0][0].get("session_id") == "2" - assert q.empty(), f"queue should be empty after both drains, got {q.qsize()}" - - # Cleanup - server._sessions.pop("sid_a", None) - server._sessions.pop("sid_b", None) -def test_session_create_does_not_persist_empty_row(monkeypatch): """session.create must NOT eagerly write a DB row. Every TUI/desktop launch opens a session here just to paint the composer; diff --git a/tests/tools/test_process_registry.py b/tests/tools/test_process_registry.py index deee514370f9..817eaabcc8e2 100644 --- a/tests/tools/test_process_registry.py +++ b/tests/tools/test_process_registry.py @@ -1336,6 +1336,66 @@ def test_drain_notifications_skips_consumed(): process_registry.completion_queue.get_nowait() +def test_drain_notifications_can_deliver_poll_observed_for_gateway(registry): + event = { + "type": "completion", + "session_id": "proc_polled", + "session_key": "session-a", + "command": "safe-test-command", + "exit_code": 0, + "output": "observed but not consumed", + } + registry._poll_observed.add(event["session_id"]) + registry.completion_queue.put(event) + + try: + results = registry.drain_notifications( + session_key="session-a", + owns_event=lambda _event: True, + skip_poll_observed=False, + ) + + assert [raw for raw, _ in results] == [event] + finally: + registry._poll_observed.discard(event["session_id"]) + + +@pytest.mark.parametrize( + "skip_state", ["_poll_observed", "_completion_consumed"] +) +def test_drain_notifications_routes_foreign_before_local_skip( + registry, skip_state +): + event = { + "type": "completion", + "session_id": f"proc_foreign_{skip_state}", + "session_key": "session-a", + "command": "safe-test-command", + "exit_code": 0, + "output": "foreign", + } + ownership_calls = [] + getattr(registry, skip_state).add(event["session_id"]) + registry.completion_queue.put(event) + + def owns_event(checked_event): + ownership_calls.append(checked_event) + return False + + try: + results = registry.drain_notifications( + session_key="session-b", + owns_event=owns_event, + ) + + assert results == [] + assert ownership_calls == [event] + assert registry.completion_queue.get_nowait() == event + assert registry.completion_queue.empty() + finally: + getattr(registry, skip_state).discard(event["session_id"]) + + def test_drain_notifications_empty_queue(): from tools.process_registry import process_registry @@ -1346,6 +1406,151 @@ def test_drain_notifications_empty_queue(): assert results == [] +@pytest.mark.parametrize("exit_code", [0, 7]) +def test_drain_notifications_filters_addressed_completion_by_owns_event( + registry, exit_code +): + owned = { + "type": "completion", + "session_id": f"proc_owned_{exit_code}", + "session_key": "session-a", + "command": "safe-test-command", + "exit_code": exit_code, + "output": "owned", + } + foreign = { + "type": "completion", + "session_id": f"proc_foreign_{exit_code}", + "session_key": "session-b", + "command": "safe-test-command", + "exit_code": exit_code, + "output": "foreign", + } + registry.completion_queue.put(owned) + registry.completion_queue.put(foreign) + + results = registry.drain_notifications( + session_key="session-a", + owns_event=lambda event: event.get("session_key") == "session-a", + ) + + assert [event["session_id"] for event, _ in results] == [ + f"proc_owned_{exit_code}" + ] + assert registry.completion_queue.get_nowait() == foreign + assert registry.completion_queue.empty() + + +def test_drain_notifications_filters_addressed_completion_by_session_key(registry): + owned = { + "type": "completion", + "session_id": "proc_owned", + "session_key": "session-a", + "command": "safe-test-command", + "exit_code": 0, + "output": "owned", + } + foreign = { + "type": "completion", + "session_id": "proc_foreign", + "session_key": "session-b", + "command": "safe-test-command", + "exit_code": 0, + "output": "foreign", + } + registry.completion_queue.put(owned) + registry.completion_queue.put(foreign) + + results = registry.drain_notifications(session_key="session-a") + + assert [event["session_id"] for event, _ in results] == ["proc_owned"] + assert registry.completion_queue.get_nowait() == foreign + assert registry.completion_queue.empty() + + +def test_drain_notifications_session_key_filter_requeues_origin_only_event(registry): + event = { + "type": "completion", + "session_id": "proc_origin_only", + "origin_ui_session_id": "ui-session-a", + "command": "safe-test-command", + "exit_code": 0, + "output": "done", + } + registry.completion_queue.put(event) + + results = registry.drain_notifications(session_key="session-a") + + assert results == [] + assert registry.completion_queue.get_nowait() == event + assert registry.completion_queue.empty() + + +def test_drain_notifications_ownerless_completion_preserves_legacy_delivery(registry): + event = { + "type": "completion", + "session_id": "proc_ownerless", + "command": "safe-test-command", + "exit_code": 0, + "output": "ownerless", + } + registry.completion_queue.put(event) + + results = registry.drain_notifications( + session_key="session-a", + owns_event=lambda _event: False, + ) + + assert [raw for raw, _ in results] == [event] + assert registry.completion_queue.empty() + + +def test_drain_notifications_ownerless_async_delegation_still_requires_proof(registry): + event = { + "type": "async_delegation", + "delegation_id": "deleg_ownerless", + "goal": "task", + "status": "completed", + "summary": "done", + "api_calls": 1, + "duration_seconds": 0.1, + } + registry.completion_queue.put(event) + + results = registry.drain_notifications( + session_key="session-a", + owns_event=lambda _event: False, + ) + + assert results == [] + assert registry.completion_queue.get_nowait() == event + assert registry.completion_queue.empty() + + +def test_drain_notifications_completion_callback_exception_fails_closed(registry): + event = { + "type": "completion", + "session_id": "proc_callback_error", + "session_key": "session-a", + "command": "safe-test-command", + "exit_code": 0, + "output": "done", + } + registry.completion_queue.put(event) + + def broken(_event): + raise RuntimeError("ownership check exploded") + + results = registry.drain_notifications( + session_key="session-a", + owns_event=broken, + ) + + assert results == [] + assert registry.completion_queue.get_nowait() == event + assert registry.completion_queue.empty() + + def test_drain_notifications_filters_async_delegation_by_session_key(): """Async-delegation events should only be consumed by the matching session's drain. diff --git a/tools/process_registry.py b/tools/process_registry.py index 62dd8494e4ae..0c7616db0630 100644 --- a/tools/process_registry.py +++ b/tools/process_registry.py @@ -1146,8 +1146,10 @@ class ProcessRegistry: return False return True - def _drain_should_skip(self, session_id: str) -> bool: - """Whether the CLI drain should skip a completion event for this session. + def _drain_should_skip( + self, session_id: str, *, skip_poll_observed: bool = True + ) -> bool: + """Whether this drain should skip a completion event for this session. Skips when the agent has either truly consumed the output (wait/log → ``_completion_consumed``) or observed the exit inline via poll() @@ -1157,32 +1159,45 @@ class ProcessRegistry: check only ``is_completion_consumed`` so a read-only poll never suppresses their autonomous delivery turn (#10156). """ - return session_id in self._completion_consumed or session_id in self._poll_observed + return session_id in self._completion_consumed or ( + skip_poll_observed and session_id in self._poll_observed + ) def drain_notifications( - self, session_key: str = "", owns_event=None, + self, + session_key: str = "", + owns_event=None, + *, + skip_poll_observed: bool = True, ) -> "list[tuple[dict, str]]": """Pop all pending notification events and return formatted pairs. Returns a list of (raw_event, formatted_text) tuples. Skips completion events the agent already consumed via wait/log or - observed inline via poll() (see ``_drain_should_skip``). + observed inline via poll() (see ``_drain_should_skip``). Gateway/TUI + callers pass ``skip_poll_observed=False`` because read-only polling must + not suppress autonomous delivery there. - Async-delegation events carry a conversation payload, so draining one - into the wrong session is a cross-chat leak (#58684, #55578). Two - filter modes, strongest wins: + When a routing filter is supplied, addressed notifications must not be + drained into the wrong session. Async-delegation events always require + conversation payload; ordinary notifications require routing when they + carry ``session_key`` or ``origin_ui_session_id`` metadata. Two filter + modes are supported, strongest first: - ``owns_event(evt) -> bool``: positive-proof ownership callback. - When provided, an async-delegation event is consumed ONLY if the - callback returns True; everything else is re-queued for its owner. + When provided, a routed event is consumed ONLY if the callback + returns True; everything else is re-queued for its owner. The TUI passes its compression-chain-aware ownership check here so a post-compression session still claims its own pre-compression dispatches. - ``session_key``: plain key equality (CLI and other single-session - callers). Non-matching async-delegation events are re-queued. + callers). Non-matching addressed events are re-queued. With neither set, all events are consumed (legacy single-session - behavior, backward compatible). + behavior, backward compatible). Ownerless ordinary notifications also + retain that legacy behavior even when a filter is provided. When a + filter is provided, ownerless async-delegation events remain + fail-closed and require positive proof. """ results: "list[tuple[dict, str]]" = [] requeue: "list[dict]" = [] @@ -1191,39 +1206,43 @@ class ProcessRegistry: evt = self.completion_queue.get_nowait() except Exception: break - _evt_sid = evt.get("session_id", "") - if evt.get("type") == "completion" and self._drain_should_skip(_evt_sid): - continue - # Filter async-delegation events so they are not delivered to the - # wrong session/thread (#58684). Positive-proof callback beats - # bare key equality when the caller can provide one. - if evt.get("type") == "async_delegation": - if owns_event is not None: - try: - owned = bool(owns_event(evt)) - except Exception: - owned = False # fail closed — never leak on a broken check - if not owned: - requeue.append(evt) - continue - elif session_key: - evt_session_key = evt.get("session_key", "") or "" - if evt_session_key != session_key: - requeue.append(evt) - continue - elif evt.get("restored"): - # Legacy unfiltered drain (no ownership callback, no - # session key). That behavior was safe when the in-memory - # queue could only hold events created by this very - # process — but durable restore (#63494) re-enqueues - # completions from PREVIOUS processes at startup, so an - # unfiltered consumer here would adopt a dead, unrelated - # session's conversation payload (#64484). Fail closed: - # leave restored events queued (still 'pending' on disk) - # for a consumer that can positively prove ownership, - # e.g. the owning session's --resume. + # Positive-proof ownership beats bare key equality. Delegation + # payloads always require proof; ordinary events require it once + # they carry routing metadata. Ownerless ordinary events preserve + # legacy single-session delivery. + is_async_delegation = evt.get("type") == "async_delegation" + evt_session_key = str(evt.get("session_key") or "") + evt_origin_sid = str(evt.get("origin_ui_session_id") or "") + requires_positive_proof = is_async_delegation or bool( + evt_session_key or evt_origin_sid + ) + if owns_event is not None and requires_positive_proof: + try: + owned = bool(owns_event(evt)) + except Exception: + owned = False # fail closed — never leak on a broken check + if not owned: requeue.append(evt) continue + elif session_key and requires_positive_proof: + if evt_session_key != session_key: + requeue.append(evt) + continue + elif is_async_delegation and evt.get("restored"): + # Durable restore can enqueue previous-process payloads into a + # fresh registry. An unfiltered legacy drain cannot prove + # ownership, so leave those events queued for the owner. + requeue.append(evt) + continue + # Local consumed/observed state may suppress only events this + # session owns (or legacy ownerless ordinary events). Routing must + # happen first so a foreign session cannot drop the owner's event. + _evt_sid = evt.get("session_id", "") + if evt.get("type") == "completion" and self._drain_should_skip( + _evt_sid, skip_poll_observed=skip_poll_observed + ): + continue + text = format_process_notification(evt) if text: results.append((evt, text)) diff --git a/tui_gateway/server.py b/tui_gateway/server.py index 8761b777122f..6fc7fe38b246 100644 --- a/tui_gateway/server.py +++ b/tui_gateway/server.py @@ -8652,9 +8652,9 @@ def _session_owns_notification_event(sid: str, session: dict, evt: dict) -> bool minus its orphan-adoption fallback. An event owns-matches when its ``origin_ui_session_id`` is this live session, or its ``session_key`` (raw or resolved through the compression chain) matches this session's - key/lineage. Used as a fail-closed gate for async-delegation payloads: - "not provably elsewhere" is NOT good enough to inject a conversation - payload into this chat (#55578). + key/lineage. Used as the fail-closed gate for every addressed notification: + "not provably elsewhere" is NOT good enough to inject a payload into this + chat (#55578). """ if session.get("_finalized"): return False @@ -8679,6 +8679,14 @@ def _session_owns_notification_event(sid: str, session: dict, evt: dict) -> bool return resolved_key in current_keys +def _notification_event_requires_owner(evt: dict) -> bool: + """Whether ``evt`` must be positively claimed before TUI delivery.""" + return evt.get("type") == "async_delegation" or bool( + str(evt.get("origin_ui_session_id") or "") + or str(evt.get("session_key") or "") + ) + + def _notification_event_dedup_key(evt: dict) -> tuple: """Return the UI-emission identity for a process notification event. @@ -8716,83 +8724,6 @@ def _notification_event_dedup_key(evt: dict) -> tuple: return (evt_sid, evt_type) -def _drain_owned_notifications( - completion_queue: "queue.Queue", - session: dict, -) -> "list[tuple[dict, str]]": - """Pop all pending notification events, keeping only those owned by this session. - - ``process_registry.drain_notifications()`` pops every event from the global - queue regardless of ``session_key`` — a turn finishing in session B would - consume an event started by session A. This wrapper applies the same - ownership check used by ``_notification_poller_loop``, so only events that - belong to *this* session (plus global/system events with no ``session_key``) - are dispatched here. Events owned by another live session are re-queued; - orphaned events (owner gone) are dropped (#42674, #35652). - - Returns the same ``[(raw_event, formatted_text)]`` shape as - ``drain_notifications()``, filtered to this session's events only. - """ - from tools.process_registry import format_process_notification, process_registry - - _my_key = str(session.get("session_key") or "") - - # Snapshot live session keys (excluding our own) for ownership routing. - # Must be computed fresh so a just-closed session isn't treated as live. - try: - with _sessions_lock: - snapshot = list(_sessions.values()) - except Exception: - snapshot = [] - live_owner_keys = { - str(s.get("session_key") or "") - for s in snapshot - if s is not session and str(s.get("session_key") or "") - } - - owned: list[tuple[dict, str]] = [] - requeue: list[dict] = [] - - while not completion_queue.empty(): - try: - evt = completion_queue.get_nowait() - except Exception: - break - - _evt_key = str(evt.get("session_key") or "") - if not _evt_key: - # Global/system event with no owner — handle here. - pass - elif _evt_key == _my_key: - # Owned by this session — handle here. - pass - elif _evt_key in live_owner_keys: - # Owned by another live session — requeue. - requeue.append(evt) - continue - else: - # Orphaned event (owner gone) — drop silently. - logger.debug( - "Dropping orphaned background notification for " - "session_key=%s (owner gone, current=%s)", - _evt_key, _my_key, - ) - continue - - _evt_sid = evt.get("session_id", "") - if evt.get("type") == "completion" and process_registry.is_completion_consumed(_evt_sid): - continue - text = format_process_notification(evt) - if text: - owned.append((evt, text)) - - # Re-queue events for live sessions so their pollers can handle them. - for evt in requeue: - completion_queue.put(evt) - - return owned - - def _notification_poller_loop( stop_event: threading.Event, sid: str, session: dict ) -> None: @@ -8802,15 +8733,12 @@ def _notification_poller_loop( status.update (kind=process) for user visibility, then chains an agent turn via _run_prompt_submit if the session is idle. - NOTE: The completion_queue is global (one per process). If multiple - TUI sessions coexist, whichever poller wakes first grabs the event, - even if the process was started by a different session. This matches - CLI/gateway behavior (single session per process). + The completion_queue is process-global. In multi-session Desktop each + poller requeues events owned by another live session and drops addressed + events whose owner is gone; ownerless legacy notifications remain global. """ from tools.process_registry import process_registry, format_process_notification - _my_key = str(session.get("session_key") or "") - _emitted = set() # dedup re-queued events so same completion isn't emitted 50 times while session is busy while not stop_event.is_set() and not session.get("_finalized"): try: @@ -8828,46 +8756,28 @@ def _notification_poller_loop( time.sleep(0.1) continue - # Fail closed for async-delegation results (#55578): these carry a - # conversation payload, and injecting one into any chat other than the - # one that commissioned it is a hard cross-session leak. The - # belongs-elsewhere check above already re-queued events owned by - # another LIVE session; what reaches here is either ours or an - # orphan whose owner is gone. Orphaned delegation payloads are - # DROPPED, not adopted — the subagent's summary is already persisted - # in the delegation records/output store, so nothing is lost, whereas - # a wrong-chat injection is unrecoverable. Non-delegation events - # (background process completions etc.) keep the historical - # adopt-orphans behavior. - if evt.get("type") == "async_delegation" and not _session_owns_notification_event( - sid, session, evt - ): - logger.warning( - "async-delegation completion %s has no live owner " - "(origin=%r key=%r); dropping from injection instead of " - "delivering to session %s (#55578 fail-closed; result " - "remains in the delegation records)", - evt.get("delegation_id", "?"), + # What reaches here is not owned by another LIVE session. Addressed + # events still require positive proof before injection: exact UI origin, + # direct durable key, or compression lineage. If none proves ownership, + # the event is orphaned and must not be adopted by this chat. Truly + # ownerless ordinary notifications retain legacy global delivery. + requires_owner = _notification_event_requires_owner(evt) + if requires_owner and not _session_owns_notification_event(sid, session, evt): + log = ( + logger.warning + if evt.get("type") == "async_delegation" + else logger.debug + ) + log( + "Dropping unowned %s notification (origin=%r key=%r) instead " + "of delivering to session %s", + evt.get("type", "completion"), str(evt.get("origin_ui_session_id") or ""), str(evt.get("session_key") or ""), sid, ) continue - # Orphan guard: _notification_event_belongs_elsewhere returns False for - # events whose owner is no longer live (session closed / /new'd away). - # Previously these orphans were consumed by whichever poller dequeued - # them — injecting an unrelated background-process notification into the - # wrong session's transcript (#42674, #35652). Drop them instead. - _evt_key = str(evt.get("session_key") or "") - if _evt_key and _evt_key != _my_key: - logger.debug( - "Dropping orphaned background notification for " - "session_key=%s (owner gone, current=%s)", - _evt_key, _my_key, - ) - continue - _evt_sid = evt.get("session_id", "") if evt.get("type") == "completion" and process_registry.is_completion_consumed(_evt_sid): continue @@ -8933,23 +8843,21 @@ def _notification_poller_loop( if _notification_event_belongs_elsewhere(sid, session, evt): deferred.append(evt) continue - # Same fail-closed rule as the live loop: an orphaned async-delegation - # payload is never adopted by a foreign session — defer it (a later - # resume of the owner's lineage can still claim it) rather than - # injecting another chat's conversation here (#55578). - if evt.get("type") == "async_delegation" and not _session_owns_notification_event( - sid, session, evt - ): - deferred.append(evt) - continue - # Orphan guard: same check as the main loop above. - _evt_key = str(evt.get("session_key") or "") - if _evt_key and _evt_key != _my_key: - logger.debug( - "Dropping orphaned background notification for " - "session_key=%s (owner gone, current=%s)", - _evt_key, _my_key, - ) + # Same positive-proof rule as the live loop. Preserve the existing + # shutdown behavior for orphaned delegation payloads by deferring them + # for a later resume; ordinary addressed orphans are dropped. + requires_owner = _notification_event_requires_owner(evt) + if requires_owner and not _session_owns_notification_event(sid, session, evt): + if evt.get("type") == "async_delegation": + deferred.append(evt) + else: + logger.debug( + "Dropping unowned %s notification during shutdown drain " + "(origin=%r key=%r)", + evt.get("type", "completion"), + str(evt.get("origin_ui_session_id") or ""), + str(evt.get("session_key") or ""), + ) continue _evt_sid = evt.get("session_id", "") if evt.get("type") == "completion" and process_registry.is_completion_consumed(_evt_sid): @@ -9515,27 +9423,28 @@ def _run_prompt_submit(rid, sid: str, session: dict, text: Any) -> None: # The background poller handles between-turn delivery; this is # the safety net for events that arrived mid-turn. # - # Ownership filter (#42674, #35652): drain_notifications() pops - # every event from the global queue regardless of session_key. - # A turn finishing in session B must not consume an event that - # belongs to session A. Owned events are re-queued so the - # owning session's poller can handle them; orphaned events - # (owner gone) are dropped silently. + # Ownership filter (#42674, #35652): a turn finishing in session B + # must not consume an event that belongs to session A. The registry + # requeues every addressed event this session cannot positively claim; + # the poller then delivers it to a live owner or drops an orphan. try: from tools.process_registry import process_registry # Positive-proof ownership (compression-chain aware) — the same # fail-closed gate the poller uses, so the post-turn drain can't - # adopt another session's (or an orphan's) delegation payload, - # while a post-compression session still claims its own - # pre-compression dispatches (#55578). - for _evt, synth in process_registry.drain_notifications( + # adopt another session's addressed notification while a + # post-compression session still claims its own pre-compression + # dispatches (#55578). + drained = process_registry.drain_notifications( session_key=session.get("session_key", ""), owns_event=lambda e: _session_owns_notification_event(sid, session, e), - ): + skip_poll_observed=False, + ) + for index, (_evt, synth) in enumerate(drained): with session["history_lock"]: if session.get("running"): - process_registry.completion_queue.put(_evt) + for pending_evt, _pending_synth in drained[index:]: + process_registry.completion_queue.put(pending_evt) break session["running"] = True from tools.async_delegation import (