SWE-Race › Tasks › celery-kombu-2492 ← prevnext →

celery-kombu-2492

celery/kombusplitsinglemerged 2026-03-26BSD-3-Clausefix: 1 file, +53 −65 fail-to-pass · 145 pass-to-pass
Results
Modelsolved / attemptsmedian stepsmedian costattempts
GPT-5.6 Luna5/612$0.0121✓ 2✓ 3✓ 4✓ 5✗ 6✓
DeepSeek V4 Flash1/255$0.0361✗ 2✓
GLM-5.3 Flash2/220$0.0061✓ 2✓
The prompt the agent sees

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_register_with_event_loop__on_disconnect__loop_cleanuptest_register_with_event_loop__on_disconnect__removes_socktest_register_with_event_loop__on_disconnect__sock_no_filenotest_register_with_event_loop__on_disconnect__sock_none_pruntest_register_with_event_loop__on_disconnect__sock_none_subc
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():

What each model did

Submitted diff first, then every command. Outputs show 12 lines; expand for the rest.
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():
step 1

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()
step 2
$ 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 = fds
8 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')
step 3
$ 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
                    )
                    return
19 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):
step 4
$ 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.transport
11 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]
step 5
$ 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():
step 6
$ 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 = None
16 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
step 7
$ python - <<'PY'
output · 4 lines
fd edge cases passed
association edge cases passed
...                                                                      [100%]
3 passed, 140 deselected in 0.54s
step 9
$ 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():
step 1

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')
step 2
$ 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):
step 3
$ 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:
                    break
41 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):
step 4
$ 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.transport
11 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]
step 5
$ 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 None
36 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():
step 6
$ 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
step 7
$ 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 = None
16 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
step 8
$ git diff --check && git status --short
output · 1 lines
M kombu/transport/redis.py
step 9
$ 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._sock
32 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():
step 11
$ 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():
step 1

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
step 2
$ find kombu -iname '*redis*' -o -iname '*event*' | head -80
output · 2 lines
kombu/transport/redis.py
kombu/utils/eventio.py
step 3
$ 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.Client
27 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):
step 4
$ 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 and
11 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)
step 5
$ 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 cycle
18 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)
step 6
$ 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
step 7
$ 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.transport
11 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]
step 8
$ 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) == 0
8 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
step 9
$ 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():
step 10
$ 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
step 11
$ 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 = ()
step 12
$ 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
step 13
$ 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(-)
step 14
$ 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():
step 1

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
step 2
$ 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
step 3
$ 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')
step 4
$ 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 ResponseError
32 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 `;`.
step 5
$ 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 for
30 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():
step 6
$ 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)
step 7
$ 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',
step 8
$ 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
step 9
$ 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():
step 10
$ 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_unlocked
12 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
step 11
$ 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 = None
16 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
step 12
$ 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():
step 1
$ 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
step 2
$ 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)
step 3
$ 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)
step 4
$ 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()
step 5
$ 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 ResponseError
23 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={})
step 6

Now implement the fix:

$ cd /repo && python - <<'PY'
step 7
$ cd /repo && cat > /tmp/repro.py <<'PY'
output · 1 lines
all ok
step 8
$ 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
step 9
$ 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
step 10

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 None
28 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
step 11
$ 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'.
step 12
$ 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
step 13
$ 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
step 14
$ 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
step 15
$ 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)
step 16
$ 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
step 17

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 availab
output · 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.