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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
38 changes: 38 additions & 0 deletions tests/test_tui_gateway_server.py
Original file line number Diff line number Diff line change
Expand Up @@ -6736,3 +6736,41 @@ def test_reap_idle_sessions_closes_only_evictable(monkeypatch):
assert closed == [("stale", "idle_timeout")]
finally:
server._sessions.clear()


def test_drain_owned_notifications_requeues_foreign_live_session_events(monkeypatch):
current = {"session_key": "agent:main:tui:session:current"}
other = {"session_key": "agent:main:tui:session:other"}
foreign_evt = {
"type": "completion",
"session_id": "proc_foreign",
"session_key": "agent:main:tui:session:other",
}
owned_evt = {
"type": "completion",
"session_id": "proc_owned",
"session_key": "agent:main:tui:session:current",
}
global_evt = {"type": "watch_overflow_global", "session_id": "proc_global"}
requeued = []

class _Queue:
def put(self, evt):
requeued.append(evt)

class _Registry:
completion_queue = _Queue()

def drain_notifications(self):
return [(foreign_evt, "foreign"), (owned_evt, "owned"), (global_evt, "global")]

monkeypatch.setitem(server._sessions, "current", current)
monkeypatch.setitem(server._sessions, "other", other)
try:
drained = server._drain_owned_notifications(_Registry(), current)
finally:
server._sessions.pop("current", None)
server._sessions.pop("other", None)

assert drained == [(owned_evt, "owned"), (global_evt, "global")]
assert requeued == [foreign_evt]
24 changes: 23 additions & 1 deletion tui_gateway/server.py
Original file line number Diff line number Diff line change
Expand Up @@ -4698,6 +4698,28 @@ def _notification_event_belongs_elsewhere(session: dict, evt: dict) -> bool:
)


def _drain_owned_notifications(process_registry, session: dict) -> list[tuple[dict, str]]:
"""Drain pending process notifications owned by this TUI session.

``process_registry.drain_notifications()`` drains the global completion
queue. In Desktop/TUI, multiple live sessions share that queue, so a
post-turn drain in session B must not consume a completion event that was
started by session A. Foreign events are put back for their owning poller.
"""
owned: list[tuple[dict, str]] = []
deferred: list[dict] = []
for evt, text in process_registry.drain_notifications():
if _notification_event_belongs_elsewhere(session, evt):
deferred.append(evt)
continue
owned.append((evt, text))

for evt in deferred:
process_registry.completion_queue.put(evt)

return owned


def _notification_event_dedup_key(evt: dict) -> tuple:
"""Return the UI-emission identity for a process notification event.

Expand Down Expand Up @@ -5254,7 +5276,7 @@ def _stream(delta):
try:
from tools.process_registry import process_registry

for _evt, synth in process_registry.drain_notifications():
for _evt, synth in _drain_owned_notifications(process_registry, session):
with session["history_lock"]:
if session.get("running"):
process_registry.completion_queue.put(_evt)
Expand Down