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
11 changes: 10 additions & 1 deletion gateway/run.py
Original file line number Diff line number Diff line change
Expand Up @@ -985,7 +985,16 @@ def _clarify_send_then_wait(fut, *, clarify_id: str, session_key: str, clarify_m
timeout = clarify_mod.get_clarify_timeout()
response = clarify_mod.wait_for_response(clarify_id, timeout=float(timeout))
if response is None or response == "":
# Timeout or session-boundary cancellation
# wait_for_response also returns on the thread-scoped interrupt
# (#83889). Report that outcome before treating the empty response as
# a timeout or session-boundary cancellation.
try:
from tools.interrupt import is_interrupted

if is_interrupted():
return "[interrupted by user]"
except Exception:
pass
return f"[user did not respond within {int(timeout / 60)}m]"
return response

Expand Down
16 changes: 16 additions & 0 deletions tests/gateway/test_clarify_send_timeout_ambiguity.py
Original file line number Diff line number Diff line change
Expand Up @@ -154,6 +154,22 @@ def test_no_response_returns_timeout_sentinel():
)


def test_interrupted_wait_returns_interrupt_sentinel(monkeypatch):
fut = MagicMock()
fut.result.return_value = _Result(True)
clarify_mod = MagicMock()
clarify_mod.get_clarify_timeout.return_value = 600
clarify_mod.wait_for_response.return_value = None
monkeypatch.setattr("tools.interrupt.is_interrupted", lambda: True)

assert (
_clarify_send_then_wait(
fut, clarify_id="cid123", session_key="sk", clarify_mod=clarify_mod
)
== "[interrupted by user]"
)


# --- Definitive failures keep their diagnostic detail in the log ----------


Expand Down
117 changes: 117 additions & 0 deletions tests/tools/test_clarify_interrupt_wait.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,117 @@
"""wait_for_response must observe the thread-scoped interrupt (#83889 RC1).

Salvage of #25506 (@liuhao1024): the wait loop checks
``tools.interrupt.is_interrupted()`` once per slice. Without it, an
interrupted agent thread holds until the full clarify timeout (600s default),
or forever in unlimited mode — the session-boundary ``clear_session`` cleanup
only runs after the turn ends, which cannot happen while the thread is
blocked here. Family E anchor in the stall triage (#84047).
"""

from __future__ import annotations

import threading
import time

import tools.clarify_gateway as cg
from tools.interrupt import set_interrupt


_ID_COUNTER = iter(range(10_000))


def _register(session_key: str = "agent:main:test:dm:1"):
clarify_id = f"clarify-test-{next(_ID_COUNTER)}"
return cg.register(
clarify_id=clarify_id,
session_key=session_key,
question="q",
choices=None,
)


def _wait_in_thread(clarify_id: str, timeout: float):
"""Run the wait in a worker thread; expose its ident for set_interrupt."""
result = {}
ready = threading.Event()

def run():
result["tid"] = threading.get_ident()
ready.set()
result["response"] = cg.wait_for_response(clarify_id, timeout)

t = threading.Thread(target=run, daemon=True)
t.start()
assert ready.wait(timeout=5.0)
return t, result


def test_interrupt_unblocks_bounded_wait():
entry = _register()
t, result = _wait_in_thread(entry.clarify_id, timeout=600.0)
try:
time.sleep(0.1)
assert t.is_alive(), "wait returned before any interrupt/resolve"
set_interrupt(True, thread_id=result["tid"])
t.join(timeout=5.0)
assert not t.is_alive(), "interrupt did not unblock the bounded wait"
assert result["response"] is None
finally:
set_interrupt(False, thread_id=result["tid"])


def test_interrupt_unblocks_unlimited_wait():
"""timeout<=0 is the forever-wedge: no deadline ever fires."""
entry = _register()
t, result = _wait_in_thread(entry.clarify_id, timeout=0.0)
try:
time.sleep(0.1)
assert t.is_alive()
set_interrupt(True, thread_id=result["tid"])
t.join(timeout=5.0)
assert not t.is_alive(), "interrupt did not unblock the unlimited wait"
assert result["response"] is None
finally:
set_interrupt(False, thread_id=result["tid"])


def test_resolve_still_wins_and_returns_response():
entry = _register()
t, result = _wait_in_thread(entry.clarify_id, timeout=600.0)
time.sleep(0.05)
assert cg.resolve_gateway_clarify(entry.clarify_id, "picked A") is True
t.join(timeout=5.0)
assert not t.is_alive()
assert result["response"] == "picked A"


def test_response_resolved_before_interrupt_wins_deterministically():
entry = _register()
assert cg.resolve_gateway_clarify(entry.clarify_id, "picked A") is True
set_interrupt(True)
try:
assert cg.wait_for_response(entry.clarify_id, timeout=600.0) == "picked A"
finally:
set_interrupt(False)


def test_timeout_still_works_without_interrupt():
entry = _register()
t, result = _wait_in_thread(entry.clarify_id, timeout=0.2)
t.join(timeout=5.0)
assert not t.is_alive()
assert result["response"] is None


def test_entry_cleaned_up_after_interrupt():
"""An interrupted wait must not leak its entry/session-index rows."""
key = "agent:main:test:dm:cleanup"
entry = _register(session_key=key)
set_interrupt(True) # current thread; wait aborts on first slice check
try:
response = cg.wait_for_response(entry.clarify_id, 600.0)
finally:
set_interrupt(False)
assert response is None
# A second resolve should find nothing — the entry is gone.
assert cg.resolve_gateway_clarify(entry.clarify_id, "late") is False
30 changes: 28 additions & 2 deletions tools/clarify_gateway.py
Original file line number Diff line number Diff line change
Expand Up @@ -105,7 +105,7 @@ def register(


def wait_for_response(clarify_id: str, timeout: float) -> Optional[str]:
"""Block on the entry's event until resolved or timeout fires.
"""Block on the entry's event until resolved, timeout, or interrupt.

Polls in 1-second slices so the agent's inactivity heartbeat keeps
firing — without this, ``Event.wait(timeout=600)`` blocks the thread
Expand All @@ -116,7 +116,17 @@ def wait_for_response(clarify_id: str, timeout: float) -> Optional[str]:
heartbeat still fires each slice so inactivity watchdogs don't kill a live
prompt.

Returns the resolved response string, or ``None`` on timeout.
The thread-scoped interrupt flag (``tools.interrupt.is_interrupted``) is
checked once per slice. This is how ``/stop`` / an interrupt-mode message
unblocks a waiting agent thread: ``AIAgent.interrupt()`` propagates the
thread-scoped signal to active tool workers, but the session-boundary
``clear_session`` cleanup only runs after the turn ends — which never
happens while this wait is blocked (#83889). Without the check, an
interrupted agent sits here until the full clarify timeout, or forever
in unlimited mode.

Returns the resolved response string, or ``None`` on timeout/interrupt.
A response that raced the interrupt and resolved first still wins.
"""
with _lock:
entry = _entries.get(clarify_id)
Expand All @@ -128,11 +138,27 @@ def wait_for_response(clarify_id: str, timeout: float) -> Optional[str]:
except Exception: # pragma: no cover - optional
touch_activity_if_due = None

try:
from tools.interrupt import is_interrupted as _is_interrupted
except Exception: # pragma: no cover - defensive: never re-wedge the wait
_is_interrupted = None

# 0 / negative → unlimited: no deadline, poll forever in 1s slices.
unlimited = timeout is None or float(timeout) <= 0.0
deadline = None if unlimited else time.monotonic() + float(timeout)
activity_state = {"last_touch": time.monotonic(), "start": time.monotonic()}
while True:
# A response that completed before the interrupt wins. Besides being
# the least surprising user-visible result, checking the event first
# makes the interrupt-vs-resolve contract deterministic.
if entry.event.is_set():
break
if _is_interrupted is not None:
try:
if _is_interrupted():
break
except Exception: # pragma: no cover - defensive
pass
if deadline is None:
slice_s = 1.0
else:
Expand Down
Loading