Skip to content

Commit b502d71

Browse files
authored
Merge pull request #70283 from bretep/fix/68660-3006x
Drop abandoned requests when draining the ZeroMQ send queue (#68660)
2 parents 707958a + 22fb248 commit b502d71

3 files changed

Lines changed: 92 additions & 0 deletions

File tree

‎changelog/68660.fixed.md‎

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,4 @@
1+
Drop abandoned requests when draining the ZeroMQ send queue in
2+
``AsyncReqMessageClient``. A request whose caller had already timed out stayed
3+
in ``self._queue`` holding its serialized payload until the drain loop reached
4+
it, which under sustained load it never did, growing the queue without bound.

‎salt/transport/zeromq.py‎

Lines changed: 20 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1001,6 +1001,26 @@ def _send_recv(self, socket, _TimeoutError=salt.ext.tornado.gen.TimeoutError):
10011001
send_recv_running = False
10021002
break
10031003

1004+
if future.done():
1005+
# The caller already abandoned this request: send() arms
1006+
# _timeout_message(), which completes the future when the
1007+
# caller's timeout expires. The reply can no longer be
1008+
# delivered to anyone, so drop the request instead of
1009+
# spending a round trip on it.
1010+
#
1011+
# Without this, self._queue grows without bound. Before
1012+
# 145a06e in-flight requests were tracked in
1013+
# self._send_future_map, a dict, so a timed-out entry was
1014+
# removed by key. A Queue has no removal-from-middle and
1015+
# self._queue has no maxsize, so an abandoned entry stays
1016+
# queued -- pinning its serialized payload -- until the
1017+
# drain loop reaches it. A REQ socket permits one
1018+
# request/reply in flight and _send_recv() restarts on
1019+
# every reconnect, so under sustained load the enqueue rate
1020+
# outruns the drain rate and it never does.
1021+
log.trace("Dropping request whose caller already timed out")
1022+
continue
1023+
10041024
try:
10051025
yield socket.send(message)
10061026
except zmq.eventloop.future.CancelledError as exc:

‎tests/pytests/unit/transport/test_zeromq.py‎

Lines changed: 68 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1825,6 +1825,16 @@ async def test_client_send_recv_on_cancelled_error(minion_opts):
18251825
client.socket = AsyncMock()
18261826
client.socket.poll.side_effect = zmq.eventloop.future.CancelledError
18271827
client._queue.put_nowait((mock_future, {"meh": "bah"}))
1828+
# The future is already done, so _send_recv drops it (see #68660) and
1829+
# returns to the queue. Queue a shutdown sentinel so the loop exits
1830+
# rather than falling through to the idle poll branch, which would
1831+
# call .result() on the AsyncMock's coroutine.
1832+
client._queue.put_nowait(
1833+
(
1834+
salt.ext.tornado.concurrent.Future(),
1835+
salt.transport.zeromq._REQ_QUEUE_SHUTDOWN,
1836+
)
1837+
)
18281838
await client._send_recv(client.socket)
18291839
mock_future.set_exception.assert_not_called()
18301840
finally:
@@ -1856,6 +1866,17 @@ async def test_client_send_recv_no_double_set_exception_after_timeout(minion_opt
18561866
client.socket = AsyncMock()
18571867
client.socket.send.side_effect = zmq.ZMQError(zmq.ETERM)
18581868
client._queue.put_nowait((future, {"meh": "bah"}))
1869+
# Since #68660 a request whose future is already done is dropped before
1870+
# it reaches socket.send, so this repro no longer exercises the send
1871+
# failure path -- the double-set it guarded against is now structurally
1872+
# unreachable here. The invariant still asserted below is that the
1873+
# original timeout exception survives untouched.
1874+
client._queue.put_nowait(
1875+
(
1876+
salt.ext.tornado.concurrent.Future(),
1877+
salt.transport.zeromq._REQ_QUEUE_SHUTDOWN,
1878+
)
1879+
)
18591880
# Before the fix this raises TypeError from tornado's _set_done.
18601881
await client._send_recv(client.socket)
18611882
# The timeout exception must be preserved, not overwritten.
@@ -1864,6 +1885,53 @@ async def test_client_send_recv_no_double_set_exception_after_timeout(minion_opt
18641885
client.close()
18651886

18661887

1888+
async def test_client_send_recv_drops_abandoned_request(minion_opts):
1889+
"""
1890+
Regression test for #68660.
1891+
1892+
``send()`` enqueues ``(future, message)`` and arms ``_timeout_message``.
1893+
When the caller's timeout fires the future is completed, but its queue
1894+
entry remains and keeps pinning the serialized payload until the drain
1895+
loop reaches it. ``self._queue`` has no maxsize and a REQ socket permits
1896+
one request/reply in flight, so under sustained load the enqueue rate
1897+
outruns the drain rate and the queue grows without bound.
1898+
1899+
``_send_recv`` must drop a request whose future is already done rather
1900+
than spend a round trip on a reply nobody can receive.
1901+
"""
1902+
client = salt.transport.zeromq.AsyncReqMessageClient(
1903+
minion_opts, "tcp://127.0.0.1:4506"
1904+
)
1905+
1906+
abandoned = salt.ext.tornado.concurrent.Future()
1907+
# Exactly what _timeout_message does when the caller's timeout expires.
1908+
client._timeout_message(abandoned)
1909+
assert abandoned.done()
1910+
1911+
# Keep our own reference: without the fix ``_send_recv`` takes its error
1912+
# path and ``_reconnect()`` swaps ``client.socket`` for a real socket, so
1913+
# asserting against ``client.socket`` afterwards would inspect the wrong
1914+
# object and fail for the wrong reason.
1915+
sock = AsyncMock()
1916+
try:
1917+
client.socket = sock
1918+
client._queue.put_nowait((abandoned, {"meh": "bah"}))
1919+
# Sentinel stops the drain loop after the abandoned entry is handled.
1920+
client._queue.put_nowait(
1921+
(
1922+
salt.ext.tornado.concurrent.Future(),
1923+
salt.transport.zeromq._REQ_QUEUE_SHUTDOWN,
1924+
)
1925+
)
1926+
await client._send_recv(sock)
1927+
# The abandoned payload must never reach the wire.
1928+
sock.send.assert_not_called()
1929+
# And its timeout exception must be left intact.
1930+
assert isinstance(abandoned.exception(), salt.exceptions.SaltReqTimeoutError)
1931+
finally:
1932+
client.close()
1933+
1934+
18671935
def test_async_req_message_client_close_never_connected(minion_opts):
18681936
"""
18691937
close() must not hang when connect() was never called (#68637).

0 commit comments

Comments
 (0)