Source code for falkordb.asyncio.cluster

import socket

import redis as sync_redis  # type: ignore[import-not-found]
import redis.asyncio as redis  # type: ignore[import-not-found]
import redis.exceptions as redis_exceptions  # type: ignore[import-not-found]
from redis.asyncio.cluster import RedisCluster  # type: ignore[import-not-found]


# detect if a connection is a cluster
[docs] def Is_Cluster(conn: redis.Redis): pool = conn.connection_pool kwargs = pool.connection_kwargs.copy() # Check if the connection is using SSL and add it # this propery is not kept in the connection_kwargs kwargs["ssl"] = pool.connection_class is redis.SSLConnection # The async Unix-domain-socket pool stores the socket path under "path", # but the synchronous redis.Redis constructor expects "unix_socket_path". # Translate the key so the sync probe can be built for unix:// connections. if pool.connection_class is redis.UnixDomainSocketConnection: kwargs["unix_socket_path"] = kwargs.pop("path") # Create a synchronous Redis client with the same parameters # as the connection pool just to keep Is_Cluster synchronous info = sync_redis.Redis(**kwargs).info(section="server") return "redis_mode" in info and info["redis_mode"] == "cluster"
# create a cluster connection from a Redis connection
[docs] def Cluster_Conn( conn, ssl, cluster_error_retry_attempts=3, startup_nodes=None, require_full_coverage=False, reinitialize_steps=5, read_from_replicas=False, address_remap=None, ): connection_kwargs = conn.connection_pool.connection_kwargs host = connection_kwargs.pop("host") port = connection_kwargs.pop("port") username = connection_kwargs.pop("username") password = connection_kwargs.pop("password") retry = connection_kwargs.pop("retry", None) retry_on_error = connection_kwargs.pop( "retry_on_error", [ ConnectionRefusedError, ConnectionError, TimeoutError, socket.timeout, redis_exceptions.ConnectionError, ], ) return RedisCluster( host=host, port=port, username=username, password=password, decode_responses=True, ssl=ssl, retry=retry, retry_on_error=retry_on_error, require_full_coverage=require_full_coverage, reinitialize_steps=reinitialize_steps, read_from_replicas=read_from_replicas, address_remap=address_remap, startup_nodes=startup_nodes, cluster_error_retry_attempts=cluster_error_retry_attempts, )