fix(streaming): interrupt abort reaches the in-flight stream socket
The /stop and /reset abort path shut down only the sockets the request client's pool sweep could see. Two shapes escaped it (#98974): the connection checked out for the in-flight body read (the pool sweep skips it — the stale-kill path already covered this via _shutdown_stale_attempt_socket, the interrupt path did not), and an abort that fires during create()'s connect/TLS window, when no socket exists yet, so the request came up afterwards with nothing to stop it. The log said "no sockets found; tcp_force_closed=0" and the self-hosted serve kept generating into a dropped consumer. Now _abort_for_interrupt also shuts down the attempt's own socket, and _chat_stream_created re-aborts when the attempt was cancelled before the response arrived. Both remain shutdown-only (never a cross-thread close, #30858): the worker unwinds and releases descriptors on its own thread. Direction shared with #98989 (@liuhao1024); the re-abort-on-create hunk is reimplemented against the current _StreamingCall shape.
This commit is contained in:
@@ -2857,6 +2857,23 @@ class _StreamingCall(StreamingWaitMonitor):
|
||||
self.agent._stream_diag_capture_response(self.clients.diag, response)
|
||||
self.agent._check_openrouter_cache_status(response)
|
||||
self._writer_token = claim_stream_writer(self.agent)
|
||||
self._reabort_if_cancelled(response)
|
||||
|
||||
def _reabort_if_cancelled(self, response: Any) -> None:
|
||||
"""Interrupt/stale abort that raced ``create()``: the one-shot pool sweep ran while
|
||||
the connect/TLS window held no socket yet (``tcp_force_closed=0``), so nothing stopped
|
||||
the request once it came up and the serve kept generating into a dropped consumer
|
||||
(#98974). Response headers prove the socket exists now — shut it down (shutdown-only,
|
||||
never a cross-thread close) so the worker unwinds as after a stale kill."""
|
||||
with self.stream_attempt_lock:
|
||||
current = int(self.stream_attempt_state["current"])
|
||||
cancelled = self._request_cancelled["value"] or current in self.stream_attempt_state["cancelled"]
|
||||
if not cancelled:
|
||||
return
|
||||
self._shutdown_stale_attempt_socket(response)
|
||||
if self._attempt_request_client is not None:
|
||||
self.agent._abort_request_openai_client(
|
||||
self._attempt_request_client, reason="cancelled_attempt_late_connect")
|
||||
|
||||
def _accept_chat_chunk(self, stream_attempt_id: int, chunk: Any) -> bool:
|
||||
with contextlib.suppress(Exception):
|
||||
@@ -3509,10 +3526,14 @@ class _StreamingCall(StreamingWaitMonitor):
|
||||
# transport error as a cancel, not a network error (#6600).
|
||||
self._request_cancelled["value"] = True
|
||||
logger.debug("Force-closing streaming httpx client due to interrupt (not a network error).")
|
||||
# Same as the stale kill: the pool sweep can miss the connection checked out for the
|
||||
# in-flight body read, so shut down the attempt's own socket too (#98974).
|
||||
_killed_response = self._attempt_stream_response
|
||||
with contextlib.suppress(Exception):
|
||||
self._cancel_current_stream_attempt("stream_interrupt_abort")
|
||||
# Kind-aware: only the request-local socket; the shared _anthropic_client is never closed here.
|
||||
self.clients.close_once("stream_interrupt_abort")
|
||||
self._shutdown_stale_attempt_socket(_killed_response)
|
||||
# Let the worker unwind Relay-managed scopes first; raising first lets
|
||||
# turn teardown race a still-open scope and corrupt the LIFO stack.
|
||||
if self.worker is not None:
|
||||
|
||||
74
tests/agent/test_stream_interrupt_abort_socket.py
Normal file
74
tests/agent/test_stream_interrupt_abort_socket.py
Normal file
@@ -0,0 +1,74 @@
|
||||
"""An interrupt abort must reach the in-flight stream's socket (#98974).
|
||||
|
||||
Report #98974: ``/stop`` / ``/reset`` logged ``OpenAI client aborted
|
||||
(stream_interrupt_abort, ..., tcp_force_closed=0, deferred_close=stranger_thread)
|
||||
— no sockets found`` and the self-hosted serve kept generating for minutes.
|
||||
Two shapes miss the socket: (a) the pool sweep skips the connection checked
|
||||
out for the in-flight body read — the stale kill already shuts down the
|
||||
attempt's own socket, the interrupt path did not; (b) the abort fires during
|
||||
``create()``'s connect/TLS window, before any socket exists, so nothing stops
|
||||
the request once headers arrive. Both stay shutdown-only (never a cross-thread
|
||||
``close()``, #30858).
|
||||
"""
|
||||
import socket as _socket
|
||||
from types import SimpleNamespace
|
||||
|
||||
import run_agent
|
||||
from agent import chat_completion_helpers as helpers
|
||||
|
||||
|
||||
def _agent():
|
||||
return run_agent.AIAgent(
|
||||
api_key="test-key", base_url="http://127.0.0.1:1/v1", model="m", provider="custom",
|
||||
quiet_mode=True, skip_context_files=True, skip_memory=True, enabled_toolsets=[], max_iterations=1,
|
||||
)
|
||||
|
||||
|
||||
def _call(agent):
|
||||
call = helpers._StreamingCall(agent, {"model": "m", "messages": [{"role": "user", "content": "hi"}]}, None)
|
||||
call._stream_stale_timeout = 5.0
|
||||
return call
|
||||
|
||||
|
||||
def _response_over(reader):
|
||||
"""Live httpx 0.28 wrapper shape down to the socket; ``close`` must never run."""
|
||||
def _no_close():
|
||||
raise AssertionError("stranger thread must never close the response")
|
||||
h11_conn = SimpleNamespace(_network_stream=SimpleNamespace(_sock=reader),
|
||||
_stream=None, _connection=None, _httpcore_stream=None)
|
||||
pool_stream = SimpleNamespace(_stream=SimpleNamespace(_connection=h11_conn), _connection=None)
|
||||
return SimpleNamespace(close=_no_close,
|
||||
stream=SimpleNamespace(_stream=SimpleNamespace(_httpcore_stream=pool_stream)))
|
||||
|
||||
|
||||
def _assert_shut_down(reader, writer):
|
||||
writer.settimeout(5)
|
||||
assert writer.recv(1) == b"", "the in-flight stream's socket was not shut down"
|
||||
|
||||
|
||||
def test_interrupt_abort_shuts_down_the_attempts_own_socket():
|
||||
reader, writer = _socket.socketpair()
|
||||
try:
|
||||
call = _call(_agent())
|
||||
call._attempt_stream_response = _response_over(reader)
|
||||
call.worker = None
|
||||
call._monitor_interrupted = {"yes": False} # set by the poll loop in production
|
||||
call._abort_for_interrupt(stale_elapsed=1.0)
|
||||
assert call._request_cancelled["value"] is True
|
||||
_assert_shut_down(reader, writer)
|
||||
finally:
|
||||
reader.close()
|
||||
writer.close()
|
||||
|
||||
|
||||
def test_stream_created_after_cancel_is_re_aborted():
|
||||
reader, writer = _socket.socketpair()
|
||||
try:
|
||||
call = _call(_agent())
|
||||
call._request_cancelled["value"] = True # abort fired while create() was still connecting
|
||||
raw_stream = SimpleNamespace(response=_response_over(reader))
|
||||
call._chat_stream_created(raw_stream)
|
||||
_assert_shut_down(reader, writer)
|
||||
finally:
|
||||
reader.close()
|
||||
writer.close()
|
||||
Reference in New Issue
Block a user