SWE-Race › Tasks › redis-redis-py-3863 ← prevnext →

redis-redis-py-3863

redis/redis-pysplitsinglemerged 2025-12-04MITfix: 2 files, +16 −22 fail-to-pass · 76 pass-to-pass
Results
Modelsolved / attemptsmedian stepsmedian costattempts
GPT-5.6 Luna5/616$0.0171✗ 2✓ 3✓ 4✓ 5✓ 6✓
DeepSeek V4 Flash2/262$0.0461✓ 2✓
GLM-5.3 Flash2/236$0.0131✓ 2✓
The prompt the agent sees

When establishing a connection, the retry mechanism only covers the socket connection and not failures that occur during the initial handshake. If the socket connects but sending the handshake commands raises a retryable `OSError` such as `ECONNREFUSED`, the client fails immediately instead of retrying the complete connection process according to its configured retry policy. The handshake must be attempted once per connection attempt, and after retries are exhausted the connection error should be reported.

Errors that are not configured as retryable, including cancellation and other Redis-related errors, must still propagate immediately without starting additional connection attempts.

Hidden tests · 2 fail-to-pass, 76 pass-to-passrun after the agent submits, in a clean verifier
test_connect_with_retriestest_connect_with_retries
Test patch · 82 lines
diff --git a/tests/test_asyncio/test_connection.py b/tests/test_asyncio/test_connection.py
index 35d404a36e..9123c32ee1 100644
--- a/tests/test_asyncio/test_connection.py
+++ b/tests/test_asyncio/test_connection.py
@@ -155,16 +155,33 @@ async def mock_connect():
     await conn.disconnect()
 
 
-async def test_connect_without_retry_on_os_error():
-    """Test that the _connect function is not being retried in case of a OSError"""
+async def test_connect_without_retry_on_non_retryable_error():
+    """
+    Test that the _connect function is not being retried in case of a CancelledError -
+    error that is not in the list of retry-able errors"""
     with patch.object(Connection, "_connect") as _connect:
-        _connect.side_effect = OSError("")
+        _connect.side_effect = asyncio.CancelledError("")
         conn = Connection(retry_on_timeout=True, retry=Retry(NoBackoff(), 2))
-        with pytest.raises(ConnectionError):
+        with pytest.raises(asyncio.CancelledError):
             await conn.connect()
         assert _connect.call_count == 1
 
 
+async def test_connect_with_retries():
+    """
+    Test that retries occur for the entire connect+handshake flow when OSError happens during the handshake phase.
+    """
+    with patch.object(asyncio.StreamWriter, "writelines") as writelines:
+        writelines.side_effect = OSError(ECONNREFUSED)
+        conn = Connection(retry_on_timeout=True, retry=Retry(NoBackoff(), 2))
+        with pytest.raises(ConnectionError):
+            await conn.connect()
+        # the handshake commands are the failing ones
+        # validate that we don't execute too many commands on each retry
+        # 3 retries --> 3 commands
+        assert writelines.call_count == 3
+
+
 async def test_connect_timeout_error_without_retry():
     """Test that the _connect function is not being retried if retry_on_timeout is
     set to False"""
diff --git a/tests/test_connection.py b/tests/test_connection.py
index 441528ca6b..07424e1ce3 100644
--- a/tests/test_connection.py
+++ b/tests/test_connection.py
@@ -124,16 +124,31 @@ def mock_connect():
         assert conn._connect.call_count == 3
         self.clear(conn)
 
-    def test_connect_without_retry_on_os_error(self):
-        """Test that the _connect function is not being retried in case of a OSError"""
+    def test_connect_without_retry_on_non_retryable_error(self):
+        """Test that the _connect function is not being retried in case of a non-retryable error"""
         with patch.object(Connection, "_connect") as _connect:
-            _connect.side_effect = OSError("")
+            _connect.side_effect = RedisError("")
             conn = Connection(retry_on_timeout=True, retry=Retry(NoBackoff(), 2))
-            with pytest.raises(ConnectionError):
+            with pytest.raises(RedisError):
                 conn.connect()
             assert _connect.call_count == 1
             self.clear(conn)
 
+    def test_connect_with_retries(self):
+        """
+        Validate that retries occur for the entire connect+handshake flow when OSError
+        happens during the handshake phase.
+        """
+        with patch.object(socket.socket, "sendall") as sendall:
+            sendall.side_effect = OSError(ECONNREFUSED)
+            conn = Connection(retry_on_timeout=True, retry=Retry(NoBackoff(), 2))
+            with pytest.raises(ConnectionError):
+                conn.connect()
+            # the handshake commands are the failing ones
+            # validate that we don't execute too many commands on each retry
+            # 3 retries --> 3 commands
+            assert sendall.call_count == 3
+
     def test_connect_timeout_error_without_retry(self):
         """Test that the _connect function is not being retried if retry_on_timeout is
         set to False"""
Reference fix · 2 files, +16 −2the 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.

redis/asyncio/connection.py, redis/connection.py

diff --git a/redis/asyncio/connection.py b/redis/asyncio/connection.py
index c5764f7343..1d50b53ee2 100644
--- a/redis/asyncio/connection.py
+++ b/redis/asyncio/connection.py
@@ -296,7 +296,14 @@ def set_parser(self, parser_class: Type[BaseParser]) -> None:
 
     async def connect(self):
         """Connects to the Redis server if not already connected"""
-        await self.connect_check_health(check_health=True)
+        # try once the socket connect with the handshake, retry the whole
+        # connect/handshake flow based on retry policy
+        await self.retry.call_with_retry(
+            lambda: self.connect_check_health(
+                check_health=True, retry_socket_connect=False
+            ),
+            lambda error: self.disconnect(),
+        )
 
     async def connect_check_health(
         self, check_health: bool = True, retry_socket_connect: bool = True
diff --git a/redis/connection.py b/redis/connection.py
index 0a87777ac3..fe25a9e3f4 100644
--- a/redis/connection.py
+++ b/redis/connection.py
@@ -843,7 +843,14 @@ def _get_parser(self) -> Union[_HiredisParser, _RESP3Parser, _RESP2Parser]:
 
     def connect(self):
         "Connects to the Redis server if not already connected"
-        self.connect_check_health(check_health=True)
+        # try once the socket connect with the handshake, retry the whole
+        # connect/handshake flow based on retry policy
+        self.retry.call_with_retry(
+            lambda: self.connect_check_health(
+                check_health=True, retry_socket_connect=False
+            ),
+            lambda error: self.disconnect(error),
+        )
 
     def connect_check_health(
         self, check_health: bool = True, retry_socket_connect: bool = True

What each model did

Submitted diff first, then every command. Outputs show 12 lines; expand for the rest.
GLM-5.3 Flash · attempt 2 · solved · 44 steps · $0.014
Submitted diff · 141 lines
diff --git a/redis/asyncio/connection.py b/redis/asyncio/connection.py
index c5764f7..ef53c3d 100644
--- a/redis/asyncio/connection.py
+++ b/redis/asyncio/connection.py
@@ -306,34 +306,23 @@ class AbstractConnection:
         try:
             if retry_socket_connect:
                 await self.retry.call_with_retry(
-                    lambda: self._connect(), lambda error: self.disconnect()
+                    lambda: self._connect_and_handshake(check_health),
+                    lambda error: self.disconnect(),
                 )
             else:
-                await self._connect()
+                await self._connect_and_handshake(check_health)
         except asyncio.CancelledError:
             raise  # in 3.7 and earlier, this is an Exception, not BaseException
         except (socket.timeout, asyncio.TimeoutError):
             raise TimeoutError("Timeout connecting to server")
         except OSError as e:
             raise ConnectionError(self._error_message(e))
-        except Exception as exc:
-            raise ConnectionError(exc) from exc
-
-        try:
-            if not self.redis_connect_func:
-                # Use the default on_connect function
-                await self.on_connect_check_health(check_health=check_health)
-            else:
-                # Use the passed function redis_connect_func
-                (
-                    await self.redis_connect_func(self)
-                    if asyncio.iscoroutinefunction(self.redis_connect_func)
-                    else self.redis_connect_func(self)
-                )
         except RedisError:
-            # clean up after any error in on_connect
+            # clean up after any error in the handshake
             await self.disconnect()
             raise
+        except Exception as exc:
+            raise ConnectionError(exc) from exc
 
         # run any user callbacks. right now the only internal callback
         # is for pubsub channel/pattern resubscription
@@ -351,6 +340,30 @@ class AbstractConnection:
     def should_reconnect(self):
         return self._should_reconnect
 
+    async def _connect_and_handshake(self, check_health: bool = True):
+        """Establish the socket connection and perform the initial handshake.
+
+        This whole process counts as a single connection attempt, so that
+        retryable errors raised during the handshake are retried together
+        with the socket connection.
+        """
+        await self._connect()
+        try:
+            if not self.redis_connect_func:
+                # Use the default on_connect function
+                await self.on_connect_check_health(check_health=check_health)
+            else:
+                # Use the passed function redis_connect_func
+                (
+                    await self.redis_connect_func(self)
+                    if asyncio.iscoroutinefunction(self.redis_connect_func)
+                    else self.redis_connect_func(self)
+                )
+        except BaseException:
+            # clean up after any error in the handshake
+            await self.disconnect()
+            raise
+
     @abstractmethod
     async def _connect(self):
         pass
diff --git a/redis/connection.py b/redis/connection.py
index 0a87777..4e1a2e3 100644
--- a/redis/connection.py
+++ b/redis/connection.py
@@ -852,26 +852,18 @@ class AbstractConnection(MaintNotificationsAbstractConnection, ConnectionInterfa
             return
         try:
             if retry_socket_connect:
-                sock = self.retry.call_with_retry(
-                    lambda: self._connect(), lambda error: self.disconnect(error)
+                self.retry.call_with_retry(
+                    lambda: self._connect_and_handshake(check_health),
+                    lambda error: self.disconnect(error),
                 )
             else:
-                sock = self._connect()
+                self._connect_and_handshake(check_health)
         except socket.timeout:
             raise TimeoutError("Timeout connecting to server")
         except OSError as e:
             raise ConnectionError(self._error_message(e))
-
-        self._sock = sock
-        try:
-            if self.redis_connect_func is None:
-                # Use the default on_connect function
-                self.on_connect_check_health(check_health=check_health)
-            else:
-                # Use the passed function redis_connect_func
-                self.redis_connect_func(self)
         except RedisError:
-            # clean up after any error in on_connect
+            # clean up after any error in the handshake
             self.disconnect()
             raise
 
@@ -884,6 +876,28 @@ class AbstractConnection(MaintNotificationsAbstractConnection, ConnectionInterfa
             if callback:
                 callback(self)
 
+    def _connect_and_handshake(self, check_health: bool = True):
+        """Establish the socket connection and perform the initial handshake.
+
+        This whole process counts as a single connection attempt, so that
+        retryable errors raised during the handshake are retried together
+        with the socket connection.
+        """
+        sock = self._connect()
+        self._sock = sock
+        try:
+            if self.redis_connect_func is None:
+                # Use the default on_connect function
+                self.on_connect_check_health(check_health=check_health)
+            else:
+                # Use the passed function redis_connect_func
+                self.redis_connect_func(self)
+        except BaseException:
+            # clean up after any error in the handshake
+            self.disconnect()
+            raise
+        return sock
+
     @abstractmethod
     def _connect(self):
         pass
step 1
$ ls /repo && ls /repo/*/ | head -50
output · 67 lines
CHANGES
CONTRIBUTING.md
LICENSE
README.md
benchmarks
codecov.yml
dev_requirements.txt
docker-compose.yml
dockers
docs
doctests
pyproject.toml
55 more lines
redis
tasks.py
tests
util
whitelist.py
/repo/benchmarks/:
__init__.py
base.py
basic_operations.py
cluster_async.py
cluster_async_pipeline.py
command_packer_benchmark.py
socket_read_size.py

/repo/dockers/:
sentinel.conf

/repo/docs/:
Makefile
_static
_templates
advanced_features.rst
backoff.rst
clustering.rst
commands.rst
conf.py
connections.rst
examples
examples.rst
exceptions.rst
genindex.rst
images
index.rst
lock.rst
lua_scripting.rst
multi_database.rst
opentelemetry.rst
redismodules.rst
requirements.txt
resp3_features.rst
retry.rst

/repo/doctests/:
README.md
cmds_cnxmgmt.py
cmds_generic.py
cmds_hash.py
cmds_list.py
cmds_servermgmt.py
cmds_set.py
cmds_sorted_set.py
cmds_string.py
data
dt_bitfield.py
dt_bitmap.py
step 2
$ grep -n "retry" /repo/redis/connection.py | head -40
output · 26 lines
56:from .retry import Retry
664:        retry_on_timeout: bool = False,
665:        retry_on_error: Union[Iterable[Type[Exception]], object] = SENTINEL,
676:        retry: Union[Any, None] = None,
694:        To specify a retry policy for specific errors, first set
695:        `retry_on_error` to a list of the error/s to retry on, then set
696:        `retry` to a valid `Retry` object.
697:        To retry on TimeoutError, `retry_on_timeout` can also be set to `True`.
722:        self.retry_on_timeout = retry_on_timeout
723:        if retry_on_error is SENTINEL:
724:            retry_on_errors_list = []
726:            retry_on_errors_list = list(retry_on_error)
14 more lines
727:        if retry_on_timeout:
728:            # Add TimeoutError to the errors list to retry on
729:            retry_on_errors_list.append(TimeoutError)
730:        self.retry_on_error = retry_on_errors_list
731:        if
…[862 chars omitted]…
   def set_retry(self, retry: Retry):
2703:    def set_retry(self, retry: Retry) -> None:
2704:        self.connection_kwargs.update({"retry": retry})
2706:            conn.retry = retry
2708:            conn.retry = retry
2713:                conn.retry.call_with_retry(
2719:                conn.retry.call_with_retry(
2738:        Dummy functions, needs to be passed as error callback to retry object.
step 3
$ sed -n 840,960p /repo/redis/connection.py
output · 36 lines
def _get_parser(self) -> Union[_HiredisParser, _RESP3Parser, _RESP2Parser]:
        return self._parser

    def connect(self):
        "Connects to the Redis server if not already connected"
        self.connect_check_health(check_health=True)

    def connect_check_health(
        self, check_health: bool = True, retry_socket_connect: bool = True
    ):
        if self._sock:
            return
24 more lines
        try:
            if retry_socket_connect:
                sock = self.retry.call_with_retry(
                    lambda: self._connect(), lambda error: self.disconnect(error)
                )
            else:
                sock = self._connect()
        except socket.timeout:
            raise TimeoutError("Timeout connecting to server")
        except OSError as e:
            raise ConnectionError(self._error_message(e))

        self._sock = sock
        try:
            i
…[3668 chars omitted]…
h to it
        elif self.protocol not in [2, "2"]:
            if isinstance(self._parser, _RESP2Parser):
                self.set_parser(_RESP3Parser)
                # update cluster exception classes
                self._parser.EXCEPTION_CLASSES = parser.EXCEPTION_CLASSES
                self._parser.on_connect(self)
            self.send_command("HELLO", self.protocol, check_health=check_health)
step 4
$ cat /repo/redis/retry.py
output · 51 lines
import abc
import socket
from time import sleep
from typing import (
    TYPE_CHECKING,
    Any,
    Callable,
    Generic,
    Iterable,
    Optional,
    Tuple,
    Type,
39 more lines
    TypeVar,
)

from redis.exceptions import ConnectionError, TimeoutError

T = TypeVar("T")
E = TypeVar("E", bound=Exception, covariant=True)

if TYPE_CHECKING:
    from redis.backoff import AbstractBackoff


class AbstractRetry(Generic[E], abc.ABC):
    """Retry a specific number of times after a failure"""

    _supported_errors: Tuple[Type[E], ...]

    def __init__(
        self,
        backoff: "AbstractBackoff",
        retries: int,
        supported_errors: Tuple[Type[E], ...],
    ):
        """
        Initialize a `Retry` object with a `Backoff` object
        that retries a maximum of `retries` times.
        `retries` can be negative to retry forever.
        You can specify the types of supported er
…[2309 chars omitted]…
t self._supported_errors as error:
                if is_retryable and not is_retryable(error):
                    raise
                failures += 1
                fail(error)
                if self._retries >= 0 and failures > self._retries:
                    raise error
                backoff = self._backoff.compute(failures)
                if backoff > 0:
                    sleep(backoff)
step 5
$ sed -n 960,1010p /repo/redis/connection.py; grep -rn "retry_socket_connect\|connect_check_health" /repo/redis/*.py | grep -v connection.py
output · 35 lines
self.send_command("HELLO", self.protocol, check_health=check_health)
            self.handshake_metadata = self.read_response()
            if (
                self.handshake_metadata.get(b"proto") != self.protocol
                and self.handshake_metadata.get("proto") != self.protocol
            ):
                raise ConnectionError("Invalid RESP version")

        # Activate maintenance notifications for this connection
        # if enabled in the configuration
        # This is a no-op if maintenance notifications are not enabled
        self.activate_maint_notifications_handling_if_enabled(check_health=check_health)
23 more lines
        # if a client_name is given, set it
        if self.client_name:
            self.send_command(
                "CLIENT",
                "SETNAME",
                self.client_name,
                check_health=check_health,
            )
 
…[614 chars omitted]…
                "CLIENT",
                    "SETINFO",
                    "LIB-VER",
                    self.lib_version,
                    check_health=check_health,
                )
                self.read_response()
        except ResponseError:
            pass

/repo/redis/sentinel.py:43:        self.connect_check_health(
/repo/redis/sentinel.py:45:            retry_socket_connect=False,
step 6
$ grep -rn "handshake\|on_connect" /repo/tests/test_connection.py | head; grep -rln "retry_socket_connect" /repo/tests
step 7
$ grep -rln "call_with_retry\|Retry(" /repo/tests | head; ls /repo/tests
output · 61 lines
/repo/tests/conftest.py
/repo/tests/test_asyncio/conftest.py
/repo/tests/test_asyncio/test_cluster.py
/repo/tests/test_asyncio/test_cluster_transaction.py
/repo/tests/test_asyncio/test_connection.py
/repo/tests/test_asyncio/test_credentials.py
/repo/tests/test_asyncio/test_multidb/test_command_executor.py
/repo/tests/test_asyncio/test_retry.py
/repo/tests/test_asyncio/test_scenario/conftest.py
/repo/tests/test_asyncio/test_scenario/test_active_active.py
__init__.py
conftest.py
49 more lines
entraid_utils.py
helpers.py
mocks.py
ssl_utils.py
test_asyncio
test_auth
test_background.py
test_backoff.py
test_bloom.py
test_cache.py
test_cluster.py
test_cluster_transaction.py
test_command_parser.py
test_command_policies.py
test_commands.py
test_connect.py
test_connection.py
test_connection_pool.py
test_credentials.py
test_data_structure.py
test_encoding.py
test_event.py
test_function.py
test_hash.py
test_he
…[15 chars omitted]…
ttp
test_json.py
test_lock.py
test_maint_notifications.py
test_maint_notifications_handling.py
test_max_connections_error.py
test_monitor.py
test_multidb
test_multiprocessing.py
test_parsers
test_pipeline.py
test_pubsub.py
test_retry.py
test_scenario
test_scripting.py
test_search.py
test_sentinel.py
test_sentinel_managed_connection.py
test_ssl.py
test_timeseries.py
test_utils.py
test_vsets.py
testdata
step 8
$ grep -n "connect_check_health\|retry_socket_connect\|on_connect" /repo/redis/asyncio/connection.py | head -20
output · 15 lines
299:        await self.connect_check_health(check_health=True)
301:    async def connect_check_health(
302:        self, check_health: bool = True, retry_socket_connect: bool = True
307:            if retry_socket_connect:
324:                # Use the default on_connect function
325:                await self.on_connect_check_health(check_health=check_health)
334:            # clean up after any error in on_connect
368:    async def on_connect(self) -> None:
370:        await self.on_connect_check_health(check_health=True)
372:    async def on_connect_check_health(self, check_health: bool = True) -> None:
373:        self._parser.on_connect(self)
392:                self._parser.on_connect(self)
3 more lines
429:                self._parser.on_connect(self)
529:            await self.connect_check_health(check_health=False)
972:        await self.on_connect()
step 9
$ sed -n 296,345p /repo/redis/asyncio/connection.py
output · 31 lines
async def connect(self):
        """Connects to the Redis server if not already connected"""
        await self.connect_check_health(check_health=True)

    async def connect_check_health(
        self, check_health: bool = True, retry_socket_connect: bool = True
    ):
        if self.is_connected:
            return
        try:
            if retry_socket_connect:
                await self.retry.call_with_retry(
19 more lines
                    lambda: self._connect(), lambda error: self.disconnect()
                )
            else:
                await self._connect()
        except asyncio.CancelledError:
            raise  # in 3.7 and earlier, this is an Exception, not BaseException
        except (socket.timeout, asyncio.TimeoutError):
            raise TimeoutError("Timeout connecting to server")
        except OSError as e:
            raise ConnectionError(self._error_message(
…[717 chars omitted]…
   # run any user callbacks. right now the only internal callback
        # is for pubsub channel/pattern resubscription
        # first, remove any dead weakrefs
        self._connect_callbacks = [ref for ref in self._connect_callbacks if ref()]
        for ref in self._connect_callbacks:
            callback = ref()
            task = callback(self)
            if task and inspect.isawaitable(task):
step 10
$ grep -n "def disconnect\|def _connect\|def is_connected\|_is_connected" /repo/redis/asyncio/connection.py /repo/redis/connection.py | head -30
output · 20 lines
/repo/redis/asyncio/connection.py:263:    def is_connected(self):
/repo/redis/asyncio/connection.py:355:    async def _connect(self):
/repo/redis/asyncio/connection.py:480:    async def disconnect(self, nowait: bool = False) -> None:
/repo/redis/asyncio/connection.py:763:    def _connection_arguments(self) -> Mapping:
/repo/redis/asyncio/connection.py:766:    async def _connect(self):
/repo/redis/asyncio/connection.py:831:    def _connection_arguments(self) -> Mapping:
/repo/redis/asyncio/connection.py:967:    async def _connect(self):
/repo/redis/asyncio/connection.py:1259:    async def disconnect(self, inuse_connections: bool = True):
/repo/redis/connection.py:194:    def disconnect(self, *args):
/repo/redis/connection.py:378:    def disconnect(self, *args):
/repo/redis/connection.py:888:    def _connect(self):
/repo/redis/connection.py:1017:    def disconnect(self, *args):
8 more lines
/repo/re
…[214 chars omitted]…
self):
/repo/redis/connection.py:1853:    def _connect(self):
/repo/redis/connection.py:2004:    def disconnect(self, inuse_connections: bool = True):
/repo/redis/connection.py:2346:    def disconnect_free_connections(
/repo/redis/connection.py:2679:    def disconnect(self, inuse_connections: bool = True) -> None:
/repo/redis/connection.py:2959:    def disconnect(self, inuse_connections: bool = True):
step 11
$ sed -n 1010,1060p /repo/redis/connection.py; sed -n 475,500p /repo/redis/asyncio/connection.py
output · 43 lines
# if a database is specified, switch to it
        if self.db:
            self.send_command("SELECT", self.db, check_health=check_health)
            if str_if_bytes(self.read_response()) != "OK":
                raise ConnectionError("Invalid Database")

    def disconnect(self, *args):
        "Disconnects from the Redis server"
        self._parser.on_disconnect()

        conn_sock = self._sock
        self._sock = None
31 more lines
        # reset the reconnect flag
        self.reset_should_reconnect()
        if conn_sock is None:
            return

        if os.getpid() == self.pid:
            try:
                conn_sock.shutdown(socket.SHUT_RDWR)
            except (OSError, TypeError):
                pass

        try:
            conn_sock.close()
        except OSError:
            pass

    def mark_for_reconnect(self):
        self._should_reconnect = True

    def s
…[1383 chars omitted]…
ait:
                        await self._writer.wait_closed()  # type: ignore[union-attr]
                except OSError:
                    pass
                finally:
                    self._reader = None
                    self._writer = None
        except asyncio.TimeoutError:
            raise TimeoutError(
                f"Timed out closing connection after {self.socket_connect_timeout}"
step 12

Now let me apply the fix to both sync and async connections.

$ cd /repo && python - <<'PY'
output · 40 lines
"Connects to the Redis server if not already connected"
        self.connect_check_health(check_health=True)

    def connect_check_health(
        self, check_health: bool = True, retry_socket_connect: bool = True
    ):
        if self._sock:
            return
        try:
            if retry_socket_connect:
                self.retry.call_with_retry(
                    lambda: self._connect_and_handshake(check_health),
28 more lines
                    lambda error: self.disconnect(error),
                )
            else:
                self._connect_and_handshake(check_health)
        except socket.timeout:
            raise TimeoutError("Timeout connecting to server")
        except OSError as e:
            raise ConnectionError(self._error_message(e))
        except RedisError:
            # clean up after any error in the handshake
            self.disconnect()
            ra
…[924 chars omitted]…
onnect_check_health(check_health=check_health)
            else:
                # Use the passed function redis_connect_func
                self.redis_connect_func(self)
        except BaseException:
            # clean up after any error in the handshake
            self.disconnect()
            raise
        return sock

    @abstractmethod
    def _connect(self):
        pass

    @abstractmethod
step 13
$ cd /repo && python - <<'PY'
output · 39 lines
async def connect(self):
        """Connects to the Redis server if not already connected"""
        await self.connect_check_health(check_health=True)

    async def connect_check_health(
        self, check_health: bool = True, retry_socket_connect: bool = True
    ):
        if self.is_connected:
            return
        try:
            if retry_socket_connect:
                await self.retry.call_with_retry(
27 more lines
                    lambda: self._connect_and_handshake(check_health),
                    lambda error: self.disconnect(),
                )
            else:
                await self._connect_and_handshake(check_health)
        except asyncio.CancelledError:
            raise  # in 3.7 and earlier, this is an Exception, not BaseException
        except (socket.timeout, asyncio.TimeoutError):
            raise TimeoutError("Timeout connecting to server")
        ex
…[1676 chars omitted]…
                else self.redis_connect_func(self)
                )
        except BaseException:
            # clean up after any error in the handshake
            await self.disconnect()
            raise

    @abstractmethod
    async def _connect(self):
        pass

    @abstractmethod
    def _host_error(self) -> str:
        pass

    def _error_message(self, exception: BaseException) -> str:
step 14
$ cd /repo && cat > /tmp/repro.py <<'PY'
output · 26 lines
Traceback (most recent call last):
  File "/repo/redis/connection.py", line 891, in _connect_and_handshake
    self.on_connect_check_health(check_health=check_health)
  File "/tmp/repro.py", line 26, in failing_handshake
    raise ConnectionError("ECONNREFUSED during handshake")
redis.exceptions.ConnectionError: ECONNREFUSED during handshake

During handling of the above exception, another exception occurred:

Traceback (most recent call last):
  File "/tmp/repro.py", line 29, in <module>
    c.connect_check_health()
14 more lines
  File "/repo/redis/connection.py", line 855, in connect_check_health
    self.retry.call_with_retry(
  File "/repo/redis/retry.py", line 116, in call_with_retry
    return do()
           ^^^^
  File "/repo/redis/connection.py", line 856, in <lambda>
    lambda: self._connect_and_handshake(check_health),
            ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^
  File "/repo/redis/connection.py", line 897, in _connect_and_handshake
    self.disconnect()
  File "/repo/redis/connection.py", line 1044, in disconnect
    conn_sock.shutdown(socket.SHUT_RDWR)
    ^^^^^^^^^^^^^^^^^^
AttributeError: 'object' object has no attribute 'shutdown'
step 15
$ cd /repo && sed -i 's/return object()/return object.__new__(object)/' /tmp/repro.py && python /tmp/repro.py
output · 26 lines
Traceback (most recent call last):
  File "/repo/redis/connection.py", line 891, in _connect_and_handshake
    self.on_connect_check_health(check_health=check_health)
  File "/tmp/repro.py", line 26, in failing_handshake
    raise ConnectionError("ECONNREFUSED during handshake")
redis.exceptions.ConnectionError: ECONNREFUSED during handshake

During handling of the above exception, another exception occurred:

Traceback (most recent call last):
  File "/tmp/repro.py", line 29, in <module>
    c.connect_check_health()
14 more lines
  File "/repo/redis/connection.py", line 855, in connect_check_health
    self.retry.call_with_retry(
  File "/repo/redis/retry.py", line 116, in call_with_retry
    return do()
           ^^^^
  File "/repo/redis/connection.py", line 856, in <lambda>
    lambda: self._connect_and_handshake(check_health),
            ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^
  File "/repo/redis/connection.py", line 897, in _connect_and_handshake
    self.disconnect()
  File "/repo/redis/connection.py", line 1044, in disconnect
    conn_sock.shutdown(socket.SHUT_RDWR)
    ^^^^^^^^^^^^^^^^^^
AttributeError: 'object' object has no attribute 'shutdown'
step 16
$ cd /repo && sed -i 's/return object.__new__(object)/import socket as _s; return _s.socket()/' /tmp/repro.py && python /tmp/repro.py
output · 7 lines
got ConnectionError: ECONNREFUSED during handshake attempts: 3
got AuthenticationError: bad password attempts: 3
Traceback (most recent call last):
  File "/tmp/repro.py", line 45, in <module>
    assert attempts["n"] == 1, attempts
           ^^^^^^^^^^^^^^^^^^
AssertionError: {'n': 3}
step 17
$ cd /repo && grep -n "class .*Error" redis/exceptions.py
output · 32 lines
4:class RedisError(Exception):
8:class ConnectionError(RedisError):
12:class TimeoutError(RedisError):
16:class AuthenticationError(ConnectionError):
20:class AuthorizationError(ConnectionError):
24:class BusyLoadingError(ConnectionError):
28:class InvalidResponse(RedisError):
32:class ResponseError(RedisError):
36:class DataError(RedisError):
40:class PubSubError(RedisError):
44:class WatchError(RedisError):
48:class NoScriptError(ResponseError):
20 more lines
52:class OutOfMemoryError(ResponseError):
64:class ExecAbortError(ResponseError):
68:class ReadOnlyError(ResponseError):
72:class NoPermissionError(ResponseError):
76:class ModuleError(ResponseError):
80:class LockError(RedisError, ValueError):
83:    # NOTE: For backwards compatibility, this class derives from ValueError.
91:class LockNotOwnedError(LockError):
97:class ChildDeadlockedError(Exception):
103:class AuthenticationWrongNumberOfAr
…[128 chars omitted]…
skError(ResponseError):
171:class TryAgainError(ResponseError):
182:class ClusterCrossSlotError(ResponseError):
192:class MovedError(AskError):
202:class MasterDownError(ClusterDownError):
211:class SlotNotCoveredError(RedisClusterException):
223:class MaxConnectionsError(ConnectionError):
232:class CrossSlotTransactionError(RedisClusterException):
250:class ExternalAuthProviderError(ConnectionError):
step 18
$ cd /repo && sed -i 's/AuthenticationError("bad password")/ResponseError("nope")/; s/from redis.exceptions import ConnectionError, AuthenticationError/from redis.exceptions import ConnectionError, ResponseError/; s/except AuthenticationError/except ResponseError/; s/got AuthenticationError/got ResponseError/' /tmp/repro.py && python /tmp/repro.py
output · 8 lines
got ConnectionError: ECONNREFUSED during handshake attempts: 3
got ResponseError: nope attempts: 1
connect-only: [Errno 111] Connection refused attempts: 0
Traceback (most recent call last):
  File "/tmp/repro.py", line 57, in <module>
    assert attempts["n"] == 3, attempts
           ^^^^^^^^^^^^^^^^^^
AssertionError: {'n': 0}
step 19
$ cd /repo && python - <<'PY'
output · 8 lines
got ConnectionError: ECONNREFUSED during handshake attempts: 3
got ResponseError: nope attempts: 1
connect-only: [Errno 111] Connection refused attempts: 0
Traceback (most recent call last):
  File "/tmp/repro.py", line 57, in <module>
    assert attempts["n"] == 3, attempts
           ^^^^^^^^^^^^^^^^^^
AssertionError: {'n': 0}
step 20
$ cd /repo && cat /tmp/repro.py
output · 46 lines
from unittest.mock import patch
import redis
from redis.connection import Connection, AbstractConnection
from redis.exceptions import ConnectionError, ResponseError
from redis.retry import Retry
from redis.backoff import NoBackoff

attempts = {"n": 0}

class FakeConn(Connection):
    def __init__(self, **kw):
        super().__init__(**kw)
34 more lines
        self.retry = Retry(NoBackoff(), 2)
    def _connect(self):
        attempts["n"] += 1
        import socket as _s; return _s.socket()
    def _host_error(self):
        return ("localhost", 6379)
    def _error_message(self, e):
        return str(e)

# 1. retryable handshake failure -> retried, then ConnectionError
attempts["n"] = 0
c = FakeConn()
def failing_handshake(check_health=True):
    raise ConnectionError("ECONNREFUSED during handshake")
c.on_connect_check_health = failing_handshake
try:
    c.connect_check_health()
    print("FAIL
…[447 chars omitted]…
, attempts["n"])
assert attempts["n"] == 1, attempts

# 3. socket connect fails -> still retried
attempts["n"] = 0
c = FakeConn()
def fail_connect():
    raise ConnectionError("[Errno 111] Connection refused")
c._connect = fail_connect
try:
    c.connect_check_health()
except ConnectionError as e:
    print("connect-only:", e, "attempts:", attempts["n"])
assert attempts["n"] == 3, attempts
print("OK")
step 21
$ cd /repo && python - <<'PY'
output · 4 lines
got ConnectionError: ECONNREFUSED during handshake attempts: 3
got ResponseError: nope attempts: 1
connect-only: [Errno 111] Connection refused attempts: 3
OK
step 22
$ cd /repo && cat > /tmp/repro_a.py <<'PY'
output · 17 lines
retryable: 'object' object has no attribute 'close' 1
Traceback (most recent call last):
  File "/tmp/repro_a.py", line 45, in <module>
    asyncio.run(main())
  File "/usr/local/lib/python3.12/asyncio/runners.py", line 195, in run
    return runner.run(main)
           ^^^^^^^^^^^^^^^^
  File "/usr/local/lib/python3.12/asyncio/runners.py", line 118, in run
    return self._loop.run_until_complete(task)
           ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^
  File "/usr/local/lib/python3.12/asyncio/base_events.py", line 691, in run_until_complete
    return future.result()
5 more lines
           ^^^^^^^^^^^^^^^
  File "/tmp/repro_a.py", line 31, in main
    assert attempts["n"] == 3
           ^^^^^^^^^^^^^^^^^^
AssertionError
step 23
$ cd /repo && python - <<'PY'
output · 17 lines
retryable: [Errno 111] Connect call failed ('127.0.0.1', 1) 1
Traceback (most recent call last):
  File "/tmp/repro_a.py", line 46, in <module>
    asyncio.run(main())
  File "/usr/local/lib/python3.12/asyncio/runners.py", line 195, in run
    return runner.run(main)
           ^^^^^^^^^^^^^^^^
  File "/usr/local/lib/python3.12/asyncio/runners.py", line 118, in run
    return self._loop.run_until_complete(task)
           ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^
  File "/usr/local/lib/python3.12/asyncio/base_events.py", line 691, in run_until_complete
    return future.result()
5 more lines
           ^^^^^^^^^^^^^^^
  File "/tmp/repro_a.py", line 32, in main
    assert attempts["n"] == 3
           ^^^^^^^^^^^^^^^^^^
AssertionError
step 24
$ cd /repo && sed -n 795,860p redis/asyncio/connection.py
output · 41 lines
sock.setsockopt(socket.SOL_TCP, k, v)

            except (OSError, TypeError):
                # `socket_keepalive_options` might contain invalid options
                # causing an error. Do not leave the connection open.
                writer.close()
                raise

    def _host_error(self) -> str:
        return f"{self.host}:{self.port}"

29 more lines
class SSLConnection(Connection):
    """Manages SSL connections to and from the Redis server(s).
    This class extends the Connection class, adding SSL functionality, and making
    use of ssl.SSLContext (https://docs.python.org/3/library/ssl.html#ssl.SSLContext)
    """

    def __init__(
        self,
        ssl_keyfile: Optional[str] = None,
        ssl_certfile: Optional[str] = None,
        ssl_cert_reqs: Union[str, ssl.VerifyMode] = "required",
        ssl_include_verify_flags: Optional[List["ssl.Verify
…[964 chars omitted]…
  def _connection_arguments(self) -> Mapping:
        kwargs = super()._connection_arguments()
        kwargs["ssl"] = self.ssl_context.get()
        return kwargs

    @property
    def keyfile(self):
        return self.ssl_context.keyfile

    @property
    def certfile(self):
        return self.ssl_context.certfile

    @property
    def cert_reqs(self):
        return self.ssl_context.cert_reqs
step 25
$ cd /repo && sed -n 760,795p redis/asyncio/connection.py
output · 36 lines
socket_type: int = 0,
        **kwargs,
    ):
        self.host = host
        self.port = int(port)
        self.socket_keepalive = socket_keepalive
        self.socket_keepalive_options = socket_keepalive_options or {}
        self.socket_type = socket_type
        super().__init__(**kwargs)

    def repr_pieces(self):
        pieces = [("host", self.host), ("port", self.port), ("db", self.db)]
24 more lines
        if self.client_name:
            pieces.append(("client_name", self.client_name))
        return pieces

    def _connection_arguments(self) -> Mapping:
        return {"host": self.host, "port": self.port}

    async def _connect(self):
        """Create a TCP socket connection"""
        async with async_timeout(self.socket_connect_timeout):
            reader, writer = await asyncio.open_connection(
                **self._connection_arguments()
            )
        self.
…[90 chars omitted]…
nfo("socket")
        if sock:
            sock.setsockopt(socket.IPPROTO_TCP, socket.TCP_NODELAY, 1)
            try:
                # TCP_KEEPALIVE
                if self.socket_keepalive:
                    sock.setsockopt(socket.SOL_SOCKET, socket.SO_KEEPALIVE, 1)
                    for k, v in self.socket_keepalive_options.items():
                        sock.setsockopt(socket.SOL_TCP, k, v)
step 26
$ cd /repo && python - <<'PY'
output · 19 lines
retryable: name 'main_srv' is not defined 1
Traceback (most recent call last):
  File "/tmp/repro_a.py", line 53, in <module>
    asyncio.run(run())
  File "/usr/local/lib/python3.12/asyncio/runners.py", line 195, in run
    return runner.run(main)
           ^^^^^^^^^^^^^^^^
  File "/usr/local/lib/python3.12/asyncio/runners.py", line 118, in run
    return self._loop.run_until_complete(task)
           ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^
  File "/usr/local/lib/python3.12/asyncio/base_events.py", line 691, in run_until_complete
    return future.result()
7 more lines
           ^^^^^^^^^^^^^^^
  File "/tmp/repro_a.py", line 50, in run
    await main(srv)
  File "/tmp/repro_a.py", line 33, in main
    assert attempts["n"] == 3
           ^^^^^^^^^^^^^^^^^^
AssertionError
step 27
$ cd /repo && cat > /tmp/repro_a.py <<'PY'
output · 19 lines
retryable: boom 1
Traceback (most recent call last):
  File "/tmp/repro_a.py", line 53, in <module>
    asyncio.run(run())
  File "/usr/local/lib/python3.12/asyncio/runners.py", line 195, in run
    return runner.run(main)
           ^^^^^^^^^^^^^^^^
  File "/usr/local/lib/python3.12/asyncio/runners.py", line 118, in run
    return self._loop.run_until_complete(task)
           ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^
  File "/usr/local/lib/python3.12/asyncio/base_events.py", line 691, in run_until_complete
    return future.result()
7 more lines
           ^^^^^^^^^^^^^^^
  File "/tmp/repro_a.py", line 50, in run
    await main()
  File "/tmp/repro_a.py", line 33, in main
    assert attempts["n"] == 3
           ^^^^^^^^^^^^^^^^^^
AssertionError
step 28
$ cd /repo && sed -i 's/from redis.retry import Retry/from redis.asyncio.retry import RedisRetry as Retry/' /tmp/repro_a.py && python /tmp/repro_a.py
output · 4 lines
Traceback (most recent call last):
  File "/tmp/repro_a.py", line 4, in <module>
    from redis.asyncio.retry import RedisRetry as Retry
ImportError: cannot import name 'RedisRetry' from 'redis.asyncio.retry' (/repo/redis/asyncio/retry.py)
step 29
$ cd /repo && grep -n "class\|def call_with_retry" redis/asyncio/retry.py | head
output · 2 lines
13:class Retry(AbstractRetry[RedisError]):
37:    async def call_with_retry(
step 30
$ cd /repo && sed -i 's/from redis.retry import Retry/from redis.asyncio.retry import Retry/' /tmp/repro_a.py && python /tmp/repro_a.py
output · 4 lines
Traceback (most recent call last):
  File "/tmp/repro_a.py", line 4, in <module>
    from redis.asyncio.retry import RedisRetry as Retry
ImportError: cannot import name 'RedisRetry' from 'redis.asyncio.retry' (/repo/redis/asyncio/retry.py)
step 31
$ cd /repo && sed -i 's/from redis.asyncio.retry import RedisRetry as Retry/from redis.asyncio.retry import Retry/' /tmp/repro_a.py && python /tmp/repro_a.py
output · 3 lines
retryable: boom 3
non-retryable: nope 1
OK async
step 32
$ cd /repo && python -m pytest tests/test_connection.py tests/test_retry.py tests/test_connect.py tests/test_connection_pool.py tests/test_sentinel.py tests/test_sentinel_managed_connection.py -q -x 2>&1 | tail -5
output · 5 lines
warnings.warn(

-- Docs: https://docs.pytest.org/en/stable/how-to/capture-warnings.html
!!!!!!!!!!!!!!!!!!!!!!!!!! stopping after 1 failures !!!!!!!!!!!!!!!!!!!!!!!!!!!
1 failed, 141 passed, 1 skipped, 148 warnings in 2.46s
step 33
$ cd /repo && python -m pytest tests/test_connection.py tests/test_retry.py tests/test_connect.py tests/test_connection_pool.py tests/test_sentinel.py tests/test_sentinel_managed_connection.py -q 2>&1 | grep -E "FAILED|failed|passed" | tail -3
output · 3 lines
FAILED tests/test_sentinel.py::test_redis_master_usage - redis.sentinel.Maste...
FAILED tests/test_sentinel.py::test_sentinel_commands_with_strict_redis_client
18 failed, 233 passed, 2 skipped, 258 warnings in 13.00s
step 34
$ cd /repo && git stash && python -m pytest tests/test_sentinel.py -q 2>&1 | tail -2; git stash pop
output · 12 lines
Saved working directory and index state WIP on master: 36a9b89 base
-- Docs: https://docs.pytest.org/en/stable/how-to/capture-warnings.html
4 failed, 18 passed, 27 warnings in 4.10s
On branch master
Changes not staged for commit:
  (use "git add <file>..." to update what will be committed)
  (use "git restore <file>..." to discard changes in working directory)
	modified:   redis/asyncio/connection.py
	modified:   redis/connection.py

no changes added to commit (use "git add" and/or "git commit -a")
Dropped refs/stash@{0} (99b1d0227f1ef22ae839ae64c66f90537fadb39f)
step 35
$ cd /repo && git stash && python -m pytest tests/test_connection.py tests/test_retry.py tests/test_connect.py tests/test_connection_pool.py tests/test_sentinel.py tests/test_sentinel_managed_connection.py -q 2>&1 | grep -E "FAILED|failed|passed" | tail -20; git stash pop
output · 21 lines
Saved working directory and index state WIP on master: 36a9b89 base
FAILED tests/test_retry.py::TestRedisClientRetry::test_get_set_retry_object_for_cluster_client
FAILED tests/test_connect.py::test_tcp_ssl_connect[771] - OSError: No SSL cer...
FAILED tests/test_connect.py::test_tcp_ssl_connect[772] - OSError: No SSL cer...
FAILED tests/test_connect.py::test_tcp_ssl_tls12_custom_ciphers[AES256-SHA:DHE-RSA-AES256-SHA:AES128-SHA:DHE-RSA-AES128-SHA]
FAILED tests/test_connect.py::test_tcp_ssl_tls12_custom_ciphers[ECDHE-ECDSA-AES256-GCM-SHA384]
FAILED tests/test_connect.py::test_tcp_ssl_tls12_custom_ciphers[ECDHE-RSA-AES128-GCM-SHA256]
FAILED tests/test_connect.py::test_tcp_ssl_version_mismatch - OSError: No SSL...
FAILED tests/test_connection_pool.py::TestConnection::test_busy_loading_disconnects_socket
FAILED tests/test_connection_pool.py::TestConnection::test_busy_loading_from_pipeline_i
…[814 chars omitted]…
rnings in 13.84s
9 more lines
On branch master
Changes not staged for commit:
  (use "git add <file>..." to update what will be committed)
  (use "git restore <file>..." to discard changes in working directory)
	modified:   redis/asyncio/connection.py
	modified:   redis/connection.py

no changes added to commit (use "git add" and/or "git commit -a")
Dropped refs/stash@{0} (617f50ebfdc1fe2430e86b2d450649a2d5369dc9)
step 36
$ cd /repo && python -m pytest tests/test_asyncio/test_connection.py tests/test_asyncio/test_retry.py -q 2>&1 | tail -3
output · 2 lines
-- Docs: https://docs.pytest.org/en/stable/how-to/capture-warnings.html
97 passed, 2 skipped, 104 warnings in 1.69s
step 37
$ cd /repo && python -m pytest tests/test_maint_notifications.py tests/test_maint_notifications_handling.py tests/test_cluster.py -q 2>&1 | tail -3; git diff --stat
output · 5 lines
-- Docs: https://docs.pytest.org/en/stable/how-to/capture-warnings.html
127 failed, 210 passed, 352 warnings, 3 errors in 87.51s (0:01:27)
 redis/asyncio/connection.py | 47 +++++++++++++++++++++++++++++----------------
 redis/connection.py         | 40 +++++++++++++++++++++++++-------------
 2 files changed, 57 insertions(+), 30 deletions(-)
step 38
$ cd /repo && git stash && python -m pytest tests/test_maint_notifications.py tests/test_maint_notifications_handling.py tests/test_cluster.py -q 2>&1 | tail -2; git stash pop
output · 12 lines
Saved working directory and index state WIP on master: 36a9b89 base
-- Docs: https://docs.pytest.org/en/stable/how-to/capture-warnings.html
127 failed, 210 passed, 352 warnings, 3 errors in 78.69s (0:01:18)
On branch master
Changes not staged for commit:
  (use "git add <file>..." to update what will be committed)
  (use "git restore <file>..." to discard changes in working directory)
	modified:   redis/asyncio/connection.py
	modified:   redis/connection.py

no changes added to commit (use "git add" and/or "git commit -a")
Dropped refs/stash@{0} (4528c31c5804188489f792ccf006577eef217a54)
step 39
$ cd /repo && git diff
output · 33 lines
diff --git a/redis/asyncio/connection.py b/redis/asyncio/connection.py
index c5764f7..ef53c3d 100644
--- a/redis/asyncio/connection.py
+++ b/redis/asyncio/connection.py
@@ -306,34 +306,23 @@ class AbstractConnection:
         try:
             if retry_socket_connect:
                 await self.retry.call_with_retry(
-                    lambda: self._connect(), lambda error: self.disconnect()
+                    lambda: self._connect_and_handshake(check_health),
+                    lambda error: self.disconnect(),
                 )
21 more lines
             else:
-                await self._connect()
+                await self._connect_and_handshake(check_health)
         except asyncio.CancelledError:
             raise  # in 3.7 and earlier, this is an Exception, not BaseException
         except (socket.timeout, asyncio.TimeoutError):
             raise TimeoutError("Timeout connecting t
…[4501 chars omitted]…
self.on_connect_check_health(check_health=check_health)
+            else:
+                # Use the passed function redis_connect_func
+                self.redis_connect_func(self)
+        except BaseException:
+            # clean up after any error in the handshake
+            self.disconnect()
+            raise
+        return sock
+
     @abstractmethod
     def _connect(self):
         pass
step 40
$ cd /repo && sed -n 955,985p redis/asyncio/connection.py; grep -rn "on_connect()" redis/*.py | grep -v asyncio
output · 31 lines
if self.certfile and self.keyfile:
                context.load_cert_chain(certfile=self.certfile, keyfile=self.keyfile)
            if self.ca_certs or self.ca_data:
                context.load_verify_locations(cafile=self.ca_certs, cadata=self.ca_data)
            if self.min_version is not None:
                context.minimum_version = self.min_version
            if self.ciphers is not None:
                context.set_ciphers(self.ciphers)
            self.context = context
        return self.context

19 more lines
class UnixDomainSocketConnection(AbstractConnection):
    "Manages UDS communication to and from a Redis server"

    def __init__(self, *, path: str = "", **kwargs):
        self.path = path
        super().__init__(**kwargs)

    def repr_pieces(self) -> Iterable[Tuple[str, Union[str, int]]]:
        pieces = [("path", self.path), ("db", self.db)]
        if self.c
…[143 chars omitted]…
h async_timeout(self.socket_connect_timeout):
            reader, writer = await asyncio.open_unix_connection(path=self.path)
        self._reader = reader
        self._writer = writer
        await self.on_connect()
redis/cluster.py:671:        # inside the RedisCluster.on_connect() function
redis/cluster.py:782:        connection.on_connect()
redis/connection.py:1419:        self._conn.on_connect()
step 41
$ cd /repo && sed -n 770,795p redis/cluster.py; sed -n 1405,1425p redis/connection.py
output · 34 lines
if node.redis_connection:
                try:
                    node.redis_connection.connection_pool.disconnect()
                except OSError:
                    # Client was already disconnected. do nothing
                    pass

    def on_connect(self, connection):
        """
        Initialize the connection, authenticate and select a database and send
         READONLY if it is set during object initialization.
        """
22 more lines
        connection.on_connect()

        if self.read_from_replicas or self.load_balancing_strategy:
            # Sending READONLY command to server to configure connection as
            # readonly. Since each cluster node may change its server type due
            # to a failover, we should establish a READONLY connection
            # regardless of the server type. If this is a primary connection,
            # READONLY would not aff
…[620 chars omitted]…
ERSION) == 1
        ):
            raise ConnectionError(
                "To maximize compatibility with all Redis products, client-side caching is supported by Redis 7.4 or later"  # noqa: E501
            )

    def on_connect(self):
        self._conn.on_connect()

    def disconnect(self, *args):
        with self._cache_lock:
            self._cache.flush()
        self._conn.disconnect(*args)
step 42
$ cd /repo && grep -n "on_connect" redis/connection.py
output · 11 lines
190:    def on_connect(self):
890:                # Use the default on_connect function
891:                self.on_connect_check_health(check_health=check_health)
912:    def on_connect(self):
913:        self.on_connect_check_health(check_health=True)
915:    def on_connect_check_health(self, check_health: bool = True):
917:        self._parser.on_connect(self)
936:                self._parser.on_connect(self)
973:                self._parser.on_connect(self)
1418:    def on_connect(self):
1419:        self._conn.on_connect()
step 43
$ cd /repo && sed -n 960,1000p redis/asyncio/connection.py
output · 40 lines
context.minimum_version = self.min_version
            if self.ciphers is not None:
                context.set_ciphers(self.ciphers)
            self.context = context
        return self.context


class UnixDomainSocketConnection(AbstractConnection):
    "Manages UDS communication to and from a Redis server"

    def __init__(self, *, path: str = "", **kwargs):
        self.path = path
28 more lines
        super().__init__(**kwargs)

    def repr_pieces(self) -> Iterable[Tuple[str, Union[str, int]]]:
        pieces = [("path", self.path), ("db", self.db)]
        if self.client_name:
            pieces.append(("client_name", self.client_name))
        return pieces

    async def _connect(self):
        async with async_timeout(self.socket_connect_timeout):
            reader, writer = await asyncio.open_unix_connection(path=self.path)
        self._reader = reader
        self._writer = writer
        await self.on_connect()

    def _host_error(self) -> str:
        return self.path


FALSE_STRINGS = ("0", "F", "FALSE", "N", "NO")


def to_bool(value) -> Optional[bool]:
    if value is None or value == "":
        return None
    if isinstance(value, str) and value.upper() in FALSE_STRINGS:
        return False
    return bool(value)
step 44

All checks pass. The fix is complete:

$ done: {"notes": "Moved the handshake into the retried connection process for both sync and async connections: new _connect_and_handshake() helper (socket connect + on_connect/redis_connect_func with cleanup on failure) is now what retry.call_with_retry wraps in connect_check_health(). Retryable errors during the handshake now cause full
output · 3 lines
…[39 chars omitted]…
; non-retryable RedisErrors/cancellation propagate immediately after disconnect. Reproducers for sync and async pass; existing tests show no regressions vs baseline."}
submitted.