fix(approval): withdraw the queue entry when no client can answer; commit choices under the lock
Why: for an old WebSocket client that never advertised server→client requests,
send_async already failed fast (on_result(None)) but _emit_approval_request ignored
it, so the approval wait — owned by tools.approval's queue, not server_requests —
still idled for the whole approvals.timeout (300s) with no prompt anywhere. The
same held for a -32601 error frame. on_result(None) now withdraws the queue entry
(withdraw_gateway_approval: cancelled cause, never a user deny) so
_await_gateway_decision returns at once; the agent sees a withdrawn prompt.
resolve_gateway_approval committed entry.result/reason/event.set() AFTER
releasing _lock, so _drop_entry (which reads result and leaves the queue under
the lock) could still pop-and-lose a choice the client was acked for. Every
committer (resolve, clear_session, unregister_gateway_notify) now commits inside
the same critical section that pops the entry, which is what _drop_entry's
docstring promised.
_drop_entry mapped a withdrawn wait to settle("set") — the raw poll-state token
went out as the request.cancel reason. It is now session_closed, and a choice
from another surface is resolved (both RequestCancelReason values).
Tests: restore test_server_request_error_response_fails_fast (an error frame
settles send() to None promptly); the approval fail-fast probe from the review
(silent WS peer, timeout 3 -> was 3.2s, now immediate); a lock-instrumented
resolve test; a settle-reason test. The test_protocol `server` fixture imports
server_requests before its sys.modules patch window so the module server.py binds
its sinks on is the one tests import — the PR's fail-fast test only passed in a
full-file run before (order dependency).
Part of #112548
This commit is contained in:
@@ -131,9 +131,8 @@ def unregister_gateway_notify(session_key: str) -> None:
|
||||
they don't hang forever (agent run finished or interrupted)."""
|
||||
with _lock:
|
||||
_gateway_notify_cbs.pop(session_key, None)
|
||||
entries = _gateway_queues.pop(session_key, [])
|
||||
for entry in entries:
|
||||
entry.event.set()
|
||||
for entry in _gateway_queues.pop(session_key, []):
|
||||
entry.event.set()
|
||||
|
||||
|
||||
def resolve_gateway_approval(session_key: str, choice: str,
|
||||
@@ -162,15 +161,34 @@ def resolve_gateway_approval(session_key: str, choice: str,
|
||||
targets = [queue.pop(0)]
|
||||
if not queue:
|
||||
_gateway_queues.pop(session_key, None)
|
||||
|
||||
for entry in targets:
|
||||
entry.result = choice
|
||||
if reason:
|
||||
entry.reason = reason
|
||||
entry.event.set()
|
||||
# Popping the entry and committing its outcome are ONE critical section: the waiter's
|
||||
# ``_drop_entry`` reads ``entry.result`` under this same lock after its deadline check, so a
|
||||
# choice acked to the client here can never be popped-and-lost as a timeout (#112548).
|
||||
for entry in targets:
|
||||
entry.result = choice
|
||||
if reason:
|
||||
entry.reason = reason
|
||||
entry.event.set()
|
||||
return len(targets)
|
||||
|
||||
|
||||
def withdraw_gateway_approval(session_key: str, request_id: str, cause: str) -> bool:
|
||||
"""Withdraw one pending approval nobody can answer (the only attached client cannot render it).
|
||||
The waiter wakes at once with ``cancelled=cause`` — a withdrawal, never a user deny — instead of
|
||||
idling for the whole approvals.timeout (#112548). False when it is no longer pending."""
|
||||
with _lock:
|
||||
queue = _gateway_queues.get(session_key, [])
|
||||
entry = next((e for e in queue if e.data.get("request_id") == request_id), None)
|
||||
if entry is None:
|
||||
return False
|
||||
queue.remove(entry)
|
||||
if not queue:
|
||||
_gateway_queues.pop(session_key, None)
|
||||
entry.cancelled = cause
|
||||
entry.event.set()
|
||||
return True
|
||||
|
||||
|
||||
def list_gateway_approvals(session_key: str) -> list[dict]:
|
||||
"""Return replay-safe snapshots of unresolved approvals for one session."""
|
||||
with _lock:
|
||||
@@ -272,12 +290,11 @@ def clear_session(session_key: str) -> None:
|
||||
_session_approved.pop(session_key, None)
|
||||
_session_yolo.discard(session_key)
|
||||
_pending.pop(session_key, None)
|
||||
entries = _gateway_queues.pop(session_key, [])
|
||||
for entry in entries:
|
||||
# Cancel blocked waits now so the old run unwinds instead of idling until timeout;
|
||||
# the prompt was withdrawn, nobody denied it.
|
||||
entry.cancelled = "the session ended before the prompt was answered"
|
||||
entry.event.set()
|
||||
for entry in _gateway_queues.pop(session_key, []):
|
||||
# Cancel blocked waits now so the old run unwinds instead of idling until timeout;
|
||||
# the prompt was withdrawn, nobody denied it.
|
||||
entry.cancelled = "the session ended before the prompt was answered"
|
||||
entry.event.set()
|
||||
_release_permission_mode_dependents(session_key)
|
||||
# Session-persistent code kernels (local and remote) share this owner key and die at the same boundary so a
|
||||
# finished conversation cannot leak a live interpreter.
|
||||
|
||||
Reference in New Issue
Block a user