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
264 changes: 263 additions & 1 deletion tests/test_tui_gateway_queue_on_busy.py
Original file line number Diff line number Diff line change
Expand Up @@ -79,6 +79,237 @@ def test_busy_interrupt_mode_redirects_active_turn(monkeypatch):
assert session.get("queued_prompt") is None


def test_successful_redirect_drops_queued_duplicate_of_inflight_user(monkeypatch):
"""#84417: correcting a live turn must not re-fire the original prompt from queue.

When the live turn's original user text is also sitting in the server queue
(e.g. a second prompt.submit of the same text while redirect was not yet
possible), a later successful redirect of a *new* correction Q must purge
that self-duplicate. Otherwise post-turn ``_drain_queued_prompt`` starts a
second agent turn with the old prompt P after Q has already been handled.
"""
monkeypatch.setattr(server, "_load_busy_input_mode", lambda: "interrupt")
agent = types.SimpleNamespace(
_supports_active_turn_redirect=True,
redirect=lambda text: True,
interrupt=lambda *a, **k: (_ for _ in ()).throw(
AssertionError("redirect must not hard-interrupt")
),
)
session = _session(agent=agent, running=True)
original = "deepseek released a new flash model β€” I changed all settings to flash"
session["inflight_turn"] = {
"user": original,
"assistant": "partial",
"streaming": True,
"error": "",
}
# Stale self-duplicate of the live turn (would re-fire after settle).
session["queued_prompt"] = {"text": original, "transport": "ws-1"}
session["queued_prompts"] = [
{"text": original, "transport": "ws-1"},
{"text": "unrelated later task", "transport": "ws-1"},
]

resp = server._handle_busy_submit(
"r1", "sid", session, "what about the pricing instead?", "ws-1"
)

assert resp["result"]["status"] == "redirected"
# Self-duplicates of the live original must be gone.
assert session.get("queued_prompt") == {
"text": "unrelated later task",
"transport": "ws-1",
}
assert not session.get("queued_prompts")


def test_successful_redirect_preserves_unrelated_queued_followups(monkeypatch):
"""A legitimate next-turn queue entry must survive a mid-turn redirect."""
monkeypatch.setattr(server, "_load_busy_input_mode", lambda: "interrupt")
agent = types.SimpleNamespace(
_supports_active_turn_redirect=True,
redirect=lambda text: True,
interrupt=lambda *a, **k: (_ for _ in ()).throw(
AssertionError("redirect must not hard-interrupt")
),
)
session = _session(agent=agent, running=True)
session["inflight_turn"] = {
"user": "live turn P",
"assistant": "",
"streaming": True,
"error": "",
}
session["queued_prompt"] = {"text": "run this after", "transport": "ws-1"}

resp = server._handle_busy_submit("r1", "sid", session, "correction Q", "ws-1")

assert resp["result"]["status"] == "redirected"
assert session.get("queued_prompt") == {
"text": "run this after",
"transport": "ws-1",
}


def test_enqueue_skips_text_duplicate_of_inflight_user():
"""#84417 defense: do not admit a self-duplicate of the live user prompt."""
session = _session()
session["inflight_turn"] = {
"user": "live turn P",
"assistant": "",
"streaming": True,
"error": "",
}

server._enqueue_prompt(session, "live turn P", "ws-1")
assert session.get("queued_prompt") is None

server._enqueue_prompt(session, "different follow-up", "ws-1")
assert session["queued_prompt"] == {
"text": "different follow-up",
"transport": "ws-1",
}


def test_enqueue_followup_does_not_merge_stale_inflight_self_duplicate():
"""#84417: scrub P before merging so drain cannot re-fire ``P\\n\\nQ``."""
session = _session()
session["inflight_turn"] = {
"user": "P",
"assistant": "",
"streaming": True,
"error": "",
}
# Pre-existing stale self-duplicate (e.g. admitted before inflight was set).
session["queued_prompt"] = {"text": "P", "transport": "ws-1"}

server._enqueue_prompt(session, "Q", "ws-1")

assert session.get("queued_prompt") == {"text": "Q", "transport": "ws-1"}
assert not session.get("queued_prompts")


def test_drop_rewrites_merged_inflight_prefix_to_followup_only():
"""Already-merged ``P\\n\\nQ`` slots keep Q and drop the live original."""
session = _session()
session["inflight_turn"] = {
"user": "P",
"assistant": "",
"streaming": True,
"error": "",
}
session["queued_prompt"] = {"text": "P\n\nQ", "transport": "ws-1"}

server._drop_queued_duplicates_of_inflight_user(session)

assert session.get("queued_prompt") == {"text": "Q", "transport": "ws-1"}


def test_hard_interrupt_queue_path_scrubs_stale_inflight_self_duplicate(monkeypatch):
"""#84417: interrupt+queue of Q must not leave P ahead of Q in the FIFO."""
monkeypatch.setattr(server, "_load_busy_input_mode", lambda: "interrupt")
interrupts = []
agent = types.SimpleNamespace(
_supports_active_turn_redirect=True,
redirect=lambda text: False, # force hard-interrupt fallback
interrupt=lambda *a, **k: interrupts.append(True),
)
session = _session(agent=agent, running=True)
session["inflight_turn"] = {
"user": "P",
"assistant": "",
"streaming": True,
"error": "",
}
session["queued_prompt"] = {"text": "P", "transport": "ws-1"}

resp = server._handle_busy_submit("r1", "sid", session, "Q", "ws-1")

assert resp["result"]["status"] == "queued"
assert session.get("queued_prompt") == {"text": "Q", "transport": "ws-1"}
assert not session.get("queued_prompts")
# Interrupt is async-threaded; policy still enqueued Q after scrubbing P.


def test_redirect_then_drain_does_not_re_fire_original_p(monkeypatch):
"""#84417 drain-level: after redirect(Q), settle must not start a second P."""
monkeypatch.setattr(server, "_load_busy_input_mode", lambda: "interrupt")
fired = []
agent = types.SimpleNamespace(
_supports_active_turn_redirect=True,
redirect=lambda text: True,
interrupt=lambda *a, **k: (_ for _ in ()).throw(
AssertionError("redirect must not hard-interrupt")
),
)
session = _session(agent=agent, running=True)
session["inflight_turn"] = {
"user": "P",
"assistant": "partial",
"streaming": True,
"error": "",
}
session["queued_prompt"] = {"text": "P", "transport": "ws-1"}

resp = server._handle_busy_submit("r1", "sid", session, "Q", "ws-1")
assert resp["result"]["status"] == "redirected"
assert session.get("queued_prompt") is None

# Turn settles (running cleared in finally) β€” drain must be a no-op.
session["running"] = False
monkeypatch.setattr(
server,
"_run_prompt_submit",
lambda rid, sid, session, text, **kwargs: fired.append(text),
)
monkeypatch.setattr(server, "_session_uses_compute_host", lambda _s: False)

assert server._drain_queued_prompt("r2", "sid", session) is False
assert fired == []


def test_compress_session_rotation_bumps_queued_prompt_generation(monkeypatch):
"""#84417 belt: rotation invalidates in-flight drain claims on the parent key.

Queue *contents* survive (a legitimate follow-up must still run after
compression); only the generation counter advances so a drain that claimed
under the pre-rotation key cannot dispatch after re-anchor.
"""
monkeypatch.setattr(server, "_transfer_active_session_slot", lambda *a, **k: True)
monkeypatch.setattr(server, "_restart_slash_worker", lambda *a, **k: None)
agent = types.SimpleNamespace(session_id="child-after-rotation")
session = _session(agent=agent, session_key="parent-before-rotation")
session["_queued_prompt_generation"] = 3
session["queued_prompt"] = {"text": "run after compress", "transport": "ws-1"}

server._sync_session_key_after_compress("sid", session, clear_pending_title=False)

assert session["session_key"] == "child-after-rotation"
assert session["_queued_prompt_generation"] == 4
# Follow-up kept β€” only the claim generation bumped.
assert session["queued_prompt"] == {
"text": "run after compress",
"transport": "ws-1",
}


def test_compress_no_rotation_does_not_bump_queue_generation(monkeypatch):
"""No-op when agent.session_id already matches session_key."""
monkeypatch.setattr(
server,
"_transfer_active_session_slot",
lambda *a, **k: (_ for _ in ()).throw(AssertionError("no transfer")),
)
agent = types.SimpleNamespace(session_id="same-key")
session = _session(agent=agent, session_key="same-key")
session["_queued_prompt_generation"] = 2

server._sync_session_key_after_compress("sid", session)

assert session["_queued_prompt_generation"] == 2





Expand Down Expand Up @@ -299,7 +530,36 @@ def _boom(*a, **k):


def test_drain_does_not_dispatch_a_prompt_cancelled_after_claim(monkeypatch):
session = _session(queued_prompt={"text": "B", "transport": None})
"""Generation cancel aborts dispatch but must restore the claimed head.

Compress re-anchor / Stop bump generation between claim and check. Dropping
the envelope would silently lose a legitimate follow-up (#84417 belt).
"""
session = _session(
queued_prompt={"text": "B", "transport": "ws-1"},
queued_prompts=[{"text": "C", "transport": "ws-1"}],
)
monkeypatch.setattr(
server,
"_session_uses_compute_host",
lambda _session: session.__setitem__("_queued_prompt_generation", 1) or False,
)
monkeypatch.setattr(
server,
"_run_prompt_submit",
lambda *args, **kwargs: (_ for _ in ()).throw(AssertionError("must not dispatch")),
)

assert server._drain_queued_prompt("r1", "sid", session) is True
assert session["running"] is False
# Claimed B restored first; C that advanced into the slot is behind it.
assert session.get("queued_prompt") == {"text": "B", "transport": "ws-1"}
assert session.get("queued_prompts") == [{"text": "C", "transport": "ws-1"}]


def test_drain_restores_claimed_prompt_when_generation_bumps_mid_claim(monkeypatch):
"""Single-item queue: generation cancel must not empty the queue."""
session = _session(queued_prompt={"text": "follow-up Q", "transport": None})
monkeypatch.setattr(
server,
"_session_uses_compute_host",
Expand All @@ -313,6 +573,8 @@ def test_drain_does_not_dispatch_a_prompt_cancelled_after_claim(monkeypatch):

assert server._drain_queued_prompt("r1", "sid", session) is True
assert session["running"] is False
assert session.get("queued_prompt") == {"text": "follow-up Q", "transport": None}
assert not session.get("queued_prompts")


def test_drain_does_not_clear_stop_after_its_final_generation_check(monkeypatch):
Expand Down
83 changes: 83 additions & 0 deletions tests/test_tui_gateway_server.py
Original file line number Diff line number Diff line change
Expand Up @@ -9033,6 +9033,89 @@ def test_session_redirect_calls_capable_core_agent(monkeypatch):
assert before is None or session["last_active"] >= before


def test_session_redirect_rpc_drops_queued_duplicate_of_inflight_user():
"""#84417: Desktop ``session.redirect`` must purge stale self-duplicates.

Production path: renderer steers via ``session.redirect`` (not
``prompt.submit``). A self-copy of the live original user text already in
the server queue must not survive a successful redirect β€” otherwise
post-turn ``_drain_queued_prompt`` restarts prompt P after Q is handled.
Unrelated next-turn envelopes stay.
"""
original = "deepseek released a new flash model β€” I changed all settings to flash"
agent = types.SimpleNamespace(
_supports_active_turn_redirect=True,
redirect=lambda text: True,
)
session = _session(agent=agent, running=True)
session["inflight_turn"] = {
"user": original,
"assistant": "partial",
"streaming": True,
"error": "",
}
session["queued_prompt"] = {"text": original, "transport": "ws-1"}
session["queued_prompts"] = [
{"text": original, "transport": "ws-1"},
{"text": "unrelated later task", "transport": "ws-1"},
]
server._sessions["sid"] = session
try:
resp = server.handle_request(
{
"id": "1",
"method": "session.redirect",
"params": {
"session_id": "sid",
"text": "what about the pricing instead?",
},
}
)
finally:
server._sessions.pop("sid", None)

assert resp["result"]["status"] == "redirected"
assert session["inflight_turn"]["user"] == original
assert session["inflight_turn"]["corrections"] == [
"what about the pricing instead?"
]
# Self-duplicates of the live original are gone; legitimate follow-up kept.
assert session.get("queued_prompt") == {
"text": "unrelated later task",
"transport": "ws-1",
}
assert not session.get("queued_prompts")


def test_session_redirect_build_window_scrubs_stale_p_when_queuing_q():
"""#84417: build-window queue of Q must not leave P ahead of Q."""
original = "live original P"
session = _session(running=True)
session["agent"] = None # async agent build window
session["inflight_turn"] = {
"user": original,
"assistant": "",
"streaming": True,
"error": "",
}
session["queued_prompt"] = {"text": original, "transport": "ws-1"}
server._sessions["sid"] = session
try:
resp = server.handle_request(
{
"id": "1",
"method": "session.redirect",
"params": {"session_id": "sid", "text": "correction Q"},
}
)
finally:
server._sessions.pop("sid", None)

assert resp["result"] == {"status": "queued", "text": "correction Q"}
assert session["queued_prompt"]["text"] == "correction Q"
assert not session.get("queued_prompts")


def test_session_redirect_records_correction_without_erasing_prompt():
"""A redirect must not overwrite the turn's original user text.

Expand Down
7 changes: 7 additions & 0 deletions tui_gateway/methods_session.py
Original file line number Diff line number Diff line change
Expand Up @@ -3245,6 +3245,10 @@ def _(rid, params: dict) -> dict:
# text has no user bubble β€” the "my message vanished on reload" loss.
with session["history_lock"]:
_record_inflight_correction(session, text)
# #84417: steer does not cancel the live original, but a server
# queue self-copy of that original must still not re-fire after
# settle (same class as redirect).
_drop_queued_duplicates_of_inflight_user(session)
session["last_active"] = time.time()
return _ok(rid, {"status": "queued" if accepted else "rejected", "text": text})

Expand Down Expand Up @@ -3281,6 +3285,9 @@ def _(rid, params: dict) -> dict:
if accepted:
with session["history_lock"]:
_record_inflight_correction(session, text)
# #84417: purge server-queue self-duplicates of the live original
# so post-turn drain cannot restart the pre-correction prompt.
_drop_queued_duplicates_of_inflight_user(session)
session["last_active"] = time.time()
return _ok(
rid,
Expand Down
Loading
Loading