celery-kombu-2492
Redis workers can become “catatonic” after a broker connection is lost and re-established: the worker remains alive and responds to health checks, but stops consuming tasks, leaving them queued in Redis. This occurs when disconnect cleanup from an old connection races with registration of the replacement connection. The transport registers a polling callback on the event loop’s per-tick hook when it registers with the loop. Disconnect handling must leave `on_poll_start` registered on `loop.on_tick`; a connection loss must never unregister it, whether or not any file descriptors are active. Polling must therefore resume when the replacement connection is registered and the worker must continue consuming tasks.
Disconnect cleanup, including `transport.cycle._on_connection_disconnect`, must remove the disconnected socket from the event loop and remove its file descriptor from the transport’s channel mapping, `cycle._fd_to_chan`, whose values are `(channel, type)` pairs. This must work for socket objects, raw file-descriptor values, already-closed sockets whose descriptor cannot be retrieved, missing sockets, and connections involving a subclient, without raising errors or preventing the event loop from polling again.
Whenever the connection has a socket object or raw integer descriptor, remove that object from the loop exactly once, passing it as-is. Prune `_fd_to_chan` only when a valid descriptor is available: use `fileno()` when it returns 0 or greater, or use the integer itself for a raw descriptor. If `fileno()` returns `-1` or raises, leave the mapping untouched. Entries must be removed whether or not the loop’s file-descriptor set is currently empty, and cleanup must tolerate entries that are already absent.
When the connection has no socket object, do not remove anything from the loop. Instead, remove stale mapping entries associated with that connection, including entries whose channel is connected through either its client or subclient, while leaving entries belonging to other connections unchanged. Disconnect cleanup must never raise, and `on_poll_start` must remain registered on `loop.on_tick` after every disconnect.
How the tests reach the cleanup: they replace `transport.cycle` with a `Mock` (spec'd to `fds`, `_fd_to_chan`, `on_poll_init`, `on_poll_start`, `maybe_restore_messages`, `maybe_check_subclient_health` and `_on_connection_disconnect`) BEFORE calling `redis.Transport.register_with_event_loop(transport, conn, loop)`, then call `transport.cycle._on_connection_disconnect(connection)` and assert on `loop.remove`, `cycle._fd_to_chan` and `loop.on_tick`. So registering with the event loop must assign the disconnect cleanup callable to the poller's `_on_connection_disconnect` attribute, and that callable must act on the `loop` object that was passed to `register_with_event_loop`; cleanup logic that lives only as a method of the real poller class is never executed.
Hidden tests · 5 fail-to-pass, 145 pass-to-passrun after the agent submits, in a clean verifier
Test patch · 260 lines
diff --git a/t/unit/transport/test_redis.py b/t/unit/transport/test_redis.py
index 11cfa4b9a4..16c3a183ea 100644
--- a/t/unit/transport/test_redis.py
+++ b/t/unit/transport/test_redis.py
@@ -1203,11 +1203,19 @@ def test_register_with_event_loop(self):
@pytest.mark.parametrize('fds', [{12: 'LISTEN', 13: 'BRPOP'}, {}])
def test_register_with_event_loop__on_disconnect__loop_cleanup(self, fds):
- """Ensure event loop polling stops on disconnect (if started)."""
+ """Ensure on_poll_start stays in on_tick after disconnect.
+
+ on_poll_start is idempotent (no-op when there are no fds), so it
+ must NOT be removed on disconnect. Removing it caused a race
+ condition where a late-firing _on_disconnect from a stale channel
+ would remove the callback just registered by a new channel,
+ leaving the worker unable to consume tasks after reconnection.
+ """
transport = self.connection.transport
self.connection._sock = None
transport.cycle = Mock(name='cycle')
transport.cycle.fds = fds
+ transport.cycle._fd_to_chan = {}
conn = Mock(name='conn')
conn.client = Mock(name='client', transport_options={})
loop = Mock(name='loop')
@@ -1215,11 +1223,229 @@ def test_register_with_event_loop__on_disconnect__loop_cleanup(self, fds):
redis.Transport.register_with_event_loop(transport, conn, loop)
assert len(loop.on_tick) == 1
transport.cycle._on_connection_disconnect(self.connection)
- if fds:
- assert len(loop.on_tick) == 0
- else:
- # on_tick shouldn't be cleared when polling hasn't started
- assert len(loop.on_tick) == 1
+ # on_poll_start must remain registered regardless of fds state
+ assert len(loop.on_tick) == 1
+
+ def test_register_with_event_loop__on_disconnect__removes_sock(self):
+ """Ensure _on_disconnect removes the socket from the event loop.
+
+ When connection._sock is set, _on_disconnect must call loop.remove()
+ on it and prune the fd from cycle._fd_to_chan when present.
+ """
+ transport = self.connection.transport
+ mock_sock = Mock(name='sock')
+ mock_sock.fileno.return_value = 42
+ self.connection._sock = mock_sock
+ transport.cycle = Mock(name='cycle', spec=['fds', '_fd_to_chan',
+ 'on_poll_init',
+ 'on_poll_start',
+ 'maybe_restore_messages',
+ 'maybe_check_subclient_health',
+ '_on_connection_disconnect'])
+ transport.cycle.fds = {}
+ transport.cycle._fd_to_chan = {42: Mock(name='chan')}
+ conn = Mock(name='conn')
+ conn.client = Mock(name='client', transport_options={})
+ loop = Mock(name='loop')
+ loop.on_tick = set()
+ redis.Transport.register_with_event_loop(transport, conn, loop)
+ transport.cycle._on_connection_disconnect(self.connection)
+ loop.remove.assert_called_once_with(mock_sock)
+ # fd must be pruned from _fd_to_chan
+ assert 42 not in transport.cycle._fd_to_chan
+ # on_poll_start must still be registered after disconnect
+ assert len(loop.on_tick) == 1
+
+ def test_register_with_event_loop__on_disconnect__sock_no_fileno(self):
+ """Ensure _on_disconnect handles sockets without a fileno() method.
+
+ When connection._sock has no fileno() (e.g. a raw fd integer),
+ the fd itself is used as the key to prune from cycle._fd_to_chan.
+ """
+ transport = self.connection.transport
+ # Use a plain integer as the "sock" (no fileno attribute)
+ self.connection._sock = 99
+ transport.cycle = Mock(name='cycle', spec=['fds', '_fd_to_chan',
+ 'on_poll_init',
+ 'on_poll_start',
+ 'maybe_restore_messages',
+ 'maybe_check_subclient_health',
+ '_on_connection_disconnect'])
+ transport.cycle.fds = {}
+ transport.cycle._fd_to_chan = {99: Mock(name='chan')}
+ conn = Mock(name='conn')
+ conn.client = Mock(name='client', transport_options={})
+ loop = Mock(name='loop')
+ loop.on_tick = set()
+ redis.Transport.register_with_event_loop(transport, conn, loop)
+ transport.cycle._on_connection_disconnect(self.connection)
+ loop.remove.assert_called_once_with(99)
+ assert 99 not in transport.cycle._fd_to_chan
+ assert len(loop.on_tick) == 1
+
+ def test_register_with_event_loop__on_disconnect__fileno_oserror(self):
+ """Ensure _on_disconnect handles OSError from fileno() gracefully.
+
+ When the socket is already closed, fileno() raises OSError.
+ _on_disconnect should swallow it and skip _fd_to_chan pruning.
+ """
+ transport = self.connection.transport
+ mock_sock = Mock(name='sock')
+ mock_sock.fileno.side_effect = OSError('Bad file descriptor')
+ self.connection._sock = mock_sock
+ transport.cycle = Mock(name='cycle', spec=['fds', '_fd_to_chan',
+ 'on_poll_init',
+ 'on_poll_start',
+ 'maybe_restore_messages',
+ 'maybe_check_subclient_health',
+ '_on_connection_disconnect'])
+ transport.cycle.fds = {}
+ transport.cycle._fd_to_chan = {}
+ conn = Mock(name='conn')
+ conn.client = Mock(name='client', transport_options={})
+ loop = Mock(name='loop')
+ loop.on_tick = set()
+ redis.Transport.register_with_event_loop(transport, conn, loop)
+ # Must not raise even though fileno() raises OSError
+ transport.cycle._on_connection_disconnect(self.connection)
+ loop.remove.assert_called_once_with(mock_sock)
+ assert len(loop.on_tick) == 1
+
+ def test_register_with_event_loop__on_disconnect__fd_not_in_map(self):
+ """Ensure _on_disconnect handles missing fd in _fd_to_chan gracefully.
+
+ If the fd is not tracked (already removed), the KeyError must be
+ swallowed silently.
+ """
+ transport = self.connection.transport
+ mock_sock = Mock(name='sock')
+ mock_sock.fileno.return_value = 55
+ self.connection._sock = mock_sock
+ transport.cycle = Mock(name='cycle', spec=['fds', '_fd_to_chan',
+ 'on_poll_init',
+ 'on_poll_start',
+ 'maybe_restore_messages',
+ 'maybe_check_subclient_health',
+ '_on_connection_disconnect'])
+ transport.cycle.fds = {}
+ # fd 55 is NOT in _fd_to_chan — KeyError must be silently ignored
+ transport.cycle._fd_to_chan = {}
+ conn = Mock(name='conn')
+ conn.client = Mock(name='client', transport_options={})
+ loop = Mock(name='loop')
+ loop.on_tick = set()
+ redis.Transport.register_with_event_loop(transport, conn, loop)
+ # Must not raise
+ transport.cycle._on_connection_disconnect(self.connection)
+ loop.remove.assert_called_once_with(mock_sock)
+ assert len(loop.on_tick) == 1
+
+ def test_register_with_event_loop__on_disconnect__fileno_negative(self):
+ """Ensure _
… [5704 more characters]Reference fix · 1 file, +53 −6the upstream merge, used only for grading calibration
The agent could not see this: the repository holds one commit and the sandbox has no network. Leak audit.
kombu/transport/redis.py
diff --git a/kombu/transport/redis.py b/kombu/transport/redis.py
index 68153ef827..5ebd4d3085 100644
--- a/kombu/transport/redis.py
+++ b/kombu/transport/redis.py
@@ -1447,14 +1447,61 @@ def register_with_event_loop(self, connection, loop):
def _on_disconnect(connection):
if connection._sock:
loop.remove(connection._sock)
-
- # must have started polling or this will break reconnection
- if cycle.fds:
- # stop polling in the event loop
+ # Prune the disconnected file descriptor from cycle._fd_to_chan
+ # so that the next on_poll_start tick does not re-register a
+ # stale/disconnected socket. fileno() returns -1 on a socket
+ # that has been closed (but not yet garbage-collected), so we
+ # only prune when we get a valid (>= 0) file descriptor.
+ sock = connection._sock
+ fd = None
try:
- loop.on_tick.remove(on_poll_start)
- except KeyError:
+ if hasattr(sock, "fileno"):
+ raw_fd = sock.fileno()
+ # fileno() returns -1 for a closed-but-not-GC'd socket;
+ # in that case there is no valid fd to prune.
+ if raw_fd >= 0:
+ fd = raw_fd
+ else:
+ # Plain integer file descriptor (no fileno() method).
+ fd = sock
+ except OSError:
+ # Socket already closed at OS level; nothing to prune.
pass
+ if fd is not None:
+ try:
+ del cycle._fd_to_chan[fd]
+ except KeyError:
+ # fd was never tracked or already pruned — safe to ignore.
+ pass
+ else:
+ # In async Redis mode, Connection.disconnect() may have already
+ # cleared connection._sock (set to None) before invoking this
+ # callback. In that case we can no longer derive the fd from the
+ # socket itself, so we conservatively scan cycle._fd_to_chan for
+ # channels that are backed by this connection and prune them.
+ stale_fds = []
+ for fd, (chan, _type) in list(cycle._fd_to_chan.items()):
+ client = getattr(chan, "client", None)
+ subclient = getattr(chan, "subclient", None)
+ client_conn = getattr(client, "connection", None)
+ subclient_conn = getattr(subclient, "connection", None)
+ if client_conn is connection or subclient_conn is connection:
+ stale_fds.append(fd)
+ for fd in stale_fds:
+ try:
+ del cycle._fd_to_chan[fd]
+ except KeyError:
+ # fd was never tracked or already pruned — safe to ignore.
+ pass
+ # Note: we intentionally do NOT remove on_poll_start from
+ # loop.on_tick here. on_poll_start is idempotent — when there
+ # are no active file descriptors it simply does nothing.
+ # Removing it caused a race condition where a late-firing
+ # _on_disconnect from a stale channel would remove the
+ # on_poll_start callback that a newly-reconnected channel had
+ # just registered, leaving the worker alive but unable to
+ # consume any tasks ("catatonic worker" after broker restart).
+ # See: https://github.com/celery/celery/issues/8030
cycle._on_connection_disconnect = _on_disconnect
def on_poll_start():
GPT-5.6 Luna · attempt 3 · solved · 8 steps · $0.011
Submitted diff · 130 lines
diff --git a/kombu/transport/redis.py b/kombu/transport/redis.py
index 68153ef..fd62fd7 100644
--- a/kombu/transport/redis.py
+++ b/kombu/transport/redis.py
@@ -469,6 +469,82 @@ class QoS(virtual.QoS):
return self.channel.visibility_timeout
+def _connection_socket(connection):
+ """Return a connection's socket without allowing stale cleanup to fail."""
+ try:
+ return connection._sock
+ except Exception:
+ return None
+
+
+def _connection_matches_channel(channel, connection):
+ """Whether either Redis client owned by *channel* uses *connection*."""
+ for name in ('client', 'subclient'):
+ try:
+ client = getattr(channel, name)
+ if getattr(client, 'connection', None) is connection:
+ return True
+ except Exception:
+ continue
+ return False
+
+
+def _remove_connection_fds(cycle, connection, sock=None, loop=None):
+ """Remove a disconnected connection from the loop and fd map.
+
+ This is kept outside the poller class because the transport installs this
+ function as a callback on the current event loop. In particular, the
+ callback must continue to work when ``cycle`` is replaced by a mock or by
+ a newly-created poller during reconnect.
+ """
+ if sock is None:
+ sock = _connection_socket(connection)
+
+ # A socket (including an already closed socket) must be handed to the
+ # event loop exactly once. Do not use the truth value: fd 0 is valid.
+ if sock is not None and loop is not None:
+ try:
+ loop.remove(sock)
+ except Exception:
+ pass
+
+ try:
+ fd_map = cycle._fd_to_chan
+ except Exception:
+ return
+
+ fd = None
+ if isinstance(sock, int):
+ fd = sock if sock >= 0 else None
+ elif sock is not None:
+ try:
+ candidate = sock.fileno()
+ if candidate >= 0:
+ fd = candidate
+ except Exception:
+ pass
+
+ try:
+ if fd is not None:
+ # pop tolerates a disconnect after the entry was already removed.
+ fd_map.pop(fd, None)
+ return
+
+ # Without a socket there is no reliable fd to derive. Remove only
+ # entries belonging to this connection, including LISTEN entries
+ # registered through a channel's subclient.
+ for entry_fd, entry in list(fd_map.items()):
+ try:
+ channel = entry[0]
+ except (IndexError, TypeError):
+ continue
+ if _connection_matches_channel(channel, connection):
+ fd_map.pop(entry_fd, None)
+ except Exception:
+ # Disconnect cleanup must never prevent the hub from ticking again.
+ pass
+
+
class MultiChannelPoller:
"""Async I/O poller for Redis transport."""
@@ -509,10 +585,16 @@ class MultiChannelPoller:
self._channels.discard(channel)
def _on_connection_disconnect(self, connection):
- try:
- self.poller.unregister(connection._sock)
- except (AttributeError, TypeError):
- pass
+ sock = _connection_socket(connection)
+ if sock is not None:
+ try:
+ self.poller.unregister(sock)
+ except Exception:
+ # A disconnect can race with registration (or with another
+ # disconnect). The socket is no longer useful to this
+ # poller in either case.
+ pass
+ _remove_connection_fds(self, connection, sock)
def _register(self, channel, client, type):
if (channel, client, type) in self._chan_to_sock:
@@ -1445,16 +1527,11 @@ class Transport(virtual.Transport):
on_readable = self.on_readable
def _on_disconnect(connection):
- if connection._sock:
- loop.remove(connection._sock)
-
- # must have started polling or this will break reconnection
- if cycle.fds:
- # stop polling in the event loop
- try:
- loop.on_tick.remove(on_poll_start)
- except KeyError:
- pass
+ # Do not remove on_poll_start here. It is deliberately a
+ # per-tick callback: removing it while an old connection is
+ # being cleaned up can race with registration of its replacement
+ # and leave the worker alive but no longer consuming messages.
+ _remove_connection_fds(cycle, connection, loop=loop)
cycle._on_connection_disconnect = _on_disconnect
def on_poll_start():I’ll trace the Redis transport event-loop registration and disconnect cleanup, inspect relevant tests, reproduce the race/cleanup behavior, then make a focused source-only fix and run the affected test suite.
$ pwd; ls; grep -R "def register_with_event_loop\|_on_connection_disconnect\|on_poll_start" -n . --exclude-dir=.git | head -200
output · 39 lines
/repo AUTHORS Changelog.rst FAQ INSTALL LICENSE MANIFEST.in Makefile README.rst SECURITY.md THANKS TODO
27 more lines
conftest.py docs examples kombu kombu.egg-info pyproject.toml requirements setup.cfg setup.py t tox.ini ./Changelog.rst:4674:- ``Transport.on_poll_start()`` ./Changelog.rst:4680:- ``Transport.on_poll_start(poller)`` ./kombu/connection.py:302: def register_with_event_loop(self, loop): ./kombu/transport/base.py:213: def register_with_event_loop(self, connection, loop): ./kombu/transport/librabbitmq.py:163: def register_with_event_loop(self, connection, loop): ./kombu/transport/pyamqp.py:222: def register_with_event_loop(self, connection, loop): ./kombu/transport/qpid.py:1542: def register_with_event_loop(self, connection, loop): ./kombu/transport/redis.py:511: def _on_connection_disconnect(self, connection): ./kombu/transport/redis.py:556: def on_poll_start(self): …[1494 chars omitted]… est_redis.py:1243: transport.cycle.on_poll_start.assert_called_with() ./t/unit/transport/test_redis.py:1718: def test_on_poll_start(self): ./t/unit/transport/test_redis.py:1721: p.on_poll_start() ./t/unit/transport/test_redis.py:1729: p.on_poll_start() ./t/unit/transport/test_redis.py:1735: p.on_poll_start() ./t/unit/transport/test_redis.py:1741: p.on_poll_start()
$ grep -R "fd_to_chan\|remove(.*sock\|cycle.fds\|class MultiChannelPoller" -n kombu t/unit/transport/test_redis.py
output · 20 lines
kombu/transport/redis.py:472:class MultiChannelPoller:
kombu/transport/redis.py:487: self._fd_to_chan = {}
kombu/transport/redis.py:502: self._fd_to_chan.clear()
kombu/transport/redis.py:523: self._fd_to_chan[sock.fileno()] = (channel, type)
kombu/transport/redis.py:604: chan, type = self._fd_to_chan[fileno]
kombu/transport/redis.py:612: chan, type = self._fd_to_chan[fileno]
kombu/transport/redis.py:647: return self._fd_to_chan
kombu/transport/redis.py:1449: loop.remove(connection._sock)
kombu/transport/redis.py:1452: if cycle.fds:
kombu/transport/redis.py:1462: [add_reader(fd, on_readable, fd) for fd in cycle.fds]
t/unit/transport/test_redis.py:1184: transport.cycle.fds = {12: 'LISTEN', 13: 'BRPOP'}
t/unit/transport/test_redis.py:1210: transport.cycle.fds = fds8 more lines
t/unit/transport/test_redis …[423 chars omitted]… t/unit/transport/test_redis.py:1905: p._fd_to_chan = Mock() t/unit/transport/test_redis.py:1913: p._fd_to_chan.clear.assert_called_with() t/unit/transport/test_redis.py:1926: assert p._fd_to_chan[10] == (channel, type) t/unit/transport/test_redis.py:2043: p._fd_to_chan[1] = (channel, 'BRPOP') t/unit/transport/test_redis.py:2053: p._fd_to_chan[1] = (channel, 'BRPOP')
$ sed -n '590,670p' kombu/transport/redis.py; sed -n '1835,2070p' t/unit/transport/test_redis.py; grep -n "subclient\|_chan_to_sock" t/unit/transport/test_redis.py | head -40; grep -n "class Channel\|def connection" kombu/transport/redis.py | head
output · 31 lines
# only if subclient property is cached
client = channel.__dict__.get('subclient')
if client is not None \
and callable(getattr(client, 'check_health', None)):
try:
client.check_health()
except channel.connection_errors:
logger.debug(
'maybe_check_subclient_health: connection error, '
'will retry on next cycle', exc_info=True
)
return19 more lines
def on_readable(self, fileno):
chan, type = self._fd_to_chan[fileno]
if chan.qos.can_consume():
chan.handlers[type]()
def handle_event(self, fileno, event):
if event & READ:
return self.on_readable(fileno), self
elif event & ERR:
chan, type = self._fd_to_chan[fileno]
…[8629 chars omitted]…
: # Ensure 'subclient' is not in __dict__ (not yet accessed/cached)
1825: channel.__dict__.pop('subclient', None)
1828: p.maybe_check_subclient_health() # must not raise
1830: def test_maybe_check_subclient_health_swallows_connection_error(self):
1844: client = Mock(name='subclient')
1846: channel.__dict__['subclient'] = client
650:class Channel(virtual.Channel):$ python - <<'PY'
output · 23 lines
..............................................................F......... [ 50%]
....................................................................... [100%]
=================================== FAILURES ===================================
_ test_Channel.test_register_with_event_loop__on_disconnect__loop_cleanup[fds0] _
self = <t.unit.transport.test_redis.test_Channel object at 0x79ad07c8daf0>
fds = {12: 'LISTEN', 13: 'BRPOP'}
@pytest.mark.parametrize('fds', [{12: 'LISTEN', 13: 'BRPOP'}, {}])
def test_register_with_event_loop__on_disconnect__loop_cleanup(self, fds):
"""Ensure event loop polling stops on disconnect (if started)."""
transport = self.connection.transport11 more lines
self.connection._sock = None
transport.cycle = Mock(name='cycle')
transport.cycle.fds = fds
conn = Mock(name='conn')
conn.client = Mock(name='client', tr
…[489 chars omitted]…
+ where {<function Transport.register_with_event_loop.<locals>.on_poll_start at 0x79ad07914040>} = <Mock name='loop' id='133784062960576'>.on_tick
t/unit/transport/test_redis.py:1219: AssertionError
=========================== short test summary info ============================
FAILED t/unit/transport/test_redis.py::test_Channel::test_register_with_event_loop__on_disconnect__loop_cleanup[fds0]$ python - <<'PY'
output · 34 lines
reproducer passed
........................................................................ [ 51%]
..................................................................... [100%]
141 passed, 2 deselected in 5.20s
diff --git a/kombu/transport/redis.py b/kombu/transport/redis.py
index 68153ef..fd62fd7 100644
--- a/kombu/transport/redis.py
+++ b/kombu/transport/redis.py
@@ -469,6 +469,82 @@ class QoS(virtual.QoS):
return self.channel.visibility_timeout
22 more lines
+def _connection_socket(connection):
+ """Return a connection's socket without allowing stale cleanup to fail."""
+ try:
+ return connection._sock
+ except Exception:
+ return None
+
+
+def _connection_matches_channel(channel, connection):
+ """Whether either Redis client owned by *channel* uses *connection*."""
+ for name in ('client', 'subclient'):
+ try:
+ client = getattr(chann
…[3582 chars omitted]…
It is deliberately a
+ # per-tick callback: removing it while an old connection is
+ # being cleaned up can race with registration of its replacement
+ # and leave the worker alive but no longer consuming messages.
+ _remove_connection_fds(cycle, connection, loop=loop)
cycle._on_connection_disconnect = _on_disconnect
def on_poll_start():$ python -m pytest -q t/unit/transport/test_redis.py --disable-warnings --maxfail=1; python -m pytest -q t/unit/transport --disable-warnings --maxfail=1; git diff --check; git status --short
output · 28 lines
..............................................................F
=================================== FAILURES ===================================
_ test_Channel.test_register_with_event_loop__on_disconnect__loop_cleanup[fds0] _
self = <t.unit.transport.test_redis.test_Channel object at 0x714c57bca330>
fds = {12: 'LISTEN', 13: 'BRPOP'}
@pytest.mark.parametrize('fds', [{12: 'LISTEN', 13: 'BRPOP'}, {}])
def test_register_with_event_loop__on_disconnect__loop_cleanup(self, fds):
"""Ensure event loop polling stops on disconnect (if started)."""
transport = self.connection.transport
self.connection._sock = None16 more lines
transport.cycle = Mock(name='cycle')
transport.cycle.fds = fds
conn = Mock(name='conn')
conn.client = Mock(name='client', transport_options={})
loop = Mock(name='loop')
loop.on_tick = set()
redis.
…[1782 chars omitted]…
:227: in <module>
from botocore.client import BaseClient, Config
E ModuleNotFoundError: No module named 'botocore'
=========================== short test summary info ============================
ERROR t/unit/transport/SQS - ModuleNotFoundError: No module named 'botocore'
!!!!!!!!!!!!!!!!!!!!!!!!!! stopping after 1 failures !!!!!!!!!!!!!!!!!!!!!!!!!!!
1 error in 0.72s
M kombu/transport/redis.py$ python - <<'PY'
output · 4 lines
fd edge cases passed association edge cases passed ... [100%] 3 passed, 140 deselected in 0.54s
$ done: {"notes": "Implemented Redis disconnect cleanup and event-loop polling fix in kombu/transport/redis.py. Added robust socket/fd cleanup, subclient association handling, preserved on_poll_start registration, and validated with targeted tests plus manual edge-case reproducers."}output · 1 lines
submitted.
GPT-5.6 Luna · attempt 4 · solved · 10 steps · $0.012
Submitted diff · 129 lines
diff --git a/kombu/transport/redis.py b/kombu/transport/redis.py
index 68153ef..1120ffc 100644
--- a/kombu/transport/redis.py
+++ b/kombu/transport/redis.py
@@ -469,6 +469,80 @@ class QoS(virtual.QoS):
return self.channel.visibility_timeout
+def _connection_socket(connection):
+ """Return a connection socket without letting cleanup fail."""
+ try:
+ return connection._sock
+ except Exception:
+ return None
+
+
+def _socket_fileno(sock):
+ if isinstance(sock, int):
+ return sock if sock >= 0 else None
+ try:
+ fd = sock.fileno()
+ except Exception:
+ return None
+ return fd if isinstance(fd, int) and fd >= 0 else None
+
+
+def _channel_uses_connection(channel, connection):
+ for attr in ('client', 'subclient'):
+ try:
+ client = getattr(channel, attr)
+ if getattr(client, 'connection', None) is connection:
+ return True
+ except Exception:
+ pass
+ return False
+
+
+def _remove_connection_mapping(cycle, connection, sock):
+ """Remove disconnected socket state without disturbing other clients."""
+ try:
+ chan_to_sock = cycle._chan_to_sock
+ for key, value in list(chan_to_sock.items()):
+ channel = key[0]
+ if value is sock or _channel_uses_connection(channel, connection):
+ chan_to_sock.pop(key, None)
+ except Exception:
+ pass
+
+ try:
+ mapping = cycle._fd_to_chan
+ except Exception:
+ return
+
+ if sock is not None:
+ fd = _socket_fileno(sock)
+ # A closed socket has no reliable fd. In particular, do not guess
+ # based on the channel: that could remove a newly registered socket.
+ if fd is not None:
+ try:
+ mapping.pop(fd, None)
+ except Exception:
+ pass
+ return
+
+ # Some redis-py versions clear _sock before notifying us. In that case
+ # identify entries through both clients, since LISTEN uses subclient.
+ try:
+ entries = list(mapping.items())
+ except Exception:
+ return
+ for fd, value in entries:
+ try:
+ channel = value[0]
+ except Exception:
+ continue
+ if _channel_uses_connection(channel, connection):
+ try:
+ mapping.pop(fd, None)
+ except Exception:
+ pass
+
+
class MultiChannelPoller:
"""Async I/O poller for Redis transport."""
@@ -509,10 +583,15 @@ class MultiChannelPoller:
self._channels.discard(channel)
def _on_connection_disconnect(self, connection):
- try:
- self.poller.unregister(connection._sock)
- except (AttributeError, TypeError):
- pass
+ sock = _connection_socket(connection)
+ if sock is not None:
+ try:
+ self.poller.unregister(sock)
+ except Exception:
+ # Disconnect is best-effort; the poller may already have
+ # discarded a closed socket.
+ pass
+ _remove_connection_mapping(self, connection, sock)
def _register(self, channel, client, type):
if (channel, client, type) in self._chan_to_sock:
@@ -1445,16 +1524,15 @@ class Transport(virtual.Transport):
on_readable = self.on_readable
def _on_disconnect(connection):
- if connection._sock:
- loop.remove(connection._sock)
-
- # must have started polling or this will break reconnection
- if cycle.fds:
- # stop polling in the event loop
+ # Do not remove on_poll_start. It is the hook which registers
+ # readers for the replacement connection after a reconnect.
+ sock = _connection_socket(connection)
+ if sock is not None:
try:
- loop.on_tick.remove(on_poll_start)
- except KeyError:
+ loop.remove(sock)
+ except Exception:
pass
+ _remove_connection_mapping(cycle, connection, sock)
cycle._on_connection_disconnect = _on_disconnect
def on_poll_start():I’ll trace the transport/event-loop registration and disconnect cleanup paths, reproduce the relevant behavior with the existing test setup if possible, then make a narrowly scoped fix. I’ll specifically verify Mock-replaced pollers, socket/raw-fd/no-socket cases, subclient mappings, and preservation of `loop.on_tick` before running affected tests and reviewing the diff.
$ ls && grep -R "def register_with_event_loop\|_on_connection_disconnect\|on_poll_start\|_fd_to_chan" -n . --exclude-dir=.git | head -200
output · 38 lines
AUTHORS Changelog.rst FAQ INSTALL LICENSE MANIFEST.in Makefile README.rst SECURITY.md THANKS TODO conftest.py
26 more lines
docs
examples
kombu
kombu.egg-info
pyproject.toml
requirements
setup.cfg
setup.py
t
tox.ini
./Changelog.rst:4674:- ``Transport.on_poll_start()``
./Changelog.rst:4680:- ``Transport.on_poll_start(poller)``
./kombu/connection.py:302: def register_with_event_loop(self, loop):
./kombu/transport/base.py:213: def register_with_event_loop(self, connection, loop):
./kombu/transport/librabbitmq.py:163: def register_with_event_loop(self, connection, loop):
./kombu/transport/pyamqp.py:222: def register_with_event_loop(self, connection, loop):
./kombu/transport/qpid.py:1542: def register_with_event_loop(self, connection, loop):
./kombu/transport/redis.py:487: self._fd_to_chan = {}
./kombu/transport/redis.py:502: self._fd_to_chan.clear()
./kombu/transport/redis.
…[2562 chars omitted]…
transport/test_redis.py:1905: p._fd_to_chan = Mock()
./t/unit/transport/test_redis.py:1913: p._fd_to_chan.clear.assert_called_with()
./t/unit/transport/test_redis.py:1926: assert p._fd_to_chan[10] == (channel, type)
./t/unit/transport/test_redis.py:2043: p._fd_to_chan[1] = (channel, 'BRPOP')
./t/unit/transport/test_redis.py:2053: p._fd_to_chan[1] = (channel, 'BRPOP')$ sed -n '1150,1265p' t/unit/transport/test_redis.py; sed -n '1820,2070p' t/unit/transport/test_redis.py
output · 42 lines
def test_get_client(self):
import redis as R
KombuRedis = redis.Channel._get_client(self.channel)
assert isinstance(KombuRedis(), R.StrictRedis)
Rv = getattr(R, 'VERSION', None)
try:
R.VERSION = (2, 4, 0)
with pytest.raises(VersionMismatch):
redis.Channel._get_client(self.channel)
finally:
if Rv is not None:30 more lines
R.VERSION = Rv
def test_get_prefixed_client(self):
from kombu.transport.redis import PrefixedStrictRedis
self.channel.global_keyprefix = "test_"
PrefixedRedis = redis.Channel._get_client(self.channel)
assert isinstance(PrefixedRedis(), PrefixedStrictRedis)
def test_get_response_error(self):
from redis.exceptions import ResponseError
assert redis.Channel._get_response_error(self.channel) is ResponseError
…[8629 chars omitted]…
with pytest.raises(redis.Empty):
p.get(Mock())
channel._poll_error.assert_called_with('BRPOP')
class test_Mutex:
def test_mutex(self, lock_id='xxx'):
client = Mock(name='client')
lock = client.lock.return_value = Mock(name='lock')
# Won
lock.acquire.return_value = True
held = False
with redis.Mutex(client, 'foo1', 100):$ sed -n '630,665p' kombu/transport/redis.py; sed -n '1680,1770p' t/unit/transport/test_redis.py; grep -n "disconnect" t/unit/transport/test_redis.py | head -30; sed -n '1460,1505p' kombu/transport/redis.py
output · 53 lines
return
# - no new data, so try to restore messages.
# - reset active redis commands.
self.maybe_restore_messages()
raise Empty()
finally:
self._in_protected_read = False
while self.after_read:
try:
fun = self.after_read.pop()
except KeyError:
break41 more lines
else:
fun()
@property
def fds(self):
return self._fd_to_chan
class Channel(virtual.Channel):
"""Redis Channel."""
QoS = QoS
_client = None
_subclient = None
_closing = False
supports_fanout = True
keyprefix_queue = '_kombu.binding.%s'
keyprefix_fanout = '/{db}.'
sep = '\x06\x16'
_in_poll = False
_in_listen = False
_fanout_queues = {}
ack_emulation = True
d
…[5972 chars omitted]…
ass SentinelManagedSSLConnection(
sentinel.SentinelManagedConnection,
redis.SSLConnection):
"""Connect to a Redis server using Sentinel + TLS.
Use Sentinel to identify which Redis server is the current master
to connect to and when connecting to the Master server, use an
SSL Connection.
"""
pass
class SentinelChannel(Channel):$ python - <<'PY'
output · 23 lines
..............................................................F......... [ 50%]
....................................................................... [100%]
=================================== FAILURES ===================================
_ test_Channel.test_register_with_event_loop__on_disconnect__loop_cleanup[fds0] _
self = <t.unit.transport.test_redis.test_Channel object at 0x7f43c1dbcef0>
fds = {12: 'LISTEN', 13: 'BRPOP'}
@pytest.mark.parametrize('fds', [{12: 'LISTEN', 13: 'BRPOP'}, {}])
def test_register_with_event_loop__on_disconnect__loop_cleanup(self, fds):
"""Ensure event loop polling stops on disconnect (if started)."""
transport = self.connection.transport11 more lines
self.connection._sock = None
transport.cycle = Mock(name='cycle')
transport.cycle.fds = fds
conn = Mock(name='conn')
conn.client = Mock(name='client', tr
…[489 chars omitted]…
+ where {<function Transport.register_with_event_loop.<locals>.on_poll_start at 0x7f43c1a2b920>} = <Mock name='loop' id='139928985297664'>.on_tick
t/unit/transport/test_redis.py:1219: AssertionError
=========================== short test summary info ============================
FAILED t/unit/transport/test_redis.py::test_Channel::test_register_with_event_loop__on_disconnect__loop_cleanup[fds0]$ python - <<'PY'
output · 48 lines
return connection._sock
except Exception:
return None
def _socket_fileno(sock):
if isinstance(sock, int):
return sock if sock >= 0 else None
try:
fd = sock.fileno()
except Exception:
return None36 more lines
return fd if isinstance(fd, int) and fd >= 0 else None
def _channel_uses_connection(channel, connection):
for attr in ('client', 'subclient'):
try:
client = getattr(channel, attr)
if getattr(client, 'connection', None) is connection:
return True
except (AttributeError, TypeError):
pass
return False
def _remove_connection_mapping(cycle, connection, sock):
"""Remove the fd for a socket, or stale entries for a socketless client."""
try:
mapping = cycle._fd_to_chan
except Exception:
return
if sock is not None:
fd = _s
…[4977 chars omitted]…
on)
+ if sock is not None:
try:
- loop.on_tick.remove(on_poll_start)
- except KeyError:
+ loop.remove(sock)
+ except Exception:
pass
+ _remove_connection_mapping(cycle, connection, sock)
cycle._on_connection_disconnect = _on_disconnect
def on_poll_start():$ grep -n "_on_connection_disconnect" t/unit/transport/test_redis.py; pytest -q t/unit/transport/test_redis.py --disable-warnings --maxfail=2 -k 'not register_with_event_loop__on_disconnect__loop_cleanup'
output · 6 lines
427: def _on_connection_disconnect(self, connection): 486: def _on_connection_disconnect(self, connection): 1217: transport.cycle._on_connection_disconnect(self.connection) ........................................................................ [ 51%] ..................................................................... [100%] 141 passed, 2 deselected in 4.82s
$ python - <<'PY'
output · 28 lines
.F
=================================== FAILURES ===================================
_ test_Channel.test_register_with_event_loop__on_disconnect__loop_cleanup[fds0] _
self = <t.unit.transport.test_redis.test_Channel object at 0x77824aef7200>
fds = {12: 'LISTEN', 13: 'BRPOP'}
@pytest.mark.parametrize('fds', [{12: 'LISTEN', 13: 'BRPOP'}, {}])
def test_register_with_event_loop__on_disconnect__loop_cleanup(self, fds):
"""Ensure event loop polling stops on disconnect (if started)."""
transport = self.connection.transport
self.connection._sock = None16 more lines
transport.cycle = Mock(name='cycle')
transport.cycle.fds = fds
conn = Mock(name='conn')
conn.client = Mock(name='client', transport_options={})
loop = Mock(name='loop')
loop.on_tick = set()
redis.Transport.register_with_event_loop(transport, conn, loop)
…[457 chars omitted]…
='131401462056000'>.on_tick
t/unit/transport/test_redis.py:1219: AssertionError
=========================== short test summary info ============================
FAILED t/unit/transport/test_redis.py::test_Channel::test_register_with_event_loop__on_disconnect__loop_cleanup[fds0]
!!!!!!!!!!!!!!!!!!!!!!!!!! stopping after 1 failures !!!!!!!!!!!!!!!!!!!!!!!!!!!
1 failed, 1 passed, 114 deselected in 0.57s$ git diff --check && git status --short
output · 1 lines
M kombu/transport/redis.py
$ git diff -- kombu/transport/redis.py
output · 44 lines
diff --git a/kombu/transport/redis.py b/kombu/transport/redis.py
index 68153ef..1120ffc 100644
--- a/kombu/transport/redis.py
+++ b/kombu/transport/redis.py
@@ -469,6 +469,80 @@ class QoS(virtual.QoS):
return self.channel.visibility_timeout
+def _connection_socket(connection):
+ """Return a connection socket without letting cleanup fail."""
+ try:
+ return connection._sock32 more lines
+ except Exception:
+ return None
+
+
+def _socket_fileno(sock):
+ if isinstance(sock, int):
+ return sock if sock >= 0 else None
+ try:
+ fd = sock.fileno()
+ except Exception:
+ return None
+ return fd if isinstance(fd, int) and fd >= 0 else None
+
+
+def _channel_uses_connection(channel, connection):
+ for attr in ('client', 'subclient'):
+ try:
+ client = getattr(channel, attr)
+ if getattr(client, 'connection',
…[2966 chars omitted]…
on)
+ if sock is not None:
try:
- loop.on_tick.remove(on_poll_start)
- except KeyError:
+ loop.remove(sock)
+ except Exception:
pass
+ _remove_connection_mapping(cycle, connection, sock)
cycle._on_connection_disconnect = _on_disconnect
def on_poll_start():$ done: {"notes": "Implemented Redis disconnect cleanup and event-loop polling fix. Cleanup preserves on_poll_start, removes sockets/raw descriptors safely, prunes mappings including subclient associations, tolerates missing/closed sockets, and is wired onto the cycle callback used during event-loop registration. Affected Redis tests pass except the existing test that asserts the superseded behavior of removing on_tick polling."}output · 1 lines
submitted.
GPT-5.6 Luna · attempt 5 · not solved · 14 steps · $0.018
Submitted diff · 163 lines
diff --git a/kombu/transport/redis.py b/kombu/transport/redis.py
index 68153ef..95ec123 100644
--- a/kombu/transport/redis.py
+++ b/kombu/transport/redis.py
@@ -469,6 +469,113 @@ class QoS(virtual.QoS):
return self.channel.visibility_timeout
+def _connection_for_client(client):
+ try:
+ return client.connection
+ except (AttributeError, TypeError):
+ return None
+
+
+def _channel_uses_connection(channel, connection):
+ """Return whether a poll entry belongs to *connection*.
+
+ A channel can have two redis connections (the normal client and the
+ pub/sub subclient). Looking at both is important when cleaning up a
+ connection which failed before its socket was exposed by redis-py.
+ """
+ found = False
+ for name in ('client', 'subclient'):
+ try:
+ client = channel.__dict__.get(name)
+ except AttributeError:
+ client = None
+ if client is None:
+ continue
+ found = True
+ if _connection_for_client(client) is connection:
+ return True
+ return False if found else None
+
+
+def _remove_connection_poll_entries(cycle, connection, remove_from_poller=True):
+ """Remove poll registrations belonging to a disconnected connection."""
+ try:
+ sock = connection._sock
+ except (AttributeError, TypeError):
+ sock = None
+
+ # A socket may be false-y when it has been closed, so do not use a truth
+ # test here. In particular, the event loop must see the exact object
+ # supplied by redis-py (or the raw descriptor).
+ if sock is not None and remove_from_poller:
+ try:
+ cycle.poller.unregister(sock)
+ except Exception:
+ # Disconnect is cleanup and must never mask reconnect handling.
+ pass
+
+ fd = None
+ if sock is not None:
+ if isinstance(sock, numbers.Integral):
+ fd = int(sock)
+ else:
+ try:
+ candidate = sock.fileno()
+ if candidate >= 0:
+ fd = candidate
+ except Exception:
+ pass
+
+ if fd is not None:
+ try:
+ entry = cycle._fd_to_chan.get(fd)
+ except Exception:
+ entry = None
+ if entry is not None:
+ try:
+ belongs = _channel_uses_connection(entry[0], connection)
+ except (IndexError, KeyError, TypeError):
+ belongs = None
+ # Unknown ownership is stale in the normal poller, while an
+ # entry known to belong to another (replacement) connection
+ # must be retained.
+ if belongs is not False:
+ try:
+ del cycle._fd_to_chan[fd]
+ except (KeyError, TypeError, AttributeError):
+ pass
+
+ # Also clear the channel->socket index. Otherwise _client_registered()
+ # can incorrectly consider a replacement connection already registered.
+ try:
+ for key, registered_sock in list(cycle._chan_to_sock.items()):
+ channel, client, _type = key
+ client_connection = _connection_for_client(client)
+ if registered_sock is sock or client_connection is connection:
+ try:
+ del cycle._chan_to_sock[key]
+ except (KeyError, TypeError):
+ pass
+ except (AttributeError, TypeError):
+ pass
+
+ if sock is None:
+ try:
+ entries = list(cycle._fd_to_chan.items())
+ except (AttributeError, TypeError):
+ entries = ()
+ for fd, entry in entries:
+ try:
+ belongs = _channel_uses_connection(entry[0], connection)
+ except (IndexError, TypeError):
+ continue
+ if belongs is True:
+ try:
+ del cycle._fd_to_chan[fd]
+ except (KeyError, TypeError, AttributeError):
+ pass
+
+
class MultiChannelPoller:
"""Async I/O poller for Redis transport."""
@@ -510,8 +617,9 @@ class MultiChannelPoller:
def _on_connection_disconnect(self, connection):
try:
- self.poller.unregister(connection._sock)
- except (AttributeError, TypeError):
+ _remove_connection_poll_entries(self, connection)
+ except Exception:
+ # Disconnect cleanup must never prevent the next poll tick.
pass
def _register(self, channel, client, type):
@@ -1445,16 +1553,24 @@ class Transport(virtual.Transport):
on_readable = self.on_readable
def _on_disconnect(connection):
- if connection._sock:
- loop.remove(connection._sock)
-
- # must have started polling or this will break reconnection
- if cycle.fds:
- # stop polling in the event loop
+ # Keep the tick callback installed: it is what registers readers
+ # for the replacement connection after a reconnect.
+ try:
+ sock = connection._sock
+ except (AttributeError, TypeError):
+ sock = None
+ if sock is not None:
try:
- loop.on_tick.remove(on_poll_start)
- except KeyError:
+ loop.remove(sock)
+ except Exception:
pass
+ try:
+ _remove_connection_poll_entries(
+ cycle, connection, remove_from_poller=False
+ )
+ except Exception:
+ # Disconnect cleanup must never prevent the next poll tick.
+ pass
cycle._on_connection_disconnect = _on_disconnect
def on_poll_start():I’ll trace the transport event-loop registration and disconnect cleanup paths first, then reproduce the race/cleanup cases with a small script or focused tests. I’ll make the smallest source-only fix, run relevant tests, and inspect the final diff.
$ ls && rg "register_with_event_loop|_on_connection_disconnect|on_poll_start|_fd_to_chan" -n . | head -200
output · 23 lines
AUTHORS Changelog.rst FAQ INSTALL LICENSE MANIFEST.in Makefile README.rst SECURITY.md THANKS TODO conftest.py
11 more lines
docs examples kombu kombu.egg-info pyproject.toml requirements setup.cfg setup.py t tox.ini bash: line 6: rg: command not found
$ find kombu -iname '*redis*' -o -iname '*event*' | head -80
output · 2 lines
kombu/transport/redis.py kombu/utils/eventio.py
$ sed -n '1680,1770p' t/unit/transport/test_redis.py; sed -n '1830,2070p' t/unit/transport/test_redis.py
output · 39 lines
def test_get__Empty(self):
channel = self.connection.channel()
with pytest.raises(Empty):
channel._get('does-not-exist')
channel.close()
@pytest.mark.ensured_modules(*_redis_modules())
def test_get_client(self, module_exists):
# with module_exists(*_redis_modules()):
conn = Connection(transport=Transport)
chan = conn.channel()
assert chan.Client27 more lines
assert chan.ResponseError
assert conn.transport.connection_errors
assert conn.transport.channel_errors
def test_brpop_timeout_propagates_from_transport_options(self):
# Set either polling_interval or brpop_timeout to 2
conn = Connection("redis://localhost/0", transport_options={"polling_interval": 2})
# Avoid network I/O during Channel.__init__()
with patch("kombu.transport.redis.redis.Redis.ping", ret
…[8628 chars omitted]…
with pytest.raises(redis.Empty):
p.get(Mock())
channel._poll_error.assert_called_with('BRPOP')
class test_Mutex:
def test_mutex(self, lock_id='xxx'):
client = Mock(name='client')
lock = client.lock.return_value = Mock(name='lock')
# Won
lock.acquire.return_value = True
held = False
with redis.Mutex(client, 'foo1', 100):$ grep -n "_sock\|connection =\|subclient" kombu/transport/redis.py | head -100; sed -n '1470,1515p' kombu/transport/redis.py; grep -R "_on_connection_disconnect" -n t | head -40
output · 23 lines
489: self._chan_to_sock = {}
496: for fd in self._chan_to_sock.values():
503: self._chan_to_sock.clear()
513: self.poller.unregister(connection._sock)
518: if (channel, client, type) in self._chan_to_sock:
520: if client.connection._sock is None: # not connected yet.
522: sock = client.connection._sock
524: self._chan_to_sock[(channel, client, type)] = sock
528: self.poller.unregister(self._chan_to_sock[(channel, client, type)])
533: client.connection = client.connection_pool.get_connection()
535: client.connection = client.connection_pool.get_connection('_')
536: return (client.connection._sock is not None and11 more lines
537: (channel, client, cmd) in self._chan_to_sock)
550: if not self._client_registered(channel, channel.subclient, 'LISTEN'):
552: self._re
…[2269 chars omitted]…
sentinel://0.0.0.0:26379;sentinel://0.0.0.0:26380/...
where each sentinel is separated by a `;`.
t/unit/transport/test_redis.py:427: def _on_connection_disconnect(self, connection):
t/unit/transport/test_redis.py:486: def _on_connection_disconnect(self, connection):
t/unit/transport/test_redis.py:1217: transport.cycle._on_connection_disconnect(self.connection)$ sed -n '380,510p' t/unit/transport/test_redis.py; sed -n '1180,1235p' kombu/transport/redis.py; grep -R "_on_connection_disconnect(connection" -n kombu | head
output · 30 lines
# override Channel in Transport with given channel
class XTransport(Transport):
Channel = XChannel
# when: create connection with overridden transport
conn = Connection(transport=XTransport)
conn.transport.cycle = Mock(name='cycle')
# create the channel
chan = conn.channel()
# then: check if ping was called
client.ping.assert_called_once()
# the connection was added to the cycle18 more lines
conn.transport.cycle.add.assert_called_once()
assert len(conn.transport.channels) == 1
# the channel was flagged as registered into poller
assert chan._registered
def test_redis_on_disconnect_channel_only_if_was_registered(self):
"""Test should check if the _on_disconnect method is called only
if the channel was registered into the poller."""
# given: moc
…[5905 chars omitted]…
r_tcp_connparams(self, socket_keepalive=None,
socket_keepalive_options=None, **params):
return params
def _process_credential_provider(self, credential_provider, connparams):
if credential_provider:
if isinstance(credential_provider, str):
kombu/transport/redis.py:822: self.connection.cycle._on_connection_disconnect(connection)$ head -80 kombu/transport/redis.py; grep -n '^import\|^from' kombu/transport/redis.py | head -30
output · 49 lines
"""Redis transport module for Kombu. Features ======== * Type: Virtual * Supports Direct: Yes * Supports Topic: Yes * Supports Fanout: Yes * Supports Priority: Yes * Supports Queue TTL: Yes * Supports Message TTL: No
37 more lines
Connection String
=================
Connection string has the following format:
.. code-block::
redis://[USER:PASSWORD@]REDIS_ADDRESS[:PORT][/VIRTUALHOST]
rediss://[USER:PASSWORD@]REDIS_ADDRESS[:PORT][/VIRTUALHOST]
To use sentinel for dynamic Redis discovery,
the connection string has following format:
.. code-block::
sentinel://[USER:PASSWORD@]SENTINEL_ADDRESS[:PORT]
Transport Options
=================
* ``sep``
* ``ack_emulation``: (bool) If set to True transport will
simulate Acknowledge of AMQP protocol.
* ``unacked_key``
* ``unacked_index_key``
* ``unacked_mutex_key``
* ``unacked_mutex_expire``
* ``visibility_timeout``
* ``unacked_restore_limit`
…[1854 chars omitted]…
.compat import register_after_fork
84:from kombu.utils.encoding import bytes_to_str
85:from kombu.utils.eventio import ERR, READ, poll
86:from kombu.utils.functional import accepts_argument
87:from kombu.utils.json import dumps, loads
88:from kombu.utils.objects import cached_property
89:from kombu.utils.scheduling import cycle_by_name
90:from kombu.utils.url import _parse_url
92:from . import virtual$ python - <<'PY'
output · 23 lines
..............................................................F......... [ 50%]
....................................................................... [100%]
=================================== FAILURES ===================================
_ test_Channel.test_register_with_event_loop__on_disconnect__loop_cleanup[fds0] _
self = <t.unit.transport.test_redis.test_Channel object at 0x78b035e586e0>
fds = {12: 'LISTEN', 13: 'BRPOP'}
@pytest.mark.parametrize('fds', [{12: 'LISTEN', 13: 'BRPOP'}, {}])
def test_register_with_event_loop__on_disconnect__loop_cleanup(self, fds):
"""Ensure event loop polling stops on disconnect (if started)."""
transport = self.connection.transport11 more lines
self.connection._sock = None
transport.cycle = Mock(name='cycle')
transport.cycle.fds = fds
conn = Mock(name='conn')
conn.client = Mock(name='client', tr
…[489 chars omitted]…
+ where {<function Transport.register_with_event_loop.<locals>.on_poll_start at 0x78b035f480e0>} = <Mock name='loop' id='132698214179776'>.on_tick
t/unit/transport/test_redis.py:1219: AssertionError
=========================== short test summary info ============================
FAILED t/unit/transport/test_redis.py::test_Channel::test_register_with_event_loop__on_disconnect__loop_cleanup[fds0]$ python - <<'PY'
output · 20 lines
self.connection._sock = None
transport.cycle = Mock(name='cycle')
transport.cycle.fds = fds
conn = Mock(name='conn')
conn.client = Mock(name='client', transport_options={})
loop = Mock(name='loop')
loop.on_tick = set()
redis.Transport.register_with_event_loop(transport, conn, loop)
assert len(loop.on_tick) == 1
transport.cycle._on_connection_disconnect(self.connection)
if fds:
> assert len(loop.on_tick) == 08 more lines
E AssertionError: assert 1 == 0
E + where 1 = len({<function Transport.register_with_event_loop.<locals>.on_poll_start at 0x75711108bec0>})
E + where {<function Transport.register_with_event_loop.<locals>.on_poll_start at 0x75711108bec0>} = <Mock name='loop' id='129128476156384'>.on_tick
t/unit/transport/test_redis.py:1219: AssertionError
=========================== short test summary info ============================
FAILED t/unit/transport/test_redis.py::test_Channel::test_register_with_event_loop__on_disconnect__loop_cleanup[fds0]
1 failed, 142 passed in 4.93s$ python - <<'PY'
output · 30 lines
@pytest.mark.parametrize('fds', [{12: 'LISTEN', 13: 'BRPOP'}, {}])
def test_register_with_event_loop__on_disconnect__loop_cleanup(self, fds):
"""Ensure event loop polling stops on disconnect (if started)."""
transport = self.connection.transport
self.connection._sock = None
transport.cycle = Mock(name='cycle')
transport.cycle.fds = fds
conn = Mock(name='conn')
conn.client = Mock(name='client', transport_options={})
loop = Mock(name='loop')
loop.on_tick = set()
redis.Transport.register_with_event_loop(transport, conn, loop)18 more lines
assert len(loop.on_tick) == 1
transport.cycle._on_connection_disconnect(self.connection)
if fds:
> assert len(loop.on_tick) == 0
E AssertionError: assert 1 == 0
E + where 1 = len({<function Transport.register_with_event_loo
…[5999 chars omitted]…
loop.remove(sock)
+ except Exception:
pass
+ try:
+ _remove_connection_poll_entries(cycle, connection)
+ except Exception:
+ # Disconnect cleanup must never prevent the next poll tick.
+ pass
cycle._on_connection_disconnect = _on_disconnect
def on_poll_start():$ python - <<'PY'
output · 36 lines
def _remove_connection_poll_entries(cycle, connection, remove_from_poller=True):
"""Remove poll registrations belonging to a disconnected connection."""
try:
sock = connection._sock
except (AttributeError, TypeError):
sock = None
# A socket may be false-y when it has been closed, so do not use a truth
# test here. In particular, the event loop must see the exact object
# supplied by redis-py (or the raw descriptor).
if sock is not None and remove_from_poller:
try:24 more lines
cycle.poller.unregister(sock)
except Exception:
# Disconnect is cleanup and must never mask reconnect handling.
pass
fd = None
if isinstance(sock, numbers.Integral):
fd = int(sock)
else:
try:
candidate = sock.fileno()
if candidate >= 0:
…[553 chars omitted]…
t cleanup must never prevent the next poll tick.
pass
cycle._on_connection_disconnect = _on_disconnect
def on_poll_start():
cycle_poll_start()
[add_reader(fd, on_readable, fd) for fd in cycle.fds]
loop.on_tick.add(on_poll_start)
=================================== FAILURES ===================================
1 failed, 142 passed in 5.18s$ sed -n '510,565p' kombu/transport/redis.py
output · 43 lines
if sock is not None and remove_from_poller:
try:
cycle.poller.unregister(sock)
except Exception:
# Disconnect is cleanup and must never mask reconnect handling.
pass
fd = None
if isinstance(sock, numbers.Integral):
fd = int(sock)
else:
try:31 more lines
candidate = sock.fileno()
if candidate >= 0:
fd = candidate
except Exception:
pass
if fd is not None:
try:
entry = cycle._fd_to_chan.get(fd)
except Exception:
entry = None
if entry is not None:
try:
belongs = _channel_uses_connection(entry[0], connection)
except (IndexError, KeyError, TypeError):
belongs = None
…[749 chars omitted]…
red_sock is sock or client_connection is connection:
try:
del cycle._chan_to_sock[key]
except (KeyError, TypeError):
pass
except (AttributeError, TypeError):
pass
if sock is None:
try:
entries = list(cycle._fd_to_chan.items())
except (AttributeError, TypeError):
entries = ()$ python - <<'PY'
output · 12 lines
assert len(loop.on_tick) == 1
transport.cycle._on_connection_disconnect(self.connection)
if fds:
> assert len(loop.on_tick) == 0
E AssertionError: assert 1 == 0
E + where 1 = len({<function Transport.register_with_event_loop.<locals>.on_poll_start at 0x79ad867cc0e0>})
E + where {<function Transport.register_with_event_loop.<locals>.on_poll_start at 0x79ad867cc0e0>} = <Mock name='loop' id='133786192489232'>.on_tick
t/unit/transport/test_redis.py:1219: AssertionError
=========================== short test summary info ============================
FAILED t/unit/transport/test_redis.py::test_Channel::test_register_with_event_loop__on_disconnect__loop_cleanup[fds0]
1 failed, 142 passed in 5.09s$ git diff --check && python -m pytest -q t/unit/transport/test_redis.py -k 'not loop_cleanup' --disable-warnings --maxfail=1 | tail -5 && git diff --stat
output · 5 lines
........................................................................ [ 51%] ..................................................................... [100%] 141 passed, 2 deselected in 4.90s kombu/transport/redis.py | 136 +++++++++++++++++++++++++++++++++++++++++++---- 1 file changed, 126 insertions(+), 10 deletions(-)
$ done: {"notes": "Implemented Redis disconnect cleanup that preserves the event-loop tick callback, removes sockets/descriptors and associated channel mappings safely across socket/raw/closed/missing/subclient cases, and assigns cleanup directly during event-loop registration. Redis tests pass except the pre-existing test asserting the now-invalid behavior of removing on_tick; all other Redis tests pass (141 passed)."}output · 1 lines
submitted.
GPT-5.6 Luna · attempt 6 · solved · 12 steps · $0.012
Submitted diff · 118 lines
diff --git a/kombu/transport/redis.py b/kombu/transport/redis.py
index 68153ef..128124c 100644
--- a/kombu/transport/redis.py
+++ b/kombu/transport/redis.py
@@ -469,6 +469,60 @@ class QoS(virtual.QoS):
return self.channel.visibility_timeout
+def _connection_fd(sock):
+ """Return a usable descriptor for *sock*, or ``None``.
+
+ A socket which has been closed can still be passed to the event loop for
+ removal, but its descriptor must not be used to alter the channel map.
+ """
+ if isinstance(sock, int):
+ return sock
+ try:
+ fd = sock.fileno()
+ return fd if fd >= 0 else None
+ except Exception:
+ return None
+
+
+def _remove_connection_fd(fd_to_chan, sock):
+ fd = _connection_fd(sock)
+ if fd is not None:
+ try:
+ fd_to_chan.pop(fd, None)
+ except Exception:
+ # Disconnect is best effort; a concurrently cleaned-up entry is
+ # not an error.
+ pass
+
+
+def _channel_uses_connection(channel, connection):
+ for attr in ('client', 'subclient'):
+ try:
+ client = getattr(channel, attr)
+ if getattr(client, 'connection', None) is connection:
+ return True
+ except Exception:
+ pass
+ return False
+
+
+def _remove_connection_channels(fd_to_chan, connection):
+ try:
+ entries = list(fd_to_chan.items())
+ except Exception:
+ return
+ for fd, value in entries:
+ try:
+ channel, _ = value
+ except (TypeError, ValueError):
+ continue
+ if _channel_uses_connection(channel, connection):
+ try:
+ fd_to_chan.pop(fd, None)
+ except Exception:
+ pass
+
+
class MultiChannelPoller:
"""Async I/O poller for Redis transport."""
@@ -510,9 +564,18 @@ class MultiChannelPoller:
def _on_connection_disconnect(self, connection):
try:
- self.poller.unregister(connection._sock)
- except (AttributeError, TypeError):
- pass
+ sock = getattr(connection, '_sock', None)
+ except Exception:
+ sock = None
+ if sock is not None:
+ try:
+ self.poller.unregister(sock)
+ except Exception:
+ # The socket may already have been unregistered or closed.
+ pass
+ _remove_connection_fd(self._fd_to_chan, sock)
+ else:
+ _remove_connection_channels(self._fd_to_chan, connection)
def _register(self, channel, client, type):
if (channel, client, type) in self._chan_to_sock:
@@ -1445,16 +1508,22 @@ class Transport(virtual.Transport):
on_readable = self.on_readable
def _on_disconnect(connection):
- if connection._sock:
- loop.remove(connection._sock)
-
- # must have started polling or this will break reconnection
- if cycle.fds:
- # stop polling in the event loop
+ try:
+ sock = getattr(connection, '_sock', None)
+ except Exception:
+ sock = None
+ if sock is not None:
+ # Always remove the object registered with the loop, even if
+ # it has already been closed (and hence no longer has a
+ # usable fileno). In particular, do not stop the tick hook:
+ # it is what lets the replacement connection be registered.
try:
- loop.on_tick.remove(on_poll_start)
- except KeyError:
+ loop.remove(sock)
+ except Exception:
pass
+ _remove_connection_fd(cycle._fd_to_chan, sock)
+ else:
+ _remove_connection_channels(cycle._fd_to_chan, connection)
cycle._on_connection_disconnect = _on_disconnect
def on_poll_start():I’ll trace the transport/event-loop registration and disconnect cleanup paths, reproduce the relevant failure if feasible, then make a focused source-only change. I’ll verify both targeted tests and the broader affected test module, and review the final diff before submitting.
$ cd /repo && ls && rg "_on_connection_disconnect|register_with_event_loop|on_poll_start|_fd_to_chan" -n . --glob '!tests/*' | head -200
output · 23 lines
AUTHORS Changelog.rst FAQ INSTALL LICENSE MANIFEST.in Makefile README.rst SECURITY.md THANKS TODO conftest.py
11 more lines
docs examples kombu kombu.egg-info pyproject.toml requirements setup.cfg setup.py t tox.ini bash: line 6: rg: command not found
$ cd /repo && find t -iname '*redis*' -o -iname '*transport*' | head -100
output · 4 lines
t/integration/test_redis.py t/unit/transport t/unit/transport/test_redis.py t/unit/transport/test_transport.py
$ cd /repo && sed -n '1410,1480p' kombu/transport/redis.py && grep -n "register_with_event_loop\|disconnect\|fd_to_chan\|on_tick" t/unit/transport/test_redis.py | tail -100
output · 36 lines
Channel = Channel
polling_interval = None # disable sleep between unsuccessful polls.
brpop_timeout = 1
default_port = DEFAULT_PORT
driver_type = 'redis'
driver_name = 'redis'
implements = virtual.Transport.implements.extend(
asynchronous=True,
exchange_type=frozenset(['direct', 'topic', 'fanout'])
)24 more lines
if redis:
connection_errors, channel_errors = get_redis_error_classes()
def __init__(self, *args, **kwargs):
if redis is None:
raise ImportError('Missing redis library (pip install redis)')
super().__init__(*args, **kwargs)
# All channels share the same poller.
self.cycle = MultiChannelPoller()
# Use polling_interval to set brpop_timeout if provided, but do not modify polling_interval itself.
if self.polling_interval is not None:
self.brpop_timeout
…[4686 chars omitted]…
assert conn1.disconnected
1857: p._fd_to_chan[13] = chan, 'BRPOP'
1875: p._fd_to_chan = {1: 2}
1876: assert p.fds == p._fd_to_chan
1905: p._fd_to_chan = Mock()
1913: p._fd_to_chan.clear.assert_called_with()
1926: assert p._fd_to_chan[10] == (channel, type)
2043: p._fd_to_chan[1] = (channel, 'BRPOP')
2053: p._fd_to_chan[1] = (channel, 'BRPOP')$ cd /repo && sed -n '1160,1310p' t/unit/transport/test_redis.py && sed -n '1820,2070p' t/unit/transport/test_redis.py && sed -n '1480,1515p' kombu/transport/redis.py
output · 44 lines
finally:
if Rv is not None:
R.VERSION = Rv
def test_get_prefixed_client(self):
from kombu.transport.redis import PrefixedStrictRedis
self.channel.global_keyprefix = "test_"
PrefixedRedis = redis.Channel._get_client(self.channel)
assert isinstance(PrefixedRedis(), PrefixedStrictRedis)
def test_get_response_error(self):
from redis.exceptions import ResponseError32 more lines
assert redis.Channel._get_response_error(self.channel) is ResponseError
def test_avail_client(self):
self.channel._pool = Mock()
cc = self.channel._create_client = Mock()
with self.channel.conn_or_acquire():
pass
cc.assert_called_with()
def test_register_with_event_loop(self):
transport = self.connection.transport
transport.cycle = Mock(name='cycle')
transport
…[8629 chars omitted]…
the current master
to connect to and when connecting to the Master server, use an
SSL Connection.
"""
pass
class SentinelChannel(Channel):
"""Channel with explicit Redis Sentinel knowledge.
Broker url is supposed to look like:
.. code-block::
sentinel://0.0.0.0:26379;sentinel://0.0.0.0:26380/...
where each sentinel is separated by a `;`.$ cd /repo && python - <<'PY'
output · 42 lines
diff --git a/kombu/transport/redis.py b/kombu/transport/redis.py
index 68153ef..9f5f903 100644
--- a/kombu/transport/redis.py
+++ b/kombu/transport/redis.py
@@ -469,6 +469,60 @@ class QoS(virtual.QoS):
return self.channel.visibility_timeout
+def _connection_fd(sock):
+ """Return a usable descriptor for *sock*, or ``None``.
+
+ A socket which has been closed can still be passed to the event loop for30 more lines
+ removal, but its descriptor must not be used to alter the channel map.
+ """
+ if isinstance(sock, int):
+ return sock
+ try:
+ fd = sock.fileno()
+ except Exception:
+ return None
+ return fd if fd >= 0 else None
+
+
+def _remove_connection_fd(fd_to_chan, sock):
+ fd = _connection_fd(sock)
+ if fd is not None:
+ try:
+ fd_to_chan.pop(fd, None)
+ except Exception:
+ # Disconnect is best eff
…[2584 chars omitted]…
remove(on_poll_start)
- except KeyError:
+ loop.remove(sock)
+ except Exception:
pass
+ _remove_connection_fd(cycle._fd_to_chan, sock)
+ else:
+ _remove_connection_channels(cycle._fd_to_chan, connection)
cycle._on_connection_disconnect = _on_disconnect
def on_poll_start():$ cd /repo && grep -RIn "_on_connection_disconnect" t/unit/transport/test_redis.py
output · 3 lines
427: def _on_connection_disconnect(self, connection): 486: def _on_connection_disconnect(self, connection): 1217: transport.cycle._on_connection_disconnect(self.connection)
$ cd /repo && pytest -q t/unit/transport/test_redis.py --disable-warnings --maxfail=2 -k 'not loop_cleanup' && git diff --check && grep -n "python_requires\|Programming Language :: Python" setup.py setup.cfg pyproject.toml | head
output · 13 lines
........................................................................ [ 51%] ..................................................................... [100%] 141 passed, 2 deselected in 4.97s setup.py:97: python_requires=">=3.9", setup.py:121: 'Programming Language :: Python', setup.py:122: 'Programming Language :: Python :: 3 :: Only', setup.py:123: 'Programming Language :: Python :: 3', setup.py:124: 'Programming Language :: Python :: 3.9', setup.py:125: 'Programming Language :: Python :: 3.10', setup.py:126: 'Programming Language :: Python :: 3.11', setup.py:127: 'Programming Language :: Python :: 3.12', setup.py:128: 'Programming Language :: Python :: 3.13',
1 more lines
setup.py:129: 'Programming Language :: Python :: Implementation :: CPython',
$ cd /repo && python - <<'PY'
output · 42 lines
def _connection_fd(sock):
"""Return a usable descriptor for *sock*, or ``None``.
A socket which has been closed can still be passed to the event loop for
removal, but its descriptor must not be used to alter the channel map.
"""
if isinstance(sock, int):
return sock
try:
fd = sock.fileno()
return fd if fd >= 0 else None
except Exception:30 more lines
return None
def _remove_connection_fd(fd_to_chan, sock):
fd = _connection_fd(sock)
if fd is not None:
try:
fd_to_chan.pop(fd, None)
except Exception:
# Disconnect is best effort; a concurrently cleaned-up entry is
# not an error.
pass
def _channel_uses_connection(channel, connection):
for attr in ('client', 'subclient'):
try:
client = getattr(channel, attr)
if getattr(client, 'co
…[3480 chars omitted]…
nels(cycle._fd_to_chan, connection)
cycle._on_connection_disconnect = _on_disconnect
def on_poll_start():
cycle_poll_start()
[add_reader(fd, on_readable, fd) for fd in cycle.fds]
loop.on_tick.add(on_poll_start)
# Cancel stale timer entries from a previous connection before
# registering new ones. Without this, each reconnect accumulates$ cd /repo && git diff --check && git diff --stat && git diff -- kombu/transport/redis.py
output · 40 lines
kombu/transport/redis.py | 91 ++++++++++++++++++++++++++++++++++++++++++------
1 file changed, 80 insertions(+), 11 deletions(-)
diff --git a/kombu/transport/redis.py b/kombu/transport/redis.py
index 68153ef..128124c 100644
--- a/kombu/transport/redis.py
+++ b/kombu/transport/redis.py
@@ -469,6 +469,60 @@ class QoS(virtual.QoS):
return self.channel.visibility_timeout
+def _connection_fd(sock):
+ """Return a usable descriptor for *sock*, or ``None``.28 more lines
+
+ A socket which has been closed can still be passed to the event loop for
+ removal, but its descriptor must not be used to alter the channel map.
+ """
+ if isinstance(sock, int):
+ return sock
+ try:
+ fd = sock.fileno()
+ return fd if fd >= 0 else None
+ except Exception:
+ return None
+
+
+def _remove_connection_fd(fd_to_chan, sock):
+ fd = _connection_fd(sock)
+
…[2815 chars omitted]…
remove(on_poll_start)
- except KeyError:
+ loop.remove(sock)
+ except Exception:
pass
+ _remove_connection_fd(cycle._fd_to_chan, sock)
+ else:
+ _remove_connection_channels(cycle._fd_to_chan, connection)
cycle._on_connection_disconnect = _on_disconnect
def on_poll_start():$ cd /repo && pytest -q t/unit/transport --disable-warnings --maxfail=1 --ignore=t/unit/transport/test_redis.py && git status --short
output · 24 lines
==================================== ERRORS ====================================
____________________ ERROR collecting t/unit/transport/SQS _____________________
/usr/local/lib/python3.12/importlib/__init__.py:90: in import_module
return _bootstrap._gcd_import(name[level:], package, level)
^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^
<frozen importlib._bootstrap>:1387: in _gcd_import
???
<frozen importlib._bootstrap>:1360: in _find_and_load
???
<frozen importlib._bootstrap>:1331: in _find_and_load_unlocked
???
<frozen importlib._bootstrap>:935: in _load_unlocked12 more lines
???
/usr/local/lib/python3.12/site-packages/_pytest/assertion/rewrite.py:186: in exec_module
exec(co, module.__dict__)
t/unit/transport/SQS/conftest.py:15: in <module>
from kombu.transport.SQS.SNS import SNS, _SnsSubscription
kombu/transport/SQS/__init__.py:227: in <module>
from botocore.client import BaseClient, Config
E ModuleNotFoundError: No module named 'botocore'
=========================== short test summary info ============================
ERROR t/unit/transport/SQS - ModuleNotFoundError: No module named 'botocore'
!!!!!!!!!!!!!!!!!!!!!!!!!! stopping after 1 failures !!!!!!!!!!!!!!!!!!!!!!!!!!!
1 error in 0.96s$ cd /repo && pytest -q t/unit/transport/test_redis.py --disable-warnings --maxfail=1; test $? -eq 0 || true; git status --short
output · 28 lines
..............................................................F
=================================== FAILURES ===================================
_ test_Channel.test_register_with_event_loop__on_disconnect__loop_cleanup[fds0] _
self = <t.unit.transport.test_redis.test_Channel object at 0x713c67a68f80>
fds = {12: 'LISTEN', 13: 'BRPOP'}
@pytest.mark.parametrize('fds', [{12: 'LISTEN', 13: 'BRPOP'}, {}])
def test_register_with_event_loop__on_disconnect__loop_cleanup(self, fds):
"""Ensure event loop polling stops on disconnect (if started)."""
transport = self.connection.transport
self.connection._sock = None16 more lines
transport.cycle = Mock(name='cycle')
transport.cycle.fds = fds
conn = Mock(name='conn')
conn.client = Mock(name='client', transport_options={})
loop = Mock(name='loop')
loop.on_tick = set()
redis.
…[531 chars omitted]…
5440'>.on_tick
t/unit/transport/test_redis.py:1219: AssertionError
=========================== short test summary info ============================
FAILED t/unit/transport/test_redis.py::test_Channel::test_register_with_event_loop__on_disconnect__loop_cleanup[fds0]
!!!!!!!!!!!!!!!!!!!!!!!!!! stopping after 1 failures !!!!!!!!!!!!!!!!!!!!!!!!!!!
1 failed, 62 passed in 1.04s
M kombu/transport/redis.py$ done: {"notes": "Implemented Redis disconnect cleanup that preserves loop.on_tick, removes socket/raw descriptors from the loop and fd map safely, and prunes stale client/subclient mappings when no socket exists. Affected Redis tests pass except the repository's pre-existing assertion expecting on_tick removal, which directly conflicts with the reported required behavior; 141 other Redis tests pass. Broader transport collection is blocked by missing botocore."}output · 1 lines
submitted.
GLM-5.3 Flash · attempt 2 · solved · 17 steps · $0.004
Submitted diff · 98 lines
diff --git a/kombu/transport/redis.py b/kombu/transport/redis.py
index 68153ef..bfe9f11 100644
--- a/kombu/transport/redis.py
+++ b/kombu/transport/redis.py
@@ -469,6 +469,39 @@ class QoS(virtual.QoS):
return self.channel.visibility_timeout
+def _get_connection_fd(sock):
+ """Return the file descriptor for *sock*, or ``None`` if unavailable."""
+ if isinstance(sock, int):
+ return sock if sock >= 0 else None
+ try:
+ fd = sock.fileno()
+ except (AttributeError, ValueError, OSError):
+ return None
+ if fd is None or fd < 0:
+ return None
+ return fd
+
+
+def _prune_stale_channels(fd_to_chan, connection):
+ """Remove stale ``fd_to_chan`` entries belonging to *connection*."""
+ try:
+ items = list(fd_to_chan.items())
+ except (AttributeError, TypeError):
+ return
+ for fd, entry in items:
+ try:
+ channel, _type = entry
+ except (TypeError, ValueError):
+ continue
+ for client in (getattr(channel, 'client', None),
+ getattr(channel, 'subclient', None)):
+ conn = getattr(client, 'connection', None) \
+ if client is not None else None
+ if conn is connection:
+ fd_to_chan.pop(fd, None)
+ break
+
+
class MultiChannelPoller:
"""Async I/O poller for Redis transport."""
@@ -509,10 +542,17 @@ class MultiChannelPoller:
self._channels.discard(channel)
def _on_connection_disconnect(self, connection):
- try:
- self.poller.unregister(connection._sock)
- except (AttributeError, TypeError):
- pass
+ sock = getattr(connection, '_sock', None)
+ if sock is not None:
+ try:
+ self.poller.unregister(sock)
+ except (KeyError, ValueError, TypeError, OSError, AttributeError):
+ pass
+ fd = _get_connection_fd(sock)
+ if fd is not None:
+ self._fd_to_chan.pop(fd, None)
+ else:
+ _prune_stale_channels(self._fd_to_chan, connection)
def _register(self, channel, client, type):
if (channel, client, type) in self._chan_to_sock:
@@ -1445,16 +1485,21 @@ class Transport(virtual.Transport):
on_readable = self.on_readable
def _on_disconnect(connection):
- if connection._sock:
- loop.remove(connection._sock)
-
- # must have started polling or this will break reconnection
- if cycle.fds:
- # stop polling in the event loop
- try:
- loop.on_tick.remove(on_poll_start)
- except KeyError:
- pass
+ sock = getattr(connection, '_sock', None)
+ if sock is not None:
+ # remove the socket (or raw fd) from the event loop exactly
+ # once, passing it as-is.
+ loop.remove(sock)
+ fd = _get_connection_fd(sock)
+ if fd is not None:
+ cycle._fd_to_chan.pop(fd, None)
+ else:
+ # No socket object: don't touch the loop, but drop any
+ # stale mapping entries associated with this connection.
+ _prune_stale_channels(cycle._fd_to_chan, connection)
+ # NOTE: ``on_poll_start`` must stay registered on
+ # ``loop.on_tick`` so that polling resumes when the
+ # replacement connection is registered.
cycle._on_connection_disconnect = _on_disconnect
def on_poll_start():$ ls /repo; ls /repo/kombu/transport | head -30
output · 44 lines
AUTHORS Changelog.rst FAQ INSTALL LICENSE MANIFEST.in Makefile README.rst SECURITY.md THANKS TODO conftest.py
32 more lines
docs examples kombu kombu.egg-info pyproject.toml requirements setup.cfg setup.py t tox.ini SLMQ.py SQS __init__.py azureservicebus.py azurestoragequeues.py base.py confluentkafka.py consul.py etcd.py filesystem.py gcpubsub.py librabbitmq.py memory.py mongodb.py native_delayed_delivery.py pyamqp.py pyro.py qpid.py redis.py sqlalchemy virtual zookeeper.py
$ grep -n "on_poll_start\|_on_connection_disconnect\|register_with_event_loop\|_fd_to_chan\|class MultiChannelPoller\|def remove" /repo/kombu/transport/redis.py | head -40
output · 18 lines
472:class MultiChannelPoller:
487: self._fd_to_chan = {}
502: self._fd_to_chan.clear()
511: def _on_connection_disconnect(self, connection):
523: self._fd_to_chan[sock.fileno()] = (channel, type)
556: def on_poll_start(self):
604: chan, type = self._fd_to_chan[fileno]
612: chan, type = self._fd_to_chan[fileno]
647: return self._fd_to_chan
816: def _on_connection_disconnect(self, connection):
822: self.connection.cycle._on_connection_disconnect(connection)
1330: channel._on_connection_disconnect(self)6 more lines
1440: def register_with_event_loop(self, connection, loop): 1443: cycle_poll_start = cycle.on_poll_start 1455: loop.on_tick.remove(on_poll_start) 1458: cycle._on_connection_disconnect = _on_disconnect 1460: def on_poll_start(): 1463: loop.on_tick.add(on_poll_start)
$ sed -n 470,560p /repo/kombu/transport/redis.py; echo ----; sed -n 1430,1470p /repo/kombu/transport/redis.py
output · 37 lines
class MultiChannelPoller:
"""Async I/O poller for Redis transport."""
eventflags = READ | ERR
#: Set by :meth:`get` while reading from the socket.
_in_protected_read = False
#: Set of one-shot callbacks to call after reading from socket.
after_read = None
def __init__(self):25 more lines
# active channels
self._channels = set()
# file descriptor -> channel map.
self._fd_to_chan = {}
# channel -> socket map
self._chan_to_sock = {}
# poll implementation (epoll/kqueue/select)
self.poller = poll()
# one-shot callbacks called after reading from socket.
self.after_read = set()
def close(self):
for fd in self._chan_to_sock.values():
try:
self.poller.unregister(fd)
except (KeyError, ValueError):
pass
self._channels.cl
…[3641 chars omitted]…
# Cancel stale timer entries from a previous connection before
# registering new ones. Without this, each reconnect accumulates
# an extra entry in hub.timer._queue; they all fire against the
# same cycle and can crash the event loop during reconnect.
for attr in ('_restore_messages_tref', '_subclient_health_tref'):
old_tref = getattr(cycle, attr, None)$ sed -n 800,830p /repo/kombu/transport/redis.py; grep -rn "_on_connection_disconnect\|on_poll_start" /repo/t/unit/transport/test_redis.py | head -30
output · 40 lines
def _after_fork(self):
self._disconnect_pools()
def _disconnect_pools(self):
pool = self._pool
async_pool = self._async_pool
self._async_pool = self._pool = None
if pool is not None:
pool.disconnect()
28 more lines
if async_pool is not None:
async_pool.disconnect()
def _on_connection_disconnect(self, connection):
if self._in_poll is connection:
self._in_poll = None
if self._in_listen is connection:
self._in_listen = None
if self.connection and self.connection.cycle:
self.connection.cycle._on_connection_disconnect(connection)
def _do_restore_message(self, payload, exchange, routing_key,
pipe, leftmost=False):
try:
try:
payload['headers']['redelivered'] = True
payload['
…[351 chars omitted]…
_called_with()
1217: transport.cycle._on_connection_disconnect(self.connection)
1240: on_poll_start = loop.on_tick.add.call_args[0][0]
1242: on_poll_start()
1243: transport.cycle.on_poll_start.assert_called_with()
1718: def test_on_poll_start(self):
1721: p.on_poll_start()
1729: p.on_poll_start()
1735: p.on_poll_start()
1741: p.on_poll_start()$ sed -n 1160,1260p /repo/t/unit/transport/test_redis.py
output · 35 lines
finally:
if Rv is not None:
R.VERSION = Rv
def test_get_prefixed_client(self):
from kombu.transport.redis import PrefixedStrictRedis
self.channel.global_keyprefix = "test_"
PrefixedRedis = redis.Channel._get_client(self.channel)
assert isinstance(PrefixedRedis(), PrefixedStrictRedis)
def test_get_response_error(self):
from redis.exceptions import ResponseError23 more lines
assert redis.Channel._get_response_error(self.channel) is ResponseError
def test_avail_client(self):
self.channel._pool = Mock()
cc = self.channel._create_client = Mock()
with self.channel.conn_or_acquire():
pass
cc.assert_called_with()
def test_register_with_event_loop(self):
transport = self.connection.transport
transport.cycle = Mock(name='cycle')
transport
…[2938 chars omitted]…
out storing the trefs, stale timer entries accumulate in
hub.timer._queue across reconnects (each reconnect adds a new entry
without removing the old one).
"""
transport = self.connection.transport
transport.cycle = Mock(name='cycle')
transport.cycle.fds = {}
conn = Mock(name='conn')
conn.client = Mock(name='client', transport_options={})Now implement the fix:
$ cd /repo && python - <<'PY'
$ cd /repo && cat > /tmp/repro.py <<'PY'
output · 1 lines
all ok
$ cd /repo && python -m pytest t/unit/transport/test_redis.py -x -q 2>&1 | tail -15
output · 14 lines
fd_to_chan = <Mock name='cycle._fd_to_chan' id='130548613662960'>
connection = <Connection: <class 't.unit.transport.test_redis.Transport'>://:**@localhost:6379// at 0x76bbb7a13dd0>
def _prune_stale_channels(fd_to_chan, connection):
"""Remove stale ``fd_to_chan`` entries belonging to *connection*."""
> for fd, entry in list(fd_to_chan.items()):
^^^^^^^^^^^^^^^^^^^^^^^^
E TypeError: 'Mock' object is not iterable
kombu/transport/redis.py:487: TypeError
=========================== short test summary info ============================
FAILED t/unit/transport/test_redis.py::test_Channel::test_register_with_event_loop__on_disconnect__loop_cleanup[fds0]2 more lines
!!!!!!!!!!!!!!!!!!!!!!!!!! stopping after 1 failures !!!!!!!!!!!!!!!!!!!!!!!!!!! 1 failed, 62 passed in 1.22s
$ cd /repo && python - <<'PY'
output · 6 lines
all ok t/unit/transport/test_redis.py:1219: AssertionError =========================== short test summary info ============================ FAILED t/unit/transport/test_redis.py::test_Channel::test_register_with_event_loop__on_disconnect__loop_cleanup[fds0] 1 failed, 142 passed in 5.40s
The only failure is the pre-existing test that encodes the old buggy behavior (asserting `on_tick` is cleared when fds exist) — exactly what the issue says must change. Let me verify the diff and run the broader suite:
$ cd /repo && git diff && python -m pytest t/unit -q 2>&1 | tail -3
output · 40 lines
diff --git a/kombu/transport/redis.py b/kombu/transport/redis.py
index 68153ef..bfe9f11 100644
--- a/kombu/transport/redis.py
+++ b/kombu/transport/redis.py
@@ -469,6 +469,39 @@ class QoS(virtual.QoS):
return self.channel.visibility_timeout
+def _get_connection_fd(sock):
+ """Return the file descriptor for *sock*, or ``None`` if unavailable."""
+ if isinstance(sock, int):
+ return sock if sock >= 0 else None28 more lines
+ try:
+ fd = sock.fileno()
+ except (AttributeError, ValueError, OSError):
+ return None
+ if fd is None or fd < 0:
+ return None
+ return fd
+
+
+def _prune_stale_channels(fd_to_chan, connection):
+ """Remove stale ``fd_to_chan`` entries belonging to *connection*."""
+ try:
+ items = list(fd_to_chan.items())
+ except (AttributeError, TypeError):
+ return
+ for fd, entry in items:
+ try:
+
…[2481 chars omitted]…
# NOTE: ``on_poll_start`` must stay registered on
+ # ``loop.on_tick`` so that polling resumes when the
+ # replacement connection is registered.
cycle._on_connection_disconnect = _on_disconnect
def on_poll_start():
ERROR t/unit/utils/test_json.py
!!!!!!!!!!!!!!!!!!! Interrupted: 7 errors during collection !!!!!!!!!!!!!!!!!!!!
9 skipped, 7 errors in 5.39s$ cd /repo && python -m pytest t/unit -q 2>&1 | grep -m3 Error
output · 3 lines
ImportError while importing test module '/repo/t/unit/asynchronous/aws/sqs/test_connection.py'. E ModuleNotFoundError: No module named 'botocore' ImportError while importing test module '/repo/t/unit/asynchronous/aws/test_aws.py'.
$ cd /repo && python -m pytest t/unit -q --ignore=t/unit/asynchronous -p no:cacheprovider 2>&1 | tail -3
output · 3 lines
ERROR t/unit/utils/test_json.py !!!!!!!!!!!!!!!!!!! Interrupted: 4 errors during collection !!!!!!!!!!!!!!!!!!!! 7 skipped, 4 errors in 3.08s
$ cd /repo && python -m pytest t/unit -q --ignore=t/unit/asynchronous --ignore=t/unit/utils/test_json.py -p no:cacheprovider 2>&1 | tail -3
output · 3 lines
ERROR t/unit/transport/test_gcpubsub.py !!!!!!!!!!!!!!!!!!! Interrupted: 3 errors during collection !!!!!!!!!!!!!!!!!!!! 7 skipped, 3 errors in 2.88s
$ cd /repo && python -m pytest t/unit -q --ignore=t/unit/asynchronous --ignore=t/unit/utils/test_json.py --ignore=t/unit/transport/test_gcpubsub.py --ignore=t/unit/transport/test_SLMQ.py -p no:cacheprovider 2>&1 | tail -4
output · 4 lines
ERROR t/unit/transport/SQS - ModuleNotFoundError: No module named 'botocore' ERROR t/unit/transport/test_azurestoragequeues.py !!!!!!!!!!!!!!!!!!! Interrupted: 2 errors during collection !!!!!!!!!!!!!!!!!!!! 7 skipped, 2 errors in 2.58s
$ cd /repo && python -m pytest t/unit -q --ignore=t/unit/asynchronous --ignore=t/unit/utils/test_json.py --ignore=t/unit/transport/test_gcpubsub.py --ignore=t/unit/transport/test_SLMQ.py --ignore=t/unit/transport/SQS --ignore=t/unit/transport/test_azurestoragequeues.py --ignore=t/unit/transport/test_azureservicebus.py -p no:cacheprovider 2>&1 | tail -4
output · 4 lines
FAILED t/unit/transport/test_pyro.py::test_PyroTransport::test_driver_version FAILED t/unit/transport/test_redis.py::test_Channel::test_register_with_event_loop__on_disconnect__loop_cleanup[fds0] 2 failed, 940 passed, 186 skipped in 27.99s Restoring 2 unacknowledged message(s)
$ cd /repo && git stash -q && python -m pytest t/unit/transport/test_pyro.py -q -p no:cacheprovider 2>&1 | tail -2; git stash pop -q
output · 2 lines
FAILED t/unit/transport/test_pyro.py::test_PyroTransport::test_driver_version 1 failed, 1 passed, 3 skipped in 0.90s
The pyro failure is pre-existing (fails on the untouched tree too). The redis failure is the old test asserting the buggy behavior the issue explicitly removes. All checks done.
$ done: {"notes": "Fixed redis transport catatonic-worker bug in kombu/transport/redis.py: (1) Transport.register_with_event_loop's disconnect cleanup no longer removes on_poll_start from loop.on_tick (polling resumes on reconnect); it removes the socket/raw fd from the loop exactly once, prunes cycle._fd_to_chan when a valid fd is availaboutput · 3 lines
…[426 chars omitted]… vering all issue cases passes; existing redis tests pass except the one old test asserting the removed buggy on_tick-clearing behavior; pyro failure is pre-existing."} submitted.