
    ng\jv             	       4   S SK r S SKrS SKrS SKrS SKrS SKrS SKrS SKrS SKrS SK	J
r
Jr  S SKJr  S SKJr  S SKJr  S SKJr  S SKJrJrJrJrJrJrJrJrJrJrJrJrJrJ r J!r!J"r"  \(       a  S SK#J$r$  S S	K%J&r&J'r'J(r(J)r)J*r*J+r+  S S
K,J-r-J.r.  S SK/J0r0J1r1J2r2  S SK3J4r4  S SK5J6r6J7r7  S SK8J9r9J:r:J;r;J<r<J=r=  S SK>J?r?  S SK@JArAJBrB  S SKCJDrD  S SKEJFrF  S SKGJHrHJIrI  S SKJJKrKJLrLJMrM  S SKNJOrOJPrPJQrQJRrRJSrSJTrTJUrUJVrVJWrWJXrXJYrYJZrZJ[r[  S SK\J]r]J^r^  S SK_J`r`Jara  S SKbJcrcJdrd  S SKeJfrfJgrg  S SKhJiri  S SKjJkrkJlrl  S SKmJnrnJoroJprpJqrq  S SKrJsrsJtrtJuruJvrvJwrwJxrxJyryJzrzJ{r{J|r|J}r}J~r~JrJrJrJrJrJr  S SKJrJrJrJrJrJr  S SKJrJrJrJrJrJrJr  \(       a  S S KJrJrJr  OSrSrSr\GR0                  " \5      r\!" S!\S"\S"   \\S"4   5      r " S# S$\M\S\^5      r " S% S"5      r " S& S'5      r " S( S)\M\S\^5      r\O H=  r\GRE                  S*S+5      GRG                  5       r\S,:X  a  M-  \" \\\V" \5      5        M?      " S- S.5      r " S/ S0\
5      r " S1 S2\5      r " S3 S4\5      r " S5 S6\5      r " S7 S8\;5      rS9S:S;\pS<\ \   S=S4S> jr " S? S@\p5      r " SA SB\65      rg)C    N)ABCabstractmethod)defaultdict)copy)chain)
MethodType)TYPE_CHECKINGAnyCallable	CoroutineDequeDict	GeneratorListLiteralMappingOptionalSetTupleTypeTypeVarUnion!AsyncClusterKeyspaceNotifications)DEFAULT_RETRY_BASEDEFAULT_RETRY_CAPDEFAULT_RETRY_COUNTDEFAULT_SOCKET_CONNECT_TIMEOUTDEFAULT_SOCKET_READ_SIZEDEFAULT_SOCKET_TIMEOUT)AsyncCommandsParserEncoder)CommandPoliciesRequestPolicyResponsePolicy)get_response_callbacks)PubSubResponseCallbackT)AbstractConnection
ConnectionConnectionPoolInterfaceSSLConnection	parse_urlLock)record_error_countrecord_operation_duration)Retry)TokenInterface)ExponentialWithJitterBackoff	NoBackoff)EMPTY_RESPONSENEVER_DECODEAbstractRedis)PIPELINE_BLOCKED_COMMANDSPRIMARYREPLICASLOT_IDAbstractRedisClusterLoadBalancerLoadBalancingStrategyblock_pipeline_commandget_node_nameparse_cluster_shardsparse_cluster_shards_unified"parse_cluster_shards_with_str_keysparse_cluster_slots)READ_COMMANDSAsyncRedisClusterCommands)list_or_argsparse_pubsub_subscriptions)AsyncPolicyResolverAsyncStaticPolicyResolver)REDIS_CLUSTER_HASH_SLOTSkey_slot)CredentialProvider)
DriverInforesolve_driver_info)#AfterAsyncClusterInstantiationEvent AsyncAfterSlotsCacheRefreshEventAsyncEventListenerInterfaceEventDispatcher)AskErrorBusyLoadingErrorClusterDownErrorClusterErrorConnectionErrorCrossSlotTransactionError	DataErrorExecAbortErrorInvalidPipelineStackMaxConnectionsError
MovedErrorRedisClusterException
RedisErrorResponseErrorSlotNotCoveredErrorTimeoutErrorTryAgainError
WatchError)AnyKeyTChannelT
EncodableTKeyTPubSubHandlerSubscription)SENTINELSSL_AVAILABLEdeprecated_argsdeprecated_functionsafe_strstr_if_bytestruncate_text)
TLSVersionVerifyFlags
VerifyModeTargetNodesTClusterNodec            b       @   \ rS rSr% Sr\S\S\SS 4S j5       rSr	\
S   \S'   S	r\" S
/SSS9\" S/SSS9\" SS/SS9SSSSSSSS\SSSSSSSSS\\\SSSS\\\S\SSSSSSSSSSSSSSS\" 5       4.S\S-  S\\-  S \S!   S-  S"\S
\S#\S-  S$\S%\S\S&\S'\S-  S(\\\      S-  S)\\-  S*\S-  S+\S-  S,\S-  S-\S-  S.\S-  S\\-  S-  S\\-  S-  S/\\-  S-  S0\S1\S2\S3\S4\S-  S5\S-  S6\S7\S8\\\\ -  4   \-  S-  S9\S:\S-  S;\S-  S<S=S>\S?   S-  S@\S?   S-  SA\S-  SB\SC\S-  SDSESF\S-  SG\S-  SH\SI\!\"\\4   /\"\\4   4   S-  SJ\#S-  SK\$SS4^SL jj5       5       5       r%  SSM\&\\"\\4         SN\&\   SS 4SO jjr'SSP jr(\)" SQSRSSST9SSU j5       r*SSV jr+S\4SW jr,S\4SX jr-SY r.S\/\SS 4   4SZ jr0S[r1\2Rf                  \4Rj                  4S\\S]\SS4S^ jjr6S_\7SS4S` jr8S\S!   4Sa jr9S\S!   4Sb jr:S\S!   4Sc jr;SSd jr<SSe jr=SSg jr>   SS\&\   S\&\   Sh\&\   S\&S!   4Si jjr? SSj\Sk\S\&S!   4Sl jjr@Sm rASSn jrBSo\4Sp jrCS\&\DS!      4Sq jrESj\FS\4Sr jrGS\H4Ss jrIS\J\\&\   4   4St jrKS'\SS4Su jrLSo\Sv\MSS4Sw jrNSSx.So\Sy\Sz\OS{\&\   S\S!   4
S| jjrPSo\Sy\S\4S} jrQS~\S\4S jrRS~\S\S!   4S jrS  SS\S_\T\7S!4   S\S\&\   4S jjrU SS\S\S_\T\7S!4   S\&\   4S jjrVSy\FS\S\4S jrWSS!Sy\T\X\F4   S\S\4S jrY SS\&\   S\&\   SS4S jjrZ   SSf\&S!   S\&\   S\&\   S\SS4
S jjr[  SS\T\\ S4   S\SS4S jjr\       SS\XS\&\   S\S\S\&\   S\&\\]      S\S\S\]4S jjr^S\_SS\4   4S jr`Srag)RedisCluster   a  
Create a new RedisCluster client.

Pass one of parameters:

  - `host` & `port`
  - `startup_nodes`

| Use ``await`` :meth:`initialize` to find cluster nodes & create connections.
| Use ``await`` :meth:`close` to disconnect connections & close client.

Many commands support the target_nodes kwarg. It can be one of the
:attr:`NODE_FLAGS`:

  - :attr:`PRIMARIES`
  - :attr:`REPLICAS`
  - :attr:`ALL_NODES`
  - :attr:`RANDOM`
  - :attr:`DEFAULT_NODE`

Note: This client is not thread/process/fork safe.

:param host:
    | Can be used to point to a startup node
:param port:
    | Port used if **host** is provided
:param startup_nodes:
    | :class:`~.ClusterNode` to used as a startup node
:param require_full_coverage:
    | When set to ``False``: the client will not require a full coverage of
      the slots. However, if not all slots are covered, and at least one node
      has ``cluster-require-full-coverage`` set to ``yes``, the server will throw
      a :class:`~.ClusterDownError` for some key-based commands.
    | When set to ``True``: all slots must be covered to construct the cluster
      client. If not all slots are covered, :class:`~.RedisClusterException` will be
      thrown.
    | See:
      https://redis.io/docs/manual/scaling/#redis-cluster-configuration-parameters
:param read_from_replicas:
    | @deprecated - please use load_balancing_strategy instead
    | Enable read from replicas in READONLY mode.
      When set to true, read commands will be assigned between the primary and
      its replications in a Round-Robin manner.
      The data read from replicas is eventually consistent with the data in primary nodes.
:param load_balancing_strategy:
    | Enable read from replicas in READONLY mode and defines the load balancing
      strategy that will be used for cluster node selection.
      The data read from replicas is eventually consistent with the data in primary nodes.
:param dynamic_startup_nodes:
    | Set the RedisCluster's startup nodes to all the discovered nodes.
      If true (default value), the cluster's discovered nodes will be used to
      determine the cluster nodes-slots mapping in the next topology refresh.
      It will remove the initial passed startup nodes if their endpoints aren't
      listed in the CLUSTER SLOTS output.
      If you use dynamic DNS endpoints for startup nodes but CLUSTER SLOTS lists
      specific IP addresses, it is best to set it to false.
:param reinitialize_steps:
    | Specifies the number of MOVED errors that need to occur before reinitializing
      the whole cluster topology. If a MOVED error occurs and the cluster does not
      need to be reinitialized on this current error handling, only the MOVED slot
      will be patched with the redirected node.
      To reinitialize the cluster on every MOVED error, set reinitialize_steps to 1.
      To avoid reinitializing the cluster on moved errors, set reinitialize_steps to
      0.
:param cluster_error_retry_attempts:
    | @deprecated - Please configure the 'retry' object instead
      In case 'retry' object is set - this argument is ignored!

      Number of times to retry before raising an error when :class:`~.TimeoutError`,
      :class:`~.ConnectionError`, :class:`~.SlotNotCoveredError`
      or :class:`~.ClusterDownError` are encountered
:param retry:
    | A retry object that defines the retry strategy and the number of
      retries for the cluster client.
      In current implementation for the cluster client (starting form redis-py version 6.0.0)
      the retry object is not yet fully utilized, instead it is used just to determine
      the number of retries for the cluster client.
      In the future releases the retry object will be used to handle the cluster client retries!
:param max_connections:
    | Maximum number of connections per node. If there are no free connections & the
      maximum number of connections are already created, a
      :class:`~.MaxConnectionsError` is raised.
:param socket_keepalive:
    | If ``True``, TCP keepalive is enabled for TCP socket connections.
:param socket_keepalive_options:
    | Mapping of TCP keepalive socket option constants to values, for
      example ``{socket.TCP_KEEPIDLE: 30}``. If left unspecified, redis-py
      uses TCP keepalive defaults when ``socket_keepalive`` is enabled:
      idle 30 seconds, interval 5 seconds, and 3 probes.
      Platform-specific options that are not available are skipped.
      Pass ``None`` or ``{}`` to avoid setting additional TCP keepalive
      options.
:param address_remap:
    | An optional callable which, when provided with an internal network
      address of a node, e.g. a `(host, port)` tuple, will return the address
      where the node is reachable.  This can be used to map the addresses at
      which the nodes _think_ they are, to addresses at which a client may
      reach them, such as when they sit behind a proxy.

| Rest of the arguments will be passed to the
  :class:`~redis.asyncio.connection.Connection` instances when created

:raises RedisClusterException:
    if any arguments are invalid or unknown. Eg:

    - `db` != 0 or None
    - `path` argument for unix socket connection
    - none of the `host`/`port` & `startup_nodes` were provided

urlkwargsreturnc                     UR                  [        U5      5        UR                  SS5      [        L a  SUS'   U " S0 UD6$ )aI  
Return a Redis client object configured from the given URL.

For example::

    redis://[[username]:[password]]@localhost:6379/0
    rediss://[[username]:[password]]@localhost:6379/0

Three URL schemes are supported:

- `redis://` creates a TCP socket connection. See more at:
  <https://www.iana.org/assignments/uri-schemes/prov/redis>
- `rediss://` creates a SSL wrapped TCP socket connection. See more at:
  <https://www.iana.org/assignments/uri-schemes/prov/rediss>

The username, password, hostname, path and all querystring values are passed
through ``urllib.parse.unquote`` in order to replace any percent-encoded values
with their corresponding characters.

All querystring options are cast to their appropriate Python types. Boolean
arguments can be specified with string values "True"/"False" or "Yes"/"No".
Values that cannot be properly cast cause a ``ValueError`` to be raised. Once
parsed, the querystring arguments and keyword arguments are passed to
:class:`~redis.asyncio.connection.Connection` when created.
In the case of conflicting arguments, querystring arguments are used.
connection_classNTssl )updater-   popr,   )clsr|   r}   s      P/home/edenadmin/noVNC/venv/lib/python3.13/site-packages/redis/asyncio/cluster.pyfrom_urlRedisCluster.from_url  s=    8 	in%::($/=@ F5M}V}    T_is_async_client)_initialize_lockretrycommand_flagscommands_parserconnection_kwargsencoder
node_flagsnodes_managerread_from_replicasreinitialize_counterreinitialize_stepsresponse_callbacksresult_callbacksr   z6Please configure the 'load_balancing_strategy' insteadz5.3.0)args_to_warnreasonversioncluster_error_retry_attemptsz+Please configure the 'retry' object insteadz6.0.0lib_namelib_versionzbUse 'driver_info' parameter instead. lib_name and lib_version will be removed in a future version.)r   r   Ni  F   d   r   utf-8strictrequiredhostportstartup_nodesrx   require_full_coverageload_balancing_strategydynamic_startup_nodesr   max_connectionsr   retry_on_errordbpathcredential_providerusernamepasswordclient_namedriver_infoencodingencoding_errorsdecode_responseshealth_check_intervalsocket_timeoutsocket_connect_timeoutsocket_read_sizesocket_keepalivesocket_keepalive_optionsr   ssl_ca_certsssl_ca_datassl_cert_reqszstr | VerifyModessl_include_verify_flagsru   ssl_exclude_verify_flagsssl_certfilessl_check_hostnamessl_keyfilessl_min_versionzTLSVersion | Nonessl_ciphersprotocollegacy_responsesaddress_remapevent_dispatcherpolicy_resolverc/                 6
  ^  U(       a  [        S5      eU(       a  [        S5      eU(       a  U(       d  U(       d  [        S5      e[        UUU5      n/0 SU
_S[        _SU_SU_SU_S	U_S
U/_SU_SU_SU_SU_SU_SU_SU_SU_SU_SU*_SU+0En0U(       a!  U0R                  [        U U!U"U#U$U%U&U'U(U)S.5        U(       d  U(       a  T R
                  U0S'   U(       a  UT l        O[        [        [        [        S9U	S9T l        U(       a  T R                  R                  U5        [        U0R                  S5      U0R                  SS5      S9U0S'   U0R                  SS5      (       d  [        U0S   S'   O+U0R                  S5      c  [        U0S   S'   O[         U0S   S'   U0T l        U(       aH  / n1U H=  n2U1R%                  ['        U2R(                  U2R*                  40 T R"                  D65        M?     U1nO/ nU(       a,  U(       a%  UR%                  ['        X40 T R"                  D65        U-c  [-        5       T l        OU-T l        UT l        [3        UUU0UU,T R.                  S9T l        [7        UUU5      T l        UT l        UT l        UT l        ST l         T RB                  RD                  [F        RH                  T RB                  RJ                  [F        RL                  T RB                  RN                  [F        RN                  T RB                  RP                  [F        RR                  T RB                  RT                  [F        RT                  [V        [F        RX                  0T l-        [F        RH                  U 4S  j[F        RX                  T R\                  [F        RT                  U 4S! j[F        RL                  T R^                  [F        RN                  T R`                  [F        RR                  T Rb                  [F        Rd                  T Rf                  [h        RH                  S" [h        RX                  S# 0	T l5        U.T l6        [o        5       T l8        S T l9        T RB                  Rt                  Rw                  5       T l<        T RB                  Rz                  Rw                  5       T l>        U0S   T l?        T RB                  R                  Rw                  5       T lA        S$ T R                  S%'   ST lB        S T lC        ST lD        [        R                  " 5       T lG        g )&Nz/Argument 'db' must be 0 or None in cluster modez3Unix domain socket is not supported in cluster modea1  RedisCluster requires at least one node to discover the cluster.
Please provide one of the following or use RedisCluster.from_url:
   - host and port: RedisCluster(host="localhost", port=6379)
   - startup_nodes: RedisCluster(startup_nodes=[ClusterNode("localhost", 6379), ClusterNode("localhost", 6380)])r   r   r   r   r   r   r   r   r   r   r   r   r   r   r   r   r   r   )r   r   r   r   r   r   r   r   r   r   r   redis_connect_func)basecap)backoffretriesT)user_protocolr   r   zCLUSTER SHARDS)r   r   r   r   c                 (   > TR                  U 5      /$ N)get_random_primary_or_all_nodes)command_nameselfs    r   <lambda>'RedisCluster.__init__.<locals>.<lambda>
  s    44\BAr   c                  &   > T R                  5       /$ r   )get_default_noder   s   r   r   r     s    1F1F1H0Ir   c                     U $ r   r   ress    r   r   r     s    r   c                     U $ r   r   r   s    r   r   r     s    cr   c                 N    [        [        UR                  5       5      S   40 UD6$ Nr   )rE   listvalues)cmdr   r}   s      r   r   r     s%    ':SZZ\"1%()/(r   CLUSTER SLOTS)Hr`   rP   r*   r   r,   
on_connectr   r2   r4   r   r   update_supported_errorsr&   getrC   rD   rB   r   appendrx   r   r   rT   _event_dispatcherr   NodesManagerr   r"   r   r   r   r   r   	__class__RANDOMr$   DEFAULT_KEYLESS	PRIMARIES
ALL_SHARDS	ALL_NODESREPLICASALL_REPLICASDEFAULT_NODEr<   DEFAULT_KEYED_command_flags_mappingget_nodes_from_slotget_primaries	get_nodesget_replicasSPECIALget_special_nodesr%   _policies_callback_mapping_policy_resolverr!   r   _aggregate_nodes
NODE_FLAGSr   r   COMMAND_FLAGSr   r   RESULT_CALLBACKSr   r   r   _usage_counterasyncior/   _usage_lock)3r   r   r   r   r   r   r   r   r   r   r   r   r   r   r   r   r   r   r   r   r   r   r   r   r   r   r   r   r   r   r   r   r   r   r   r   r   r   r   r   r   r   r   r   r   r   r   computed_driver_infor}   passed_nodesnodes3   `                                                  r   __init__RedisCluster.__init__7  s   N 'A  'E  D-'S   3;+V"
"

"
 "#6	"

 "
 "
 ;"
 /"
 "
 "
  0"
 $%:"
 %&<"
   0!"
" '(@#"
$  0%"
& n'"
( )"
*  0+"
0 MM(5$0#.%20H0H$0*<#.'6#.  !8+/??F'(DJ4+1B 5	DJ JJ..~>'= **Z0#ZZ(:DA(
#$ zz,d33, '()9: ZZ
#+2 '()9: >RF'()9:!'L%##		499O8N8NO & )MMD  T!R4;Q;Q!RS#%4%6D"%5D"*)!"7'!33
 x:JK"4'>$"4$%! NN!!=#@#@NN$$m&>&>NN$$m&=&=NN##]%?%?NN'')C)C]00X
# )) , '')A)A&&(I$$d&8&8##T^^&&(9(9!!4#9#9**O((/
 	'  !024 $..3388:!^^99>>@"()="> $ ? ? D D F 	o.  -1

  "<<>r   additional_startup_nodes_infolast_failed_node_namec                   #    U R                   (       a  U R                  (       d  [        R                  " 5       U l        U R                   ISh  vN   U R                   (       aa   U R                  R                  UUS9I Sh  vN   U R                  R                  U R                  R                  5      I Sh  vN   SU l         SSS5      ISh  vN   U $ U $  N NX N#! [         aI    U R                  R                  5       I Sh  vN    U R                  R                  S5      I Sh  vN    e f = f Nd! , ISh  vN  (       d  f       U $ = f7f)zJGet all nodes from startup nodes & creates connections if not initialized.N)r  r  Fr   )
r   r   r  r/   r   
initializer   default_nodeBaseExceptionaclose)r   r  r  s      r   r  RedisCluster.initialize-  s     ::$\\^
zzz##"00;;:W2G <    #22== ..;;   ,1( "z t " ) "0077999"0077HHH "zz s   AECED8&C C6C :C;C ED6	EC C  'D3D
#D3+D.,D33D86E8E>E?EEc                   #    U R                   (       d  U R                  (       d  [        R                  " 5       U l        U R                   ISh  vN   U R                   (       dL  SU l         U R                  R                  5       I Sh  vN   U R                  R                  S5      I Sh  vN   SSS5      ISh  vN   gg Ns N; N N! , ISh  vN  (       d  f       g= f7f)z.Close all connections & client if initialized.NTr   )r   r   r  r/   r   r  r   s    r   r  RedisCluster.acloseG  s     ::$\\^
zzz'''+D$,,33555,,33ODDD	 "zz   " 6D	 "zzzsl   AC%CC%6C	C
"C,C-C1C%<C	=C%CC	C%C"CC"C%z5.0.0zUse aclose() insteadclose)r   r   namec                 @   #    U R                  5       I Sh  vN   g N7f)z.alias for aclose() for backwards compatibilityN)r  r   s    r   r  RedisCluster.closeR  s      kkm   c                    #    U R                  5       I Sh  vN    U R                  5       I Sh  vN $  N N! [         a    U R                  5       I Sh  vN    e f = f7f)z
Async context manager entry. Increments a usage counter so that the
connection pool is only closed (via aclose()) when no context is using
the client.
N)_increment_usager  	Exception_decrement_usager   s    r   
__aenter__RedisCluster.__aenter__W  s[      ##%%%	*** 	& + 	'')))	s:   A 4A 8 68 A 8 AAAA c                    #    U R                    ISh  vN   U =R                  S-  sl        U R                  sSSS5      ISh  vN   $  N6 N! , ISh  vN  (       d  f       g= f7f)zu
Helper coroutine to increment the usage counter while holding the lock.
Returns the new value of the usage counter.
N   r  r  r   s    r   r  RedisCluster._increment_usagef  A     
 ###1$&& $#####C   A(A
A(!AA(AA(A(A%AA%!A(c                    #    U R                    ISh  vN   U =R                  S-  sl        U R                  sSSS5      ISh  vN   $  N6 N! , ISh  vN  (       d  f       g= f7f)zu
Helper coroutine to decrement the usage counter while holding the lock.
Returns the new value of the usage counter.
Nr!  r"  r   s    r   r  RedisCluster._decrement_usageo  r$  r%  c                    #    [         R                  " U R                  5       5      I Sh  vN nUS:X  a-  [         R                  " U R                  5       5      I Sh  vN   gg N8 N7f)z
Async context manager exit. Decrements a usage counter. If this is the
last exit (counter becomes zero), the client closes its connection pool.
Nr   )r  shieldr  r  )r   exc_type	exc_value	tracebackcurrent_usages        r   	__aexit__RedisCluster.__aexit__x  sQ     
 &nnT-B-B-DEEA..///  F 0s!   (A'A#1A'A%A'%A'c                 >    U R                  5       R                  5       $ r   r  	__await__r   s    r   r2  RedisCluster.__await__       **,,r   zUnclosed RedisCluster client_warn_grlc                     [        U S5      (       aT  U R                  (       dB  U" U R                   SU < 3[        U S9   X R                  S.nU" 5       R	                  U5        g g g ! [
         a     g f = f)Nr    sourceclientmessage)hasattrr   _DEL_MESSAGEResourceWarningcall_exception_handlerRuntimeError)r   r5  r6  contexts       r   __del__RedisCluster.__del__  sw    
 4''0@0@T&&'q1?4P%)6G6GH--g6	 1A'
   s    $A' '
A43A4
connectionc                    #    UR                  5       I S h  vN   UR                  S5      I S h  vN   [        UR                  5       I S h  vN 5      S:w  a  [	        S5      eg  NN N7 N7f)NREADONLYOKzREADONLY command failed)r   send_commandrr   read_responserY   r   rF  s     r   r   RedisCluster.on_connect  se     ##%%% %%j111j66889TA!";<< B 	& 	28s1   A+A%A+A'A+A)A+'A+)A+c                 \    [        U R                  R                  R                  5       5      $ )zGet all nodes of the cluster.)r   r   nodes_cacher   r   s    r   r   RedisCluster.get_nodes  s"    D&&2299;<<r   c                 @    U R                   R                  [        5      $ )z%Get the primary nodes of the cluster.)r   get_nodes_by_server_typer:   r   s    r   r   RedisCluster.get_primaries      !!::7CCr   c                 @    U R                   R                  [        5      $ )z%Get the replica nodes of the cluster.)r   rR  r;   r   s    r   r   RedisCluster.get_replicas  rT  r   c                     [         R                  " [        U R                  R                  R                  5       5      5      $ )z!Get a random node of the cluster.)randomchoicer   r   rO  r   r   s    r   get_random_nodeRedisCluster.get_random_node  s+    }}T$"4"4"@"@"G"G"IJKKr   c                 .    U R                   R                  $ )z#Get the default node of the client.)r   r  r   s    r   r   RedisCluster.get_default_node  s    !!...r   r  c                     U(       a  U R                  UR                  S9(       d  [        S5      eXR                  l        g)zn
Set the default node of the client.

:raises DataError: if None is passed or node does not exist in cluster.
	node_namez1The requested node does not exist in the cluster.N)get_noder  r[   r   r  r   r  s     r   set_default_nodeRedisCluster.set_default_node  s0     4==499==OPP*.'r   r`  c                 :    U R                   R                  XU5      $ )z&Get node by (host, port) or node_name.)r   ra  r   r   r   r`  s       r   ra  RedisCluster.get_node  s     !!**4yAAr   keyreplicac                    U R                  U5      nU R                  R                  R                  U5      nU(       d  [	        SU S35      eU(       a-  [        U R                  R                  U   5      S:  a  gSnXE   $ SnXE   $ )a  
Get the cluster node corresponding to the provided key.

:param key:
:param replica:
    | Indicates if a replica should be returned
    |
      None will returned if no replica holds this key

:raises SlotNotCoveredError: if the key is not covered by any slot.
Slot "z " is not covered by the cluster.   Nr!  r   )keyslotr   slots_cacher   rc   len)r   rh  ri  slot
slot_cachenode_idxs         r   get_node_from_keyRedisCluster.get_node_from_key  s     ||C ''3377=
%tf4T&UVV4%%11$781<H ## H##r   c                 x    U R                   (       a  U[        ;   a  U R                  5       $ U R                  5       $ )z?
Returns random primary or all nodes depends on READONLY mode.
)r   rF   rZ  get_random_primary_node)r   r   s     r   r   ,RedisCluster.get_random_primary_or_all_nodes  s2     ""|}'D''))++--r   c                 J    [         R                  " U R                  5       5      $ )z
Returns a random primary node
)rX  rY  r   r   s    r   rv  $RedisCluster.get_random_primary_node  s     }}T//122r   commandc                    #    U R                   R                  U R                  " U/UQ76 I Sh  vN U R                  =(       a	    U[        ;   U[        ;   a  U R
                  5      /$ S5      /$  N@7f)z>
Returns a list of nodes that hold the specified keys' slots.
N)r   get_node_from_slot_determine_slotr   rF   r   )r   rz  argss      r   r    RedisCluster.get_nodes_from_slot  sw      11**7:T::''DG},D07=0H,,
 	
 OS
 	
:s   ,A1A/AA1c                 R    U R                   (       d  [        S5      eU R                   $ )z=
Returns a list of nodes for commands with a special policy.
z6Cannot execute FT.CURSOR commands without FT.AGGREGATE)r   r`   r   s    r   r   RedisCluster.get_special_nodes   s+     $$'H  $$$r   c                 J    [        U R                  R                  U5      5      $ )zk
Find the keyslot for a given key.

See: https://redis.io/docs/manual/scaling/#redis-cluster-data-sharding
)rM   r   encode)r   rh  s     r   rm  RedisCluster.keyslot  s     ++C011r   c                     U R                   $ )z%Get the encoder object of the client.)r   r   s    r   get_encoderRedisCluster.get_encoder  s    ||r   c                     U R                   $ )zGGet the kwargs passed to :class:`~redis.asyncio.connection.Connection`.)r   r   s    r   get_connection_kwargs"RedisCluster.get_connection_kwargs  s    %%%r   c                     Xl         g r   )r   r   r   s     r   	set_retryRedisCluster.set_retry  s    
r   callbackc                      X R                   U'   g)zSet a custom response callback.N)r   r   rz  r  s      r   set_response_callback"RedisCluster.set_response_callback  s    +3(r   )	node_flagr~  request_policyr  c                x  #    U(       d  U R                   R                  U5      nX0R                  ;   a  U R                  U   nU R                  U   nU[        R
                  :X  a  U" U/UQ76 I S h  vN nO$U[        R                  :X  a	  U" U5      nOU" 5       nUR                  5       S:X  a  X`l        U$  NE7f)Nzft.aggregate)	r   r   r   r   r$   r   r   lowerr   )r   rz  r  r  r~  policy_callbacknodess          r   _determine_nodesRedisCluster._determine_nodes"  s      **..w7I333!88CN99.I]888)'9D99E}<<<#G,E#%E==?n,$)! :s   A0B:2B83AB:c                   #    U R                   R                  U5      [        :X  a  [        US   5      $ UR	                  5       S;   aX  [        U5      S:  a  [        SU/UQ7 35      eUSS[        US   5      -    nU(       d  [        R                  " S[        5      $ OiU R                  R                  " U/UQ76 I S h  vN nU(       d=  UR	                  5       S;   a  [        R                  " S[        5      $ [        SU 35      e[        U5      S:X  a  U R                  US   5      $ U Vs1 s H  o@R                  U5      iM     nn[        U5      S:w  a  [        U S35      eUR                  5       $  Ns  snf 7f)	Nr   )EVALEVALSHArl  zInvalid args in command: r!  )FCALLFCALL_ROzNo way to dispatch this command to Redis Cluster. Missing key.
You can execute the command by specifying target nodes.
Command: z) - all keys must map to the same key slot)r   r   r<   intupperro  r`   rX  	randrangerL   r   get_keysrm  r   )r   rz  r~  keysrh  slotss         r   r}  RedisCluster._determine_slot@  s|    !!'*g5tAw< ==?114y1}+/$/?@  ADG,-D ''+CDD  --66wFFFD ==?&;;!++A/GHH+//3f6  t9><<Q(( /33dsc"d3u:?')DE  yy{1 G$ 4s%   B>F  E9A/F 0E;
0F ;F target_nodesc                 L    [        U[        5      =(       a    XR                  ;   $ r   )
isinstancestrr   )r   r  s     r   _is_node_flagRedisCluster._is_node_flags  s    ,,P1PPr   c                     [        U[        5      (       a  UnU$ [        U[        5      (       a  U/nU$ [        U[        5      (       a  [        UR	                  5       5      nU$ [        S[        U5       35      e)Nztarget_nodes type can be one of the following: node_flag (PRIMARIES, REPLICAS, RANDOM, ALL_NODES),ClusterNode, list<ClusterNode>, or dict<any, ClusterNode>. The passed type is )r  r   rx   dictr   	TypeErrortype)r   r  r  s      r   _parse_target_nodes RedisCluster._parse_target_nodesv  s    lD)) E   k22!NE  d++ ,,./E  & '+<&8%9; r   erroris_internalretry_attemptsc           
         #    [        UR                  UR                  UR                  UR                  UUb  UOSUS9I Sh  vN   g N7f)zY
Records error count metric directly.
Accepts either a Connection or ClusterNode object.
Nr   server_addressserver_portnetwork_peer_addressnetwork_peer_port
error_typer  r  )r0   r   r   )r   r  rF  r  r  s        r   _record_error_metric!RedisCluster._record_error_metric  sE      !%??"!+(oo-;-G>Q#
 	
 	
s   A A
AA
r   duration_secondsc           	         #    [        US5      (       a  UR                  nOUR                  R                  SS5      n[	        UUUR
                  UR                  Ub  [        U5      OSUS9I Sh  vN   g N7f)z`
Records operation duration metric directly.
Accepts either a Connection or ClusterNode object.
r   r   Nr   r  r  r  db_namespacer  )r>  r   r   r   r1   r   r   r  )r   r   r  rF  r  r   s         r   _record_command_metric#RedisCluster._record_command_metric  sk      :t$$B--11$:B'%-%??"$&NR
 	
 	
s   A.A80A61A8c           
      H	  ^ ^^#    TS   n/ nSnT R                   R                  5       nTR                  SS5      nU(       a+  T R                  U5      (       d  T R	                  U5      nSnSnT R
                  R                  TS   R                  5       5      I Sh  vN nU(       d  U(       d  T R                  R                  U5      n	U	(       dd  T R                  5       (       d  Sn
OT R                  " T6 I Sh  vN n
U
c  [        5       nOq[        [        R                  [        R                  S9nOJU	T R                   ;   a  [        T R                   U	   S9nO#[        5       nOU(       d  U(       a
  [        5       nSU-   nSn["        R$                  " 5       nSn['        U5       GH  nT R(                  (       aO  T R+                  US	9I Sh  vN   Sn[-        U5      S:X  a'  US   T R                  5       :X  a  T R/                  5          U(       d;  T R0                  " TUR2                  US
.6I Sh  vN nU(       d  [5        ST S35      e[-        U5      S:X  aw  T R6                  " US   /TQ70 TD6I Sh  vN nUT R8                  ;   a%  T R8                  U   " X4S   R:                  U040 TD6nT R<                  UR>                     " U5      s  $ U Vs/ s H  nUR:                  PM     nn[@        RB                  " UUU 4S jU 5       6 I Sh  vN nUT R8                  ;   a,  T R8                  U   " U[E        [G        UU5      5      40 TD6s  $ T R<                  UR>                     " [E        [G        UU5      5      5      s  $    g GN GN GN GN` GNs  snf  N! [H         a  nUS:  a  [K        U5      T RL                  RN                  ;   a  US-  nUS-  n[Q        USS5      n[S        US5      (       a_  T RU                  U["        R$                  " 5       U-
  URV                  US9I Sh  vN    T RY                  UURV                  US9I Sh  vN     SnAGM  [S        US5      (       a%  T RY                  UURV                  USS9I Sh  vN    UeSnAff = f7f)a!  
Execute a raw command on the appropriate cluster node or target_nodes.

It will retry the command as specified by the retries property of
the :attr:`retry` & then raise an exception.

:param args:
    | Raw command args
:param kwargs:

    - target_nodes: :attr:`NODE_FLAGS` or :class:`~.ClusterNode`
      or List[:class:`~.ClusterNode`] or Dict[Any, :class:`~.ClusterNode`]
    - Rest of the kwargs are passed to the Redis connection

:raises RedisClusterException: if target_nodes is not provided & the command
    can't be mapped to a slot
r   Fr  NTr  response_policyr  r!  )r  r  r  !No targets were found to execute  command onc              3   x   >#    U  H/  n[         R                  " TR                  " U/TQ70 TD65      v   M1     g 7fr   )r  create_task_execute_command).0r  r~  r}   r   s     r   	<genexpr>/RedisCluster.execute_command.<locals>.<genexpr>  sD       )5 $// $ 5 5d LT LV L  )5s   7:r  rF  r   r  rF  r  )r  rF  r  )r  rF  r  r  )-r   get_retriesr   r  r  r   resolver  r   r   r   r}  r#   r$   r   r%   r   time	monotonicranger   r  ro  replace_default_noder  r  r`   r  r   r  r   r  r  gatherr  zipr  r  r   ERRORS_ALLOW_RETRYgetattrr>  r  rF  r  )r   r~  r}   rz  r  target_nodes_specifiedr  passed_targetscommand_policiescommand_flagrp  execute_attemptsfailure_count
start_timer  _retr  r  r   es   ```                  r   execute_commandRedisCluster.execute_command  s    $ q'!&//1ND9$"4"4^"D"D33NCL%)"N!%!6!6!>!>tAw}}!OO(>--11':L,,..D!%!5!5t!<<D<'6'8$'6'4'B'B(6(D(D($
  4#>#>>'6'+'B'B<'P($ (7'8$!&<.0 ~- ^^%
 $'(Aoo<QoRRR(,%%*$Q4+@+@+BB --/F-)-)>)>'7'F'F"0* $L
 (3?v[Q  |$) $ 5 5l1o W WPV WWC$"7"77"33G<#1o&:&:C%@DJ  ::(88  3??,$DII,D?#*>> )5	$ F $"7"77#44W=#T#dF*;%< @F    ::(883tV,-/ /_ )G P =8 S$ X @  !A%$q'T^^5V5V*V #a'N!Q&M,3A7NPT,U)q,//"99)0-1^^-=
-J'(||"#	 :    #77"#'(||+8 8   
  q,//"77"#'(||+8(-	 8    G=s   BR"NAR";N<CR"N<R"'N!4N5AN!:N;AN!R"N!N/$N!N=N!R"1N!R"R"R"N!N!N!!
R+BR1P42#RQRR""0RRRRR"target_nodec                   #    S=pES nU R                   nUS   n[        R                  " 5       n	US:  a  US-  n U(       a+  U R                  US9nUR	                  S5      I S h  vN   SnOsU(       al  U R
                  " U6 I S h  vN n
U R                  R                  U
U R                  =(       a    US   [        ;   US   [        ;   a  U R                  OS 5      nSnUR                  " U0 UD6I S h  vN nU R                  U[        R                  " 5       U	-
  US9I S h  vN   U$ [Q        S5      nXl        U R                  U[        R                  " 5       U	-
  UUS9I S h  vN   Ue GN	 N N} NP! [         a=  nXl        U R                  U[        R                  " 5       U	-
  UUS9I S h  vN    e S nAf[         a=  nXl        U R                  U[        R                  " 5       U	-
  UUS9I S h  vN    e S nAf[        [         4 a  nUR#                  5         UR%                  5       I S h  vN    U R                  R'                  UR(                  5        UR(                  Ul        SU l        Xl        U R                  U[        R                  " 5       U	-
  UUS9I S h  vN    e S nAf[.        [0        4 au  nU R3                  5       I S h  vN    [4        R6                  " S	5      I S h  vN    Xl        U R                  U[        R                  " 5       U	-
  UUS9I S h  vN    e S nAf[8         a  nU =R:                  S-  sl        U R<                  (       a>  U R:                  U R<                  -  S:X  a!  U R3                  5       I S h  vN    SU l        O$U R                  R?                  U5      I S h  vN    SnU R                  U[        R                  " 5       U	-
  UUS9I S h  vN    U RA                  UUS
9I S h  vN     S nAGOS nAf[B         au  n[E        URF                  URH                  S9nSnU R                  U[        R                  " 5       U	-
  UUS9I S h  vN    U RA                  UUS
9I S h  vN     S nAGOS nAf[J         a  nXpR                   S-  :  a  [4        R6                  " S5      I S h  vN    U R                  U[        R                  " 5       U	-
  UUS9I S h  vN    U RA                  UUS
9I S h  vN     S nAOS nAf[L         a=  nXl        U R                  U[        R                  " 5       U	-
  UUS9I S h  vN    e S nAf[N         a=  nXl        U R                  U[        R                  " 5       U	-
  UUS9I S h  vN    e S nAff = fUS:  a  GM  GN GNY7f)NFr   r!  r_  ASKING)r   r  rF  r  T      ?)r  rF  r   r   rl  g?zTTL exhausted.))RedisClusterRequestTTLr  r  ra  r  r}  r   r|  r   rF   r   r  rV   rF  r^   rY   rd   'update_active_connections_for_reconnectdisconnect_free_connections move_node_to_end_of_cached_nodesr  r  r   rW   rc   r  r  sleepr_   r   r   	move_slotr  rU   rA   r   r   re   rb   r  rX   )r   r  r~  r}   askingmovedredirect_addrttlrz  r  rp  responser  s                r   r  RedisCluster._execute_commandJ  s     ))q'^^%
Ag1HCd"&----"HK%55h???"F "&!5!5t!<<D"&"4"4"G"G//LDG}4L7m3 44!#K "E!,!<!<d!Mf!MM11!(%)^^%5
%B* 2   
  \ )*")) !^^-
:"	 * 
 	
 	
 W @
 = N $ *11!(%)^^%5
%B*	 2    & 
  +11!(%)^^%5
%B*	 2    #\2 
 CCE!==??? ""CCKDTDTU*5*:*:' $( *11!(%)^^%5
%B*	 2    $&9:  kkm##mmD)))*11!(%)^^%5
%B*	 2      ))Q.)++11D4K4KKqP++-''01D-,,66q99911!(%)^^%5
%B*	 2    //* 0      -166 G11!(%)^^%5
%B*	 2    //* 0    ! 44q88!-----11!(%)^^%5
%B*	 2    //* 0    ! *11!(%)^^%5
%B*	 2     *11!(%)^^%5
%B*	 2    } AggT	
s  7W*E6 $E-% E6 E0A.E6 4E25.E6 #E4$E6 )=W&V?'W-E6 0E6 2E6 4E6 6
V1 1F81F42F88V11G=6G97G==V1#J.3H64A3J.'J*(J..V1L1KL13K646L1*L-+L11V1>APN+P O2P3O64PPPWV1%AR0Q31R	R
RWV1"+T"S0T">T?T"TT"W"V1/1U' U#!U''V141V,%V(&V,,V11	W?Wtransaction
shard_hintClusterPipelinec                 <    U(       a  [        S5      e[        X5      $ )z
Create & return a new :class:`~.ClusterPipeline` object.

Cluster implementation of pipeline does not support transaction or shard_hint.

:raises RedisClusterException: if transaction or shard_hint are truthy values
z(shard_hint is deprecated in cluster mode)r`   r  )r   r  r  s      r   pipelineRedisCluster.pipeline  s     '(RSSt11r   ClusterPubSubc                      [        U 4XUS.UD6$ )a_  
Create and return a ClusterPubSub instance.

Allows passing a ClusterNode, or host&port, to get a pubsub instance
connected to the specified node

:param node: ClusterNode to connect to
:param host: Host of the node to connect to
:param port: Port of the node to connect to
:param kwargs: Additional keyword arguments
:return: ClusterPubSub instance
)r  r   r   )r  )r   r  r   r   r}   s        r   pubsubRedisCluster.pubsub  s    & TMdMfMMr   
key_prefixignore_subscribe_messagesr   c                     SSK Jn  U" U UUS9$ )a\  
Return an
:class:`~redis.asyncio.keyspace_notifications.AsyncClusterKeyspaceNotifications`
object for subscribing to keyspace and keyevent notifications across
all primary nodes in the cluster.

Note: Keyspace notifications must be enabled on all Redis cluster nodes
via the ``notify-keyspace-events`` configuration option.

Args:
    key_prefix: Optional prefix to filter and strip from keys in
                notifications.
    ignore_subscribe_messages: If True, subscribe/unsubscribe
                              confirmations are not returned by
                              get_message/listen.
r   r   )r  r  )$redis.asyncio.keyspace_notificationsr   )r   r  r  r   s       r   keyspace_notifications#RedisCluster.keyspace_notifications)  s    *	
 1!&?
 	
r   r  timeoutr  blockingblocking_timeout
lock_classthread_localraise_on_release_errorc	                 .    Uc  [         nU" U UUUUUUUS9$ )a  
Return a new Lock object using key ``name`` that mimics
the behavior of threading.Lock.

If specified, ``timeout`` indicates a maximum life for the lock.
By default, it will remain locked until release() is called.

``sleep`` indicates the amount of time to sleep per loop iteration
when the lock is in blocking mode and another client is currently
holding the lock.

``blocking`` indicates whether calling ``acquire`` should block until
the lock has been acquired or to fail immediately, causing ``acquire``
to return False and the lock not being acquired. Defaults to True.
Note this value can be overridden by passing a ``blocking``
argument to ``acquire``.

``blocking_timeout`` indicates the maximum amount of time in seconds to
spend trying to acquire the lock. A value of ``None`` indicates
continue trying forever. ``blocking_timeout`` can be specified as a
float or integer, both representing the number of seconds to wait.

``lock_class`` forces the specified lock implementation. Note that as
of redis-py 3.0, the only lock class we implement is ``Lock`` (which is
a Lua-based lock). So, it's unlikely you'll need this parameter, unless
you have created your own custom lock class.

``thread_local`` indicates whether the lock token is placed in
thread-local storage. By default, the token is placed in thread local
storage so that a thread only sees its token, not a token set by
another thread. Consider the following timeline:

    time: 0, thread-1 acquires `my-lock`, with a timeout of 5 seconds.
             thread-1 sets the token to "abc"
    time: 1, thread-2 blocks trying to acquire `my-lock` using the
             Lock instance.
    time: 5, thread-1 has not yet completed. redis expires the lock
             key.
    time: 5, thread-2 acquired `my-lock` now that it's available.
             thread-2 sets the token to "xyz"
    time: 6, thread-1 finishes its work and calls release(). if the
             token is *not* stored in thread local storage, then
             thread-1 would see the token value as "xyz" and would be
             able to successfully release the thread-2's lock.

``raise_on_release_error`` indicates whether to raise an exception when
the lock is no longer owned when exiting the context manager. By default,
this is True, meaning an exception will be raised. If False, the warning
will be logged and the exception will be suppressed.

In some use cases it's necessary to disable thread local storage. For
example, if you have code where one thread acquires a lock and passes
that lock instance to a worker thread to release later. If thread
local storage isn't disabled in this case, the worker thread won't see
the token set by the thread that acquired the lock. Our assumption
is that these cases aren't common and as such default to using
thread local storage.)r  r  r  r  r  r  r.   )	r   r  r  r  r  r  r  r  r  s	            r   lockRedisCluster.lockH  s5    H J-%#9	
 		
r   funcc                   #    UR                  SS5      nUR                  SS5      nUR                  SS5      nU R                  SU5       ISh  vN n  U(       a  UR                  " U6 I Sh  vN   U" U5      I Sh  vN nUR                  5       I Sh  vN n	U(       a  UOU	 sSSS5      ISh  vN   $  Ni NK N= N' N! [         a#    Ub  US:  a  [
        R                  " U5         M  f = f! , ISh  vN  (       d  f       g= f7f)z
Convenience method for executing the callable `func` as a transaction
while watching all keys specified in `watches`. The 'func' callable
should expect a single argument which is a Pipeline object.
r  Nvalue_from_callableFwatch_delayTr   )r   r  watchexecuterf   r  r  )
r   r  watchesr}   r  r  r  pipe
func_value
exec_values
             r   r  RedisCluster.transaction  s      ZZd3
$jj)>Fjj5==z22d	"jj'222'+Dz!1J'+||~!5J)<:*L 322 3!1!5 3 " ".;?

;/ 322s   ADB9DC3C/B;0C?B= CB?C'D3C4D;C=C?CD)C0,C3/C00C33D
9C<:D
D)r   r   r   r   r   r   r   r  r  r   r   r   r   r   r   r   r   r   r   r   r   r   r   NNr~   N)r~   rz   )r~   rx   r  rx   r~   NNNNF)TNr   )NT)Ng?TNNTT)b__name__
__module____qualname____firstlineno____doc__classmethodr  r
   r   r   r   __annotations__	__slots__ro   r   rm   r    r   r   rK   r  r   boolr?   r2   r   r  rN   objectrO   floatr   bytesr   r   rT   rJ   r	  r   r  r  rp   r  r  r  r  r.  r   r2  r?  warningswarnr  get_running_looprD  r*   r   r   r   r   rZ  r   rc  ra  rs  r   rv  r   r   r   ri   rm  r"   r  r   r  r  r(   r  r$   r  r}  r  r  r   r  r  r  rj   r  r  r  r	  r/   r  r   r  __static_attributes__r   r   r   rz   rz      s7	   m^ 3 # .  B '+gdm*I" *+G
 *
 =  -0H  48&*#(@D&*"#,?""7;9=##"&(0+32:'!&'('=/M 8!%NV#'"&,6?C?C#'#'"&/3"&#!%MQ37/H/Jic*Djc* Cic*
 M*T1c*  $c* !c* "7!=c*  $c*  c* '*c* c* t|c* T)_-4c*  #I!c*" Dj#c*$ 0$6%c*& *'c*( *)c** 4Z+c*, ,%-c*. 6\D(/c*0  &(4/1c*4 5c*6 7c*8 9c*<  %=c*> ?c*@ !&Ac*B Cc*D Ec*F #*#sU{*:";f"Dt"KGc*J Kc*L DjMc*N 4ZOc*P *Qc*R #'}"5"<Sc*T #'}"5"<Uc*V DjWc*X !Yc*Z 4Z[c*\ -]c*^ 4Z_c*` *ac*b cc*d  sCx 15c? BCdJec*f *D0gc*h -ic*j 
kc*"c*N JN/3'/U38_0E'F  (} 
	4	E 1GgV W' '' '0-9S$%>? - 2L ]],,  
	
=: 
=$ 
==4. =DtM2 DDd=1 DL/	/ #"#'	BsmB smB C=	B
 
-	 B ).$$!%$	-	 $8.3
 
	%8D,?#@ 	%2: 2# 2W &tC#,>'? &u  4S 4<M 4RV 4 $(  &	
 C= 
m	<1S 1 1 1fQ# Q$ Q ]8K 0 !(,

 *m34
 	

 !
4 &*

  
 *m34	

 	"
2P: P P Pdy(y16tZ7G1HyTWy	yx NR2#C=2=Ec]2	2" )-""	N}%N smN sm	N
 N 
N. /3*.
#ud*+
 $(
 
-	
D $(,0+/!'+O
O
 %O
 	O

 O
 #5/O
 T$Z(O
 O
 !%O
 
O
bd$5s:;r   rz   c                      \ rS rSrSrSr S-S\S.S\S\\\	4   S	\
\   S
\	S\\   S\SS4S jjjrS\4S jrS\S\4S jrS\	4S jrSr\R(                  \R,                  4S\S\SS4S jjrS.S jrS\4S jrS\SS4S jrS\SS4S jrS\SS4S jrS\4S jrS.S jrS.S jr S\S \S!\S\4S" jr!S#\S!\S\4S$ jr"S%\#S&   S\4S' jr$S(\%4S) jr&S*\'4S+ jr(S,r)g)/rx   i  z
Create a new ClusterNode.

Each ClusterNode manages multiple :class:`~redis.asyncio.connection.Connection`
objects for the (host, port).
)_background_tasks_connections_freer   r   r   r   r   r   r  r   r   server_typeNr   )r   r   r   r   r8  r   r   r   r~   c                   US:X  a  [         R                  " U5      nXS'   X&S'   Xl        X l        [	        X5      U l        X0l        X@l        XPl        X`l	        UR                  S0 5      U l        / U l        [        R                  " U R                  S9U l        [!        5       U l        U R                  R%                  SS 5      U l        U R&                  c  [)        5       U l        g g )N	localhostr   r   r   )maxlenr   )socketgethostbynamer   r   rA   r  r8  r   r   r   r   r   r6  collectionsdequer7  setr5  r   r   rT   )r   r   r   r8  r   r   r   s          r   r	  ClusterNode.__init__  s     ;''-D$(&!$(&!		!$-	&. 0!2"3"7"78Lb"Q.0(3(9(9AUAU(V
47E!%!7!7!;!;<NPT!U!!)%4%6D" *r   c           	      p    SU R                    SU R                   SU R                   SU R                   S3	$ )Nz[host=z, port=z, name=z, server_type=])r   r   r  r8  r   s    r   __repr__ClusterNode.__repr__  s?    TYYKwtyyk 2II;nT-=-=,>aA	
r   objc                 b    [        U[        5      =(       a    UR                  U R                  :H  $ r   )r  rx   r  )r   rF  s     r   __eq__ClusterNode.__eq__  s!    #{+EDII0EEr   c                 ,    [        U R                  5      $ r   )hashr  r   s    r   __hash__ClusterNode.__hash__  s    DIIr   zUnclosed ClusterNode objectr5  r6  c                     U R                    HW  nUR                  (       d  M  U" U R                   SU < 3[        U S9   X R                  S.nU" 5       R	                  U5          g    g ! [
         a     Nf = f)Nr8  r9  r;  )r6  is_connectedr?  r@  rA  rB  )r   r5  r6  rF  rC  s        r   rD  ClusterNode.__del__  sz    
 ++J&&&**+1TH5tT)-:K:KLGF11':  , $ s    $A))
A65A6c                    #    [         R                  " S U R                   5       SS06I S h  vN n[        S U 5       S 5      nU(       a  Ueg  N!7f)Nc              3   j   #    U  H)  n[         R                  " UR                  5       5      v   M+     g 7fr   r  r  
disconnectr  rF  s     r   r  )ClusterNode.disconnect.<locals>.<genexpr>
  s.      "3J ##J$9$9$;<<"3   13return_exceptionsTc              3   T   #    U  H  n[        U[        5      (       d  M  Uv   M      g 7fr   )r  r  )r  r   s     r   r  rV    s     E3C*S)*DCC3s   (	()r  r  r6  next)r   r  excs      r   rT  ClusterNode.disconnect  s[     NN"&"3"3

 #
 
 E3EtLI 
s   +AA"Ac                 |    U R                   R                  5       $ ! [         a    [        U R                  5      U R
                  :  ag  [        [        5       S[        4S9nU R                  R                  5       nXS'   U R                  " S0 UD6nU R                  R                  U5        Us $ [        5       ef = f)Nr   )r   r   supported_errorsr   r   )r7  popleft
IndexErrorro  r6  r   r2   r5   rY   r   r   r   r   r^   )r   r   r   rF  s       r   acquire_connectionClusterNode.acquire_connection  s    	(::%%'' 	(4$$%(<(<< %K&5%7
 %)$:$:$?$?$A!-2'*!22G5FG
!!((4!!%''-	(s    BB;0B;rF  c                 l   #    UR                  5       (       a  UR                  5       I Sh  vN   gg N7f)z
Disconnect a connection if it's marked for reconnect.
This implements lazy disconnection to avoid race conditions.
The connection will auto-reconnect on next use.
N)should_reconnectrT  rL  s     r   disconnect_if_needed ClusterNode.disconnect_if_needed/  s0      &&((''))) ))s   )424c                 0   UR                  5       (       af  [        R                  " U R                  U5      5      nU R                  R                  U5        UR                  U R                  R                  5        gU R                  R                  U5        g)z
Release connection back to free queue.
If the connection is marked for reconnect, disconnect it before
returning it to the free queue.
N)
rd  r  r  _disconnect_and_releaser5  addadd_done_callbackdiscardr7  r   )r   rF  tasks      r   releaseClusterNode.release8  sq     &&((&&t'C'CJ'OPD""&&t,""4#9#9#A#AB

*%r   c                 *  #     UR                  5       I S h  vN   U R                  R                  U5        g  N ! [         aL  n[        R                  SUSS9   U R                  R                  U5        O! [         a     Of = f S nAg S nAff = f7f)Nz4disconnecting released cluster connection failed: %rTexc_info)	rT  r  loggerdebugr6  remove
ValueErrorr7  r   )r   rF  r[  s      r   rh  #ClusterNode._disconnect_and_releaseE  s     	''))) 	

*% * 
	LLF  
!!((4 
	s[   B: 8: B: 
BBA65B6
B BBBBBBc                     U R                   nUR                  S[        5      nU" UR                  SS5      UR                  SS5      UR                  SS5      S9$ )	zFReturn an :class:`Encoder` derived from this node's connection kwargs.encoder_classr   r   r   r   r   F)r   r   r   )r   r   r"   )r   r}   rx  s      r   r  ClusterNode.get_encoderV  sV    ''

?G<ZZ
G4"JJ'8(C#ZZ(:EB
 	
r   c                     [        U R                  5      nU R                   H  nX!;  d  M
  UR                  5         M     g)z
Mark all in-use (active) connections for reconnect.
In-use connections are those in _connections but not currently in _free.
They will be disconnected after their current operation completes.
N)r@  r7  r6  mark_for_reconnect)r   free_setrF  s      r   r  3ClusterNode.update_active_connections_for_reconnect`  s3     tzz?++J)--/ ,r   c                    #    U R                   (       a9  [        R                  " S [        U R                   5       5       SS06I Sh  vN   gg N7f)z
Disconnect all free/idle connections in the pool.
This is useful after topology changes (e.g., failover) to clear
stale connection state like READONLY mode.
The connections remain in the pool and will reconnect on next use.
c              3   @   #    U  H  oR                  5       v   M     g 7fr   )rT  rU  s     r   r  :ClusterNode.disconnect_free_connections.<locals>.<genexpr>u  s     N<Mj''))<Ms   rX  TN)r7  r  r  tupler   s    r   r  'ClusterNode.disconnect_free_connectionsk  sG      ::..NE$**<MN"&   s   AAAArz  r}   c                   #     [         U;   a-  UR                  SS9I S h  vN nUR                  [         5        OUR                  5       I S h  vN n [        U;   a  UR                  [        5        UR                  SS 5        X R
                  ;   a  U R
                  U   " U40 UD6$ U$  N N_! [         a    [        U;   a  U[           s $ e f = f7f)NT)disable_decodingr  )r7   rK  r   rb   r6   r   )r   rF  rz  r}   r  s        r   parse_responseClusterNode.parse_responsey  s     		v%!+!9!94!9!PP

<(!+!9!9!;; V#JJ~& 	

64  ---**73HGGG' Q < 	'n--	sU   CB0 B,B0 CB0 B.B0 AC,B0 .B0 0CCCCr~  c                   #    U R                  5       n U R                  U5      I S h  vN   UR                  UR                  " U6 5      I S h  vN   U R                  " X1S   40 UD6I S h  vN  U R                  U5      I S h  vN   U R                  U5        $  Ns NO N3 N! U R                  U5        f = f!  U R                  U5      I S h  vN    U R                  U5        f ! U R                  U5        f = f= f7fr   )ra  re  send_packed_commandpack_commandr  rm  )r   r~  r}   rF  s       r   r  ClusterNode.execute_command  s     ,,.
	)++J777 001H1H$1OPPP ,,ZaKFKK)//
;;; Z( 8 Q L < Z(	)//
;;; Z(Z(s   DB: B%B: BB: ,B -B: 1B$B"B$
DB: B:  B: "B$$B77D:C><C(CC(C>(C;;C>>DcommandsPipelineCommandc                   #    U R                  5       n U R                  U5      I S h  vN   UR                  UR                  S U 5       5      5      I S h  vN   SnU H;  n U R                  " X$R
                  S   40 UR                  D6I S h  vN Ul        M=     U U R                  U5      I S h  vN   U R                  U5        $  N Nv N>! [         a  nXTl        Sn S nAM  S nAff = f N<! U R                  U5        f = f!  U R                  U5      I S h  vN    U R                  U5        f ! U R                  U5        f = f= f7f)Nc              3   8   #    U  H  oR                   v   M     g 7fr   )r~  )r  r   s     r   r  /ClusterNode.execute_pipeline.<locals>.<genexpr>  s     (FXcXs   Fr   T)
ra  re  r  pack_commandsr  r~  r}   resultr  rm  )r   r  rF  r  r   r  s         r   execute_pipelineClusterNode.execute_pipeline  s=    ,,.
	)++J777 00(((FX(FF  
 C'+':':"HHQK(36::( "CJ   )//
;;; Z(1 8" ! !"JC < Z(	)//
;;; Z(Z(s   ED
 C.D
 CD
 $-CC	CD
 !C45C26C4:ED
 D
 C
C/C*$D
 *C//D
 2C44DE
ED8 D#!D8&E8EEEtokenc                   ^ ^^#    [         R                  " 5       nT R                  (       a  T R                  R                  5       mTR                  R                  UU4S jU 4S j5      I S h  vN   TR                  R                  U4S jU 4S j5      I S h  vN   UR                  T5        T R                  (       a  M  U(       a5  UR                  5       mT R                  R                  T5        U(       a  M4  g g  N Ng7f)Nc                  d   > T R                  STR                  S5      TR                  5       5      $ )NAUTHoid)rJ  try_get	get_value)connr  s   r   r   .ClusterNode.re_auth_callback.<locals>.<lambda>  s'    ))EMM%0%//2Cr   c                 &   > TR                  U 5      $ r   _mockr  r   s    r   r   r    s    djj/r   c                  $   > T R                  5       $ r   )rK  )r  s   r   r   r    s    **,r   c                 &   > TR                  U 5      $ r   r  r  s    r   r   r    s    DJJu<Mr   )r>  r?  r7  r_  r   call_with_retryr   )r   r  	tmp_queuer  s   `` @r   re_auth_callbackClusterNode.re_auth_callback  s     %%'	jj::%%'D**,, 0	   **,,,.M   T" jjj $$&DJJd# is0   A)D.D /+DD&D9D>DDr  c                    #    g7f)z_
Dummy functions, needs to be passed as error callback to retry object.
:param error:
:return:
Nr   )r   r  s     r   r  ClusterNode._mock  s
      	   )r5  r6  r   r7  r   r   r   r   r  r   r   r8  r   r   )*r$  r%  r&  r'  r(  r+  r*   r  r   r  r   r   r
   r	  rD  r,  rH  rL  r?  r0  r1  r  r2  rD  rT  ra  re  rm  rh  r"   r  r  r  r  r  r   r  r3   r  ra   r  r3  r   r   r   rx   rx     s   I( &*	7  #-777 CHo7 c]	7 7 z*7 !7 
7@
# 
F# F$ F#  1L ]],,  
	 
(J (6*Z *D *&* & &&
 &t &"
W 
	0$/2>A	4)3 )# )# )&)t4E/F )4 )>$N $& r   c                      \ rS rSrSr   S$S\S   S\S\\\	4   S\S	\
\\\\4   /\\\4   4      S
\
\   SS4S jjr   S%S\
\   S\
\   S\
\   S\
S   4S jjr S&S\\S4   S\\S4   S\SS4S jjrS\SS4S jrS\\-  4S jr  S'S\S\SS4S jjrS\S\S   4S jr  S(S\
\\\\4         S\
\   SS4S jjrS)S \SS4S! jjrS\S\S\\\4   4S" jrS#rg)*r   i  )_dynamic_startup_nodesr   r5  r   r  rO  _epochread_load_balancer_initialize_lockr   rn  r   r   Nr   rx   r   r   r   r   r   r~   c                 \   U Vs0 s H  owR                   U_M     snU l        X l        X0l        XPl        S U l        0 U l        0 U l        SU l        [        5       U l
        [        R                  " 5       U l        [        5       U l        X@l        Uc  [#        5       U l        g X`l        g s  snf r   )r  r   r   r   r   r  rO  rn  r  r>   r  r  r/   r  r@  r5  r  rT   r   )r   r   r   r   r   r   r   r  s           r   r	  NodesManager.__init__  s     ;HH-$iio-H%:"!2*+/57;="...5lln47E,A##%4%6D"%5"# Is   B)r   r   r`  c                     U(       aE  U(       a>  US:X  a  [         R                  " U5      nU R                  R                  [	        XS95      $ U(       a  U R                  R                  U5      $ [        S5      e)Nr:  r  zEget_node requires one of the following: 1. node name 2. host and port)r<  r=  rO  r   rA   r[   rf  s       r   ra  NodesManager.get_node  sg     D{"++D1##''4(KLL##''	22W r   oldnew
remove_oldc                 B   U(       a  [        UR                  5       5       H  nXB;  d  M
  UR                  U5      nUR                  5         [        R
                  " UR                  5       5      nU R                  R                  U5        UR                  U R                  R                  5        M     UR                  5        HX  u  pGXA;   aJ  X   nUR                  Ul        UR                  5         UR                   H  n	U	R                  5         M     MT  XqU'   MZ     g r   )r   r  r   r  r  r  r  r5  ri  rj  rk  itemsr8  r7  r{  )
r   r  r  r  r  removed_noderl  r  existing_noder  s
             r   	set_nodesNodesManager.set_nodes   s     SXXZ(? $'774=L HHJ"..$@@BD **..t4**4+A+A+I+IJ )  ))+JD{ !$	,0,<,<)EEG)//D++- 0I! &r   c                 L   XR                   ;   aB  [        U R                   5      S:  a)  U R                   R                  U5      nX R                   U'   XR                  ;   aD  [        U R                  5      S:  a*  U R                  R                  U5      nX R                  U'   ggg)z
Move a failing node to the end of startup_nodes and nodes_cache so it's
tried last during reinitialization and when selecting the default node.
If the node is not in the respective list, nothing is done.
r!  N)r   ro  r   rO  )r   r`  r  s      r   r  -NodesManager.move_node_to_end_of_cached_nodesI  s     ***s43E3E/F/J%%)))4D,0y) (((S1A1A-BQ-F##''	2D*.Y' .G(r   r  c                 d  #    SnU R                  UR                  UR                  S9nU(       a   UR                  [        :w  a  [        Ul        OX[        UR                  UR                  [        40 U R                  D6nU R                  U R                  UR                  U05        U R                  UR                     nX4;  a  U/U R                  UR                  '   SnOUX4S   LaN  US   n[        Ul        UR                  U5        UR                  U5        X4S'   U R                  U:X  a  X0l        SnU(       a-   U R                   R#                  [%        5       5      I S h  vN   g g  N! [&         a4  n[(        R+                  S[-        U5      R.                  U5         S nAg S nAff = f7f)NFr  Tr   2listener raised during slots-cache refresh: %s: %s)ra  r   r   r8  r:   rx   r   r  rO  r  rn  slot_idr;   r   rt  r  r   dispatch_asyncrR   r  rr  	exceptionr  r$  )r   r  node_changedredirected_node
slot_nodesold_primaryr[  s          r   r  NodesManager.move_slotZ  s    --QVV!&&-A**g5.5+ *+/+A+AO NN4++o.B.BO-TU%%aii0
, ,;*;DQYY'LqM1 %Q-K '.K#k* o.+qM  K/$3!L ,,;;46     	
   HI&& 	sB   D=F0 &E/ &E-'E/ +F0-E/ /
F-9*F(#F0(F--F0rp  r   c                    USL a  Uc  [         R                  n [        U R                  U   5      S:  ah  U(       aa  U R                  U   S   R                  nU R
                  R                  U[        U R                  U   5      U5      nU R                  U   U   $ U R                  U   S   $ ! [        [        4 a    [        SU SU R                   S35      ef = f)NTr!  r   rk  z5" not covered by the cluster. "require_full_coverage=")r?   ROUND_ROBINro  rn  r  r  get_server_indexr`  r  rc   r   )r   rp  r   r   primary_namerr  s         r   r|  NodesManager.get_node_from_slot  s     %*A*I&;&G&G#	4##D)*Q.3J#//5a8==22CC #d&6&6t&<"=?V ''-h77##D)!,,I& 	% **.*D*D)EQH 	s   BB0 B0 0-Cr8  c                     U R                   R                  5        Vs/ s H  nUR                  U:X  d  M  UPM     sn$ s  snf r   )rO  r   r8  )r   r8  r  s      r   rR  %NodesManager.get_nodes_by_server_type  sF     ((//1
1;. 1
 	
 
s   >>r  r  c                   #    U R                   R                  5         0 n0 n/ nSnSnS nU R                  n	Uc  / nU R                   IS h  vN   U R                  U	:w  a   S S S 5      IS h  vN   g [	        U R
                  R                  5       5      n
/ nUbF  [        U
5       H7  u  pUR                  U:X  d  M  UR                  U
R                  U5      5          O   [        U
5      S:  a  [        R                  " U
5        U VVs/ s H  u  p[        X40 U R                  D6PM     nnnUbO  [        U5       H@  u  pUR                  U:X  d  M  U(       d  UR                  U5        UR                  U5          O   [!        U
UU5       GH  n  U R"                  R%                  ['        U R(                  U R                  R+                  SS 5      5      5        UR-                  S5      I S h  vN nSn[        U5      S:X  a>  US   S   S   (       d.  [        U R
                  5      S:X  a  UR4                  US   S   S'   U GH  n[7        S[        U5      5       H%  nUU    Vs/ s H  n[9        U5      PM     snUU'   M'     US   nUS   nUS	:X  a  UR4                  n[;        US   5      nU R=                  X5      u  p/ nUR+                  [?        X5      5      nU(       d  [        X[@        40 U R                  D6nUUUR                  '   UR                  U5        US
S  nU H|  nUS   nUS   nU R=                  X5      u  pUR+                  [?        X5      5      nU(       d  [        X[B        40 U R                  D6nUUUR                  '   UR                  U5        M~     [7        [;        US   5      [;        US   5      S-   5       H  nUU;  a  UUU'   M  UU   S   nUR                  UR                  :w  d  M4  UR                  UR                   SUR                   SU 35        [        U5      S:  d  Mr  [1        SSRE                  U5       35      e   GM     Sn[7        [F        5       H  nUU;  d  M  Sn  O   U(       d  GM    O   U(       d  [1        S[I        U5       35      UeU(       d0  U RJ                  (       a  [1        S[        U5       S[F         S35      eU RM                  U R(                  USS9  0 n0 nURO                  5        HW  u  nn [Q        U 5      n!UR+                  U!5      n"U"c-  U  Vs/ s H  oR(                  UR                     PM     n"nU"UU!'   U"UU'   MY     UU l)        U RT                  (       a%  U RM                  U R
                  U R(                  SS9  U RW                  [@        5      S   U l,        U =R                  S-  sl        S S S 5      IS h  vN    U R"                  R[                  []        5       5      I S h  vN   g  GN GNs  snnf  GN?! [.         a    [1        S5      ef = f! [2         a  nUn S nAGM  S nAff = fs  snf s  snf  N~! , IS h  vN  (       d  f       N= f Nn! [2         a4  n[^        Ra                  S[c        U5      Rd                  U5         S nAg S nAff = f7f)NFr!  r   r   z(Cluster mode is not enabled on this nodeTr   rl      z vs z
 on slot: r   z6startup_nodes could not agree on a valid slots cache: z, zORedis Cluster cannot be connected. Please provide at least one reachable node: z9All slots are not covered after query all startup_nodes. z of z covered...)r  r  )3r  resetr  r  r   r   r   	enumerater  r   r   ro  rX  shufflerx   r   r   r   dispatchrQ   rO  r   r  rb   r`   r  r   r  rr   r  remap_host_portrA   r:   r;   joinrL   r  r   r  r  idrn  r  rR  r  r  rR   rr  r  r  r$  )#r   r  r  tmp_nodes_cache	tmp_slotsdisagreementsstartup_nodes_reachablefully_coveredr  epochr   deferred_failed_nodesindexr  r   r   additional_startup_nodesstartup_nodecluster_slotsr  rp  ivalprimary_nodenodes_for_slotr  replica_nodesreplica_nodetarget_replica_nodetmp_slotnode_lists_by_idnew_slots_cacher  node_list_idr  s#                                      r   r  NodesManager.initialize  s    
 	%%'4646	"'	(0,.)((({{e#  )(( !!3!3!:!:!<=M$&!$0#,]#;KEyy$99-44]5F5Fu5MN $< =!A% }- #@("?JD DA$*@*@A"? % ( %0#,-E#FKEyy$994188>044U; $G !&(%!
..77? $ 0 0 $ 6 6 : :;PRV W /;.J.J+/ ) /3+ &!+)!,Q/2D../14-9->->M!$Q'*)D"1c$i0@DQ"H<#4"HQ 1#'7L'?Drz+00|A/D!%!5!5d!AJD%'N"1"5"5mD6O"PK&&1 '373I3I' 9DOK$4$45"))+6$(HM(5+A+A%)%9%9$%E
.=.A.A)$5/+  32= $G37;7M7M3/ EX(;(@(@A&--.AB )6  #3tAw<T!W1ABI-+9IaL (1|AH'}}0@0@@ - 4 4'/}}oT+:J:J9K:VWUX$Y!" $'}#5#9*?+88<		-8P7Q)S+& %& CM *r !%78A	)(- 9 !=S!V ++++.y>*:< !! !T%?%? ,O9~&d+C*D E!"  NN4++_NN @B>@O(0e!%y-11,?
%JO!P%$"2"2499"=%J!P5?$\2(2%  1  /D**t1143C3CPTU !% = =g Fq IDKK1Kg )(p		((7702  s )((2) ) 3F 
 !  !"I	, #Iz "QQ )(((r  	DQ   	sE  A[
X#[Y2 [+X&,[1AY2;AY2	"X)+$Y2>Y2AX20X/1X25Y7A2Y2)Y&>FY2:Y2<Y2Y2B2Y2"Y+'A?Y2&[1Y02[7&Z ZZ "[&[)Y2/X22YY
Y#YY2Y##Y20[2Z	8Y;9Z	[Z 
[*[[[[attrc                    #    S U l         [        R                  " S [        X5      R	                  5        5       6 I S h  vN   g  N7f)Nc              3   j   #    U  H)  n[         R                  " UR                  5       5      v   M+     g 7fr   rS  r  r  s     r   r  &NodesManager.aclose.<locals>.<genexpr>  s,      8D ##DOO$5668rW  )r  r  r  r  r   )r   r  s     r   r  NodesManager.aclose  s>      nn#D/668
 	
 	
s   <AAAc                 N    U R                   (       a  U R                  X45      $ X4$ )z
Remap the host and port returned from the cluster to a different
internal value.  Useful if the client is not connecting directly
to the cluster.
)r   )r   r   r   s      r   r  NodesManager.remap_host_port  s(     %%tl33zr   )r5  r  r  r   r  r   r   r  rO  r  r   rn  r   )TNNr"  r#  )FNr  )rO  )r$  r%  r&  r'  r+  r   r,  r   r  r
   r   r   r   r  rT   r	  ra  r  r  rU   r_   r  r|  rR  r  r  r  r3  r   r   r   r   r     s   I* '+PT6:6M*6  $6  S>	6
  $6  %S/):E#s(O)K LM6 #?36 
6< #"#'	sm sm C=	
 
-	 , !	'#}$%' #}$%' 	'
 
'R/# /$ /"9J!6 9| $) $	 !
 
0
C 
D<O 
 JN/3Q'/U38_0E'FQ  (}Q 
	Qf
 
 
C s uS#X r   r   c                   j   \ rS rSr% SrSrSr\S   \S'    S&S\	S\
\   S	S4S
 jjr\S'S j5       rS\S\S	S4S jrS(S jrS(S jrS)S jrS	\\SS 4   4S jrS	\4S jrS	\4S jrS\\\4   S\S	S 4S jr S*S\S\S	\\   4S jjr S\S\S	S 4S jr!S r"S r#S r$S  r%S! r&S" r'S#\(\)\4   S	S 4S$ jr*S%r+g)+r  i  aG  
Create a new ClusterPipeline object.

Usage::

    result = await (
        rc.pipeline()
        .set("A", 1)
        .get("A")
        .hset("K", "F", "V")
        .hgetall("K")
        .mset_nonatomic({"A": 2, "B": 3})
        .get("A")
        .get("B")
        .delete("A", "B", "K")
        .execute()
    )
    # result = [True, "1", 1, {"F": "V"}, True, True, "2", "3", 1, 1, 1]

Note: For commands `DELETE`, `EXISTS`, `TOUCH`, `UNLINK`, `mset_nonatomic`, which
are split across multiple nodes, you'll get multiple results for them in the array.

Retryable errors:
    - :class:`~.ClusterDownError`
    - :class:`~.ConnectionError`
    - :class:`~.TimeoutError`

Redirection errors:
    - :class:`~.TryAgainError`
    - :class:`~.MovedError`
    - :class:`~.AskError`

:param client:
    | Existing :class:`~.RedisCluster` client
)cluster_client_transaction_execution_strategyTr   Nr<  r  r~   c                     Xl         X l        U R                  (       d  [        U 5      U l        g [        U 5      U l        g r   )r  r  PipelineStrategyTransactionStrategyr  )r   r<  r  s      r   r	  ClusterPipeline.__init__  s?     %' $$ T" 	  %T* 	 r   c                 .    U R                   R                  $ )z.Get the nodes manager from the cluster client.)r  r   r   s    r   r   ClusterPipeline.nodes_manager  s     ""000r   rz  r  c                 :    U R                   R                  X5        g)z5Set a custom response callback on the cluster client.N)r  r  r  s      r   r  %ClusterPipeline.set_response_callback  s    11'Dr   c                 V   #    U R                   R                  5       I S h  vN   U $  N7fr   )r  r  r   s    r   r  ClusterPipeline.initialize  s'     &&11333 	4s   )')c                 >   #    U R                  5       I S h  vN $  N7fr   )r  r   s    r   r  ClusterPipeline.__aenter__  s     __&&&&s   c                 @   #    U R                  5       I S h  vN   g  N7fr   r  )r   r*  r+  r,  s       r   r.  ClusterPipeline.__aexit__       jjlr  c                 >    U R                  5       R                  5       $ r   r1  r   s    r   r2  ClusterPipeline.__await__  r4  r   c                     g)z?Pipeline instances should  always evaluate to True on Python 3+Tr   r   s    r   __bool__ClusterPipeline.__bool__  s    r   c                 ,    [        U R                  5      $ r   )ro  r  r   s    r   __len__ClusterPipeline.__len__  s    4++,,r   r~  r}   c                 :    U R                   R                  " U0 UD6$ )a$  
Append a raw command to the pipeline.

:param args:
    | Raw command args
:param kwargs:

    - target_nodes: :attr:`NODE_FLAGS` or :class:`~.ClusterNode`
      or List[:class:`~.ClusterNode`] or Dict[Any, :class:`~.ClusterNode`]
    - Rest of the kwargs are passed to the Redis connection
)r  r  r   r~  r}   s      r   r  ClusterPipeline.execute_command  s      ''77HHHr   raise_on_errorallow_redirectionsc                    #     U R                   R                  X5      I Sh  vN U R                  5       I Sh  vN   $  N N! U R                  5       I Sh  vN    f = f7f)a  
Execute the pipeline.

It will retry the commands as specified by retries specified in :attr:`retry`
& then raise an exception.

:param raise_on_error:
    | Raise the first error if there are any errors
:param allow_redirections:
    | Whether to retry each failed command individually in case of redirection
      errors

:raises RedisClusterException: if target_nodes is not provided & the command
    can't be mapped to a slot
N)r  r  r  r   r  r  s      r   r  ClusterPipeline.execute	  sP     $	1199  **,	 $**,sE   A!A ?A A!AA!A A!AAAA!r  c                     U R                   R                  U5      R                  5        H  nU R                  " U/UQ76   M     U $ r   )r  _partition_keys_by_slotr   r  )r   rz  r  	slot_keyss       r   _split_command_across_slots+ClusterPipeline._split_command_across_slots	  sC     ,,DDTJQQSI  595 T r   c                 T   #    U R                   R                  5       I Sh  vN   g N7fz
Reset back to empty pipeline.
N)r  r  r   s    r   r  ClusterPipeline.reset"	  s      &&,,...   (&(c                 8    U R                   R                  5         g)zz
Start a transactional block of the pipeline after WATCH commands
are issued. End the transactional block with `execute`.
N)r  multir   s    r   r'  ClusterPipeline.multi(	  s    
 	  &&(r   c                 T   #    U R                   R                  5       I Sh  vN   g N7f)r  N)r  rk  r   s    r   rk  ClusterPipeline.discard/	       &&..000r%  c                 R   #    U R                   R                  " U6 I Sh  vN   g N7f)z$Watches the values at keys ``names``N)r  r  r   namess     r   r  ClusterPipeline.watch3	  s     &&,,e444   '%'c                 T   #    U R                   R                  5       I Sh  vN   g N7f)z'Unwatches all previously specified keysN)r  unwatchr   s    r   r2  ClusterPipeline.unwatch7	  r+  r%  c                 R   #    U R                   R                  " U6 I S h  vN   g  N7fr   )r  unlinkr-  s     r   r5  ClusterPipeline.unlink;	  s     &&--u555r0  mappingc                 8    U R                   R                  U5      $ r   )r  mset_nonatomicr   r7  s     r   r9  ClusterPipeline.mset_nonatomic>	  s     ''66w??r   )r  r  r  r   )r~   r   r~   r  )r*  Nr+  Nr,  Nr~   NTT),r$  r%  r&  r'  r(  r+  r   r   r*  rz   r   r,  r	  propertyr   r  r(   r  r  r  r.  r   r
   r2  r  r  r  r   rj   ri   r  r   r  r   r  r'  rk  r  r2  r5  r   rg   r9  r3  r   r   r   r  r    sl   "HI '+gdm* CG	
"	
19$	
		
 1 1ES E<M ERV E'-9S$0A%AB -$ - -I4+,I8;I	I" GK"?C	c2#'	/)1516@w
23@	@r   r  r8  r  r9  c                   <    \ rS rSrS\S\S\SS4S jrS\4S jrS	r	g)
r  iL	  positionr~  r}   r~   Nc                 D    X l         X0l        Xl        S U l        S U l        g r   )r~  r}   r@  r  r  )r   r@  r~  r}   s       r   r	  PipelineCommand.__init__M	  s     	 -1;?r   c                 V    SU R                    SU R                   SU R                   S3$ )N[z]  ())r@  r~  r}   r   s    r   rD  PipelineCommand.__repr__T	  s)    4==/DII;bQ??r   )r~  r  r}   r@  r  )
r$  r%  r&  r'  r  r
   r	  r  rD  r3  r   r   r   r  r  L	  s6    @ @S @C @D @@# @r   r  c            	          \ rS rSr\SS j5       r\S\\\4   S\	SS4S j5       r
\ SS\S	\S\\	   4S
 jj5       r\S\\\4   SS4S j5       r\S 5       r\S 5       r\S 5       r\S 5       r\S 5       r\S 5       r\S\4S j5       rSrg)ExecutionStrategyiX	  r~   r  c                    #    g7f)zF
Initialize the execution strategy.

See ClusterPipeline.initialize()
Nr   r   s    r   r  ExecutionStrategy.initializeY	  
      	r  r~  r}   c                     g)zN
Append a raw command to the pipeline.

See ClusterPipeline.execute_command()
Nr   r  s      r   r  !ExecutionStrategy.execute_commandb	       	r   r  r  c                    #    g7f)z
Execute the pipeline.

It will retry the commands as specified by retries specified in :attr:`retry`
& then raise an exception.

See ClusterPipeline.execute()
Nr   r  s      r   r  ExecutionStrategy.executem	  s
      	r  r7  c                     g)zu
Executes multiple MSET commands according to the provided slot/pairs mapping.

See ClusterPipeline.mset_nonatomic()
Nr   r:  s     r   r9   ExecutionStrategy.mset_nonatomic{	  rO  r   c                    #    g7f)zB
Resets current execution strategy.

See: ClusterPipeline.reset()
Nr   r   s    r   r  ExecutionStrategy.reset	  rL  r  c                     g)z=
Starts transactional context.

See: ClusterPipeline.multi()
Nr   r   s    r   r'  ExecutionStrategy.multi	  s     	r   c                    #    g7f)z1
Watch given keys.

See: ClusterPipeline.watch()
Nr   r-  s     r   r  ExecutionStrategy.watch	  rL  r  c                    #    g7f)zI
Unwatches all previously specified keys

See: ClusterPipeline.unwatch()
Nr   r   s    r   r2  ExecutionStrategy.unwatch	  rL  r  c                    #    g 7fr   r   r   s    r   rk  ExecutionStrategy.discard	       r  c                    #    g7f)zF
"Unlink a key specified by ``names``"

See: ClusterPipeline.unlink()
Nr   r-  s     r   r5  ExecutionStrategy.unlink	  rL  r  c                     g r   r   r   s    r   r  ExecutionStrategy.__len__	      r   r   Nr<  r=  )r$  r%  r&  r'  r   r  r   rj   ri   r
   r  r,  r   r  r   rg   r9  r  r'  r  r2  rk  r5  r  r  r3  r   r   r   rI  rI  X	  s>     4+,8;	  FJ"?C	c  w
23	                r   rI  c            	          \ rS rSrS\SS4S jrSS jrS\\\	4   S	\
SS4S
 jrS r\S\\\	4   SS4S j5       r\ SS\S\S\\
   4S jj5       r\S 5       r\S 5       r\S 5       r\S 5       r\S 5       r\S 5       rS\4S jrSrg)AbstractStrategyi	  r  r~   Nc                     Xl         / U l        g r   )_pipe_command_queue)r   r  s     r   r	  AbstractStrategy.__init__	  s    &*
79r   r  c                    #    U R                   R                  R                  (       a,  U R                   R                  R                  5       I S h  vN   / U l        U R                   $  N7fr   )rg  r  r   r  rh  r   s    r   r  AbstractStrategy.initialize	  sK     ::$$00**++66888 zz 9s   AA)A'A)r~  r}   c                     U R                   R                  [        [        U R                   5      /UQ70 UD65        U R                  $ r   )rh  r   r  ro  rg  r  s      r   r   AbstractStrategy.execute_command	  sA     	""C 3 34FtFvF	
 zzr   c                     SR                  [        [        U5      5      nSU S[        U5       SUR                  S    3nU4UR                  SS -   Ul        g)zC
Provides extra context to the exception prior to it being handled
r8  
Command # rE  ) of pipeline caused error: r   r!  N)r  maprq   rs   r~  )r   r  numberrz  r   msgs         r   _annotate_exception$AbstractStrategy._annotate_exception	  se     hhs8W-.=#5"6 7&^^A./1 	 ).."44	r   r7  c                     g r   r   r:  s     r   r9  AbstractStrategy.mset_nonatomic	  s     	r   r  r  c                    #    g 7fr   r   r  s      r   r  AbstractStrategy.execute	  s
      	r  c                    #    g 7fr   r   r   s    r   r  AbstractStrategy.reset	  r^  r  c                     g r   r   r   s    r   r'  AbstractStrategy.multi	  rc  r   c                    #    g 7fr   r   r-  s     r   r  AbstractStrategy.watch	  r^  r  c                    #    g 7fr   r   r   s    r   r2  AbstractStrategy.unwatch	  r^  r  c                    #    g 7fr   r   r   s    r   rk  AbstractStrategy.discard	  r^  r  c                    #    g 7fr   r   r-  s     r   r5  AbstractStrategy.unlink	  r^  r  c                 ,    [        U R                  5      $ r   )ro  rh  r   s    r   r  AbstractStrategy.__len__	  s    4&&''r   )rh  rg  r<  r=  )r$  r%  r&  r'  r  r	  r  r   rj   ri   r
   r  rt  r   r   rg   r9  r,  r   r  r  r'  r  r2  rk  r5  r  r  r3  r   r   r   re  re  	  s+   :_ : :4+,8;		5 w
23	 
 FJ"?C	c 
            ( (r   re  c                      ^  \ rS rSrS\SS4U 4S jjrS\\\4   SS4S jr	 SS	\
S
\
S\\   4S jjr  SSSS\S   S	\
S
\
S\\   4
S jjrS rS rS rS rS rS rSrU =r$ )r  i
  r  r~   Nc                 $   > [         TU ]  U5        g r   )superr	  r   r  r   s     r   r	  PipelineStrategy.__init__
  s    r   r7  r  c                 \   U R                   R                  R                  n0 nUR                  5        HA  n[	        UR                  US   5      5      nUR                  U/ 5      R                  U5        MC     UR                  5        H  nU R                  " S/UQ76   M     U R                   $ )Nr   MSET)
rg  r  r   r  rM   r  
setdefaultextendr   r  )r   r7  r   slots_pairspairrp  pairss          r   r9  PipelineStrategy.mset_nonatomic
  s     **++33MMODGNN4734D""4,33D9 $ !'')E  0%0 * zzr   r  r  c                   #    U R                   (       d  / $  U R                  R                  R                  R	                  5       n  U R                  R                  R
                  (       a,  U R                  R                  R                  5       I S h  vN   U R                  U R                  R                  U R                   UUS9I S h  vN U R                  5       I S h  vN   $  NT N N	! [        R                   ac  nUS:  aV  US-  nU R                  R                  R                  5       I S h  vN    [        R                  " S5      I S h  vN     S nAO
UeS nAff = fGM%  ! U R                  5       I S h  vN    f = f7f)N)r  r  r   r!  r  )rh  rg  r  r   r  r   r  _executer  rz   r  r  r  r  )r   r  r  r  r  s        r   r  PipelineStrategy.execute
  s:     ""I	!ZZ66<<HHJN zz00<<"jj77BBDDD!%

11++'5+=	 "/ " $ **,' E$  $66 	 %) '!+"jj77>>@@@%mmD111  	  , **,s   F	/E+ AC. C(7C. C*C. F	"C,#F	(C. *C. ,F	.E%2E 4D75E EE E+ E  E%%E+ +F?F FF	r<  rz   stackr  c           
      <  #    U Vs/ s H6  oUR                   (       a!  [        UR                   [        5      (       d  M4  UPM8     nn0 nU GH4  nUR                  R	                  SS 5      nUR
                  R                  UR                  S   R                  5       5      I S h  vN n	U(       a:  UR                  U5      (       d$  UR                  U5      n
U	(       d
  [        5       n	GO#U	(       d  UR                  R                  UR                  S   5      nU(       dn  UR                  5       (       d  S nO!UR                  " UR                  6 I S h  vN nUc  [        5       n	OW[        [         R"                  [$        R"                  S9n	O0XR&                  ;   a  [        UR&                  U   S9n	O
[        5       n	UR(                  " UR                  U	R*                  US.6I S h  vN n
U
(       d  [-        SUR                   S35      eXl        [1        U
5      S:  a  [-        S	UR                   35      eU
S   nUR2                  U;  a  U/ 4X}R2                  '   X}R2                     S   R5                  U5        GM7     [6        R8                  " 5       n[:        R<                  " S
 UR?                  5        5       6 I S h  vN nURA                  5        H  u  nu  nnS nU H0  n[        UR                   [        5      (       d  M$  UR                   n  O   URB                  R                  SS5      n[E        S[6        R8                  " 5       U-
  URF                  URH                  Ub  [K        U5      OS US9I S h  vN   M     [M        U5      (       Ga  U(       a  U H  n[        UR                   [N        [P        [R        45      (       d  M/   URT                  UR.                  RV                     " URX                  " UR                  0 UR                  D6I S h  vN 5      Ul         M     U(       a  U H  nUR                   n[        U[        5      (       d  M&  SR[                  []        [^        UR                  5      5      nSUR`                  S-    S[c        U5       SUR                   3nU4UR                  SS  -   Ul        Ue   UR                  5       nUbc  UR                  UR2                  5      nUbE  US    H<  n[e        UR                   5      [f        Rh                  ;   d  M,  URk                  5           O   U Vs/ s H  oUR                   PM     sn$ s  snf  GN GN GN GN GN GNW! [         a  nUUl          S nAGM  S nAff = fs  snf 7f)Nr  r   r  r  r  r  r  r!  zToo many targets for command c              3   x   #    U  H0  n[         R                  " US    R                  US   5      5      v   M2     g7f)r   r!  N)r  r  r  r  s     r   r  ,PipelineStrategy._execute.<locals>.<genexpr>{
  s8      *D ##DG$<$<T!W$EFF*s   8:r   PIPELINEr  r8  ro  rE  rp  )6r  r  r  r}   r   r   r  r~  r  r  r  r#   r   r   r   r}  r$   r   r%   r   r  r  r`   r  ro  r  r   r  r  r  r  r   r  r   r1   r   r   r  anyre   r_   rU   r   r  r  r  rq  rq   r@  rs   r  rz   r  r  )r   r<  r  r  r  r   todor  r  r  r  r  rp  r  r  errorsr`  r  
node_errorr   r  r  rz  rs  default_cluster_noder  s                             r   r  PipelineStrategy._execute5
  s     !
 C

jY6WC5 	 
 C ZZ^^NDAN%+%<%<%D%D!!#&   f&:&:>&J&J%99.I''6'8$'#)#7#7#;#;CHHQK#HL'%6688#'D)/)?)?)J#JD</>/@,/>/</J/J0>0L0L0,
 (+H+HH/>/5/L/L$00"0, 0?/@,%+%<%<XX#3#B#B,&  
 $/;CHH:[Q  $4 < 1$+.KCHH:,VWW?Dyy%$(":ii ))Q&&s+i n ^^%
~~!LLN
 
 ,1;;='I'hJcjj)44!$J  
 ''++D!4B+'!%!1J!>#yy II(*SWD    ,9$ v;;!C!#**}j(.STT+)/)J)J # 4 4 D D*$*$:$:CHH$S

$SS*UCJ   C ZZF!&)44"%((3x+B"C()9(: ;  -g 67 899?G 
 (+fv{{12&>$   $*#:#:#< 
 $/$yy)=)B)BC+
  ,A  

+|/N/NN"779!  / ',,es

e,,I
 " $K$ (
"" T( +)*CJJ+F -s   V3U U AV#U%$B/VU(BV'U+(CVU.AV	A5V>U1?AVAU7U4
U7+/VCV0VV	V(V+V.V1V4U77
VVVVVc                    #    / U l         g7fr#  rh  r   s    r   r  PipelineStrategy.reset
  s      !s   	c                     [        S5      e)Nz@method multi() is not supported outside of transactional contextr`   r   s    r   r'  PipelineStrategy.multi
  s    #N
 	
r   c                     #    [        S5      e7f)Nz@method watch() is not supported outside of transactional contextr  r-  s     r   r  PipelineStrategy.watch
  s     #N
 	
   c                     #    [        S5      e7f)NzBmethod unwatch() is not supported outside of transactional contextr  r   s    r   r2  PipelineStrategy.unwatch
       #P
 	
r  c                     #    [        S5      e7f)NzBmethod discard() is not supported outside of transactional contextr  r   s    r   rk  PipelineStrategy.discard
  r  r  c                 h   #    [        U5      S:w  a  [        S5      eU R                  SUS   5      $ 7f)Nr!  z>unlinking multiple keys is not implemented in pipeline commandUNLINKr   )ro  r`   r  r-  s     r   r5  PipelineStrategy.unlink
  s9     u:?'P  ##HeAh77s   02r  r=  )r$  r%  r&  r'  r  r	  r   rg   ri   r9  r,  r   r
   r  r  r  r'  r  r2  rk  r5  r3  __classcell__r   s   @r   r  r  
  s    _  w
23	  GK"?C	cH  $#'K-K- %&K- 	K-
 !K- 
cK-Z!







8 8r   r  c                   ~  ^  \ rS rSrS1rSS1r1 Skr\\4r	\
\\\4rS\SS4U 4S jjrS\\\4   4S	 jrS
\\\4   S\SS4S jrS
\\\4   S\S\4U 4S jjrS rS rS rS\S\4S jrS rS r S\!\"\4   SS4S jr# S(S\$S\$S\%\   4S jjr&S\%S   S\$4S jr'S\%S   S\$4S  jr(S! r)S" r*S# r+S$ r,S% r-S& r.S'r/U =r0$ ))r  i
  UNWATCHWATCH>   EXECDISCARDr  r  r~   Nc                 P  > [         TU ]  U5        SU l        SU l        [	        5       U l        S U l        S U l        SU l        [        U R                  R                  R                  5      U l        U R                  R                  [        R                   U R"                  -   5        g )NF)r  r	  _explicit_transaction	_watchingr@  _pipeline_slots_transaction_node_transaction_connection
_executingr   rg  r  r   _retryr   rz   r  SLOT_REDIRECT_ERRORSr  s     r   r	  TransactionStrategy.__init__
  s    %*"),8<=A$4::44::;++++d.G.GG	
r   c                 d   U R                   (       d  [        S5      eU R                  R                  R                  R                  [        U R                   5      S   S5      nXl        U R                  (       d   U R                  R                  5       nX l        U R                  U R                  4$ )a?  
Find a connection for a pipeline transaction.

For running an atomic transaction, watch keys ensure that contents have not been
altered as long as the watch commands for those keys were sent over the same
connection. So once we start watching a key, we fetch a connection to the
node that owns that slot and reuse it.
z:At least a command with a key is needed to identify a noder   F)
r  r`   rg  r  r   r|  r   r  r  ra  )r   r  rF  s      r   *_get_client_and_connection_for_transaction>TransactionStrategy._get_client_and_connection_for_transaction
  s     ##'L  !JJ55CCVV%%&q)5
 "&++%)%;%;%N%N%PJ+5(%%t'C'CCCr   r~  r}   r
   c                    ^ ^^^^ S mS mUUUUU 4S jn[         R                  " US9nUR                  5         UR                  5         T(       a  TeT$ )Nc                     >  [         R                  " TR                  " T0 TD65      mg ! [         a  n U m S n A g S n A ff = fr   )r  runr  r  )r  r~  r  r}   r  r   s    r   runner3TransactionStrategy.execute_command.<locals>.runner  s;    ";;t'<'<d'Mf'MN s   &* 
A ;A )target)	threadingThreadstartr  )r   r~  r}   r  threadr  r  s   ```  @@r   r  #TransactionStrategy.execute_command  sH    	 	 !!0Kr   c                 4  >#    U R                   R                  R                  (       a,  U R                   R                  R                  5       I S h  vN   S nUS   U R                  ;  a+  U R                   R                  R
                  " U6 I S h  vN nU R                  (       d  US   U R                  ;   a  U R                  (       d  US   S:X  a  U R                  5         UbG  U R                  (       a  X0R                  ;  a  [        S5      eU R                  R                  U5        O%US   U R                  ;  a  [        SUS    S35      eU R                  " U0 UD6$ Ub  U R                  R                  U5        [        TU ]@  " U0 UD6$  GNB GN7f)Nr   r  z0Cannot watch or send commands on different slotsz)Cannot identify slot number for command: z(,it cannot be triggered in a transaction)rg  r  r   r  NO_SLOTS_COMMANDSr}  r  IMMEDIATE_EXECUTE_COMMANDSr  _validate_watchr  rZ   ri  r`   _immediate_execute_commandr  r  )r   r~  r}   slot_numberr   s       r   r  $TransactionStrategy._execute_command/  sm     ::$$00**++66888%)7$000 $

 9 9 I I4 PPK NNd1g)H)HH,,Aw'!$$&&''K?S?S,S3J  $$((5a 6 66+?Qy I> > 
 22DCFCC&$$((57*D;F;;= 9 Qs%   AFF?FFDFFc                 J    U R                   (       a  [        S5      eSU l        g )N"Cannot issue a WATCH after a MULTIT)r  ra   r  r   s    r   r  #TransactionStrategy._validate_watchS  s    %%ABBr   c                 z   ^ ^^#    T R                   R                  UUU 4S jT R                  SS9I S h  vN $  N7f)Nc                  (   > TR                   " T 0 TD6$ r   ) _get_connection_and_send_command)r~  optionsr   s   r   r   @TransactionStrategy._immediate_execute_command.<locals>.<lambda>[  s    D994K7Kr   Twith_failure_count)r  r  _reinitialize_on_error)r   r~  r  s   ```r   r  .TransactionStrategy._immediate_execute_commandY  s;     [[00K''# 1 
 
 	
 
s   /;9;c           
        #    U R                  5       u  p4U R                  (       d  UR                  U5      I S h  vN   [        R                  " 5       n U R
                  " XCUS   /UQ70 UD6I S h  vN n[        US   [        R                  " 5       U-
  UR                  UR                  [        UR                  5      S9I S h  vN   U$  N N] N
! [         ac  nXGl        [        US   [        R                  " 5       U-
  UR                  UR                  [        UR                  5      US9I S h  vN    e S nAff = f7f)Nr   r   r  r  r  r  r  )r  r  re  r  r  _send_command_parse_responser1   r   r   r  r   r  rF  )r   r~  r  
redis_noderF  r  r  r  s           r   r  4TransactionStrategy._get_connection_and_send_command`  s,    !%!P!P!R
~~11*=== ^^%
	!>>Q26:A H ,!!W!%!1J!>)&OO /   O% >  
	%L+!!W!%!1J!>)&OO /   
	se   8ECEC 2C3AC CC EC C 
EAD<5D86D<<EErF  r  c                    #    UR                   " U6 I Sh  vN   UR                  " X40 UD6I Sh  vN nX0R                  ;   a  SU l        U$  N7 N7f)z'
Send a command and parse the response
NF)rJ  r  UNWATCH_COMMANDSr  )r   rF  r  r   r~  r  outputs          r   r  0TransactionStrategy._send_command_parse_response  sS      %%t,,,!00UWUU000"DN 	-Us   AAAAAAc           
      *  #    [        US5      (       ag  [        UR                  R                  UR                  R                  UR                  R                  UR                  R                  UUSS9I S h  vN   U R
                  (       a5  [        U5      U R                  ;   a  U R                  (       a  [        S5      e[        U5      U R                  ;   d  [        U5      U R                  ;   Ga  U R                  (       a_  U R                  (       aN  U R                  R                  5       I S h  vN   U R                  R                  U R                  5        S U l        U R                  R                   =R"                  S-  sl        U R                  R                   R$                  (       a  U R                  R                   R"                  U R                  R                   R$                  -  S:X  a>  U R                  R                   R&                  R)                  5       I S h  vN   SU l        OL[+        U[,        5      (       a7  U R                  R                   R&                  R/                  U5      I S h  vN   SU l        g  GN GNY Nf N7f)NrF  Tr  z-Slot rebalancing occurred while watching keysr!  r   F)r>  r0   rF  r   r   r  r  r  r  rf   CONNECTION_ERRORSr  r  rT  rm  rg  r  r   r   r   r  r  rU   r  )r   r  r  s      r   r  *TransactionStrategy._reinitialize_on_error  s    5,''$$//44!,,11%*%5%5%:%:"'"2"2"7"7 ,    >>E{d777DOO !PQQ K4444E{d444++0F0F22==???&&..t/K/KL/3,JJ%%::a?:

))<<JJ--BB**++>>? jj//==HHJJJ,-)eX..**33AAKKERRRK( @ K SsK   A4J6J	7B<J3J4C4J(J)AJ<J=JJJJc           
        #    [        X5       H  u  pE[        U[        5      (       d  M  U R                  XER                  S-   UR
                  5        [        S[        R                  " 5       U-
  U R                  R                  U R                  R                  [        U R                  R                  5      US9I Sh  vN   Ue   g N	7f)z(
Raise the first exception on the stack
r!  TRANSACTIONr  N)r  r  r  rt  r@  r~  r1   r  r  r  r   r   r  r   )r   	responsesr  r  rr   s         r   _raise_first_error&TransactionStrategy._raise_first_error  s      )+FA!Y''((LL1,<chhG/!.%)^^%5
%B#'#?#?#D#D $ < < A A!$T%A%A%D%D!E    ,s   $CBC?C	 
Cr7  r  c                     [        S5      e)Nz1Method is not supported in transactional context.)NotImplementedErrorr:  s     r   r9  "TransactionStrategy.mset_nonatomic  s     ""UVVr   r  r  c                    #    U R                   nU(       d$  U R                  (       a  U R                  (       d  / $ U R                  X15      I S h  vN $  N7fr   )rh  r  r  !_execute_transaction_with_retries)r   r  r  r  s       r   r  TransactionStrategy.execute  s?      ##dnnD4H4HI;;ERRRRs   AAAAr  r  c                 n   ^ ^^#    T R                   R                  UU U4S jU 4S jSS9I S h  vN $  N7f)Nc                  (   > TR                  TT 5      $ r   )_execute_transaction)r  r   r  s   r   r   GTransactionStrategy._execute_transaction_with_retries.<locals>.<lambda>  s    D--e^Dr   c                 &   > TR                  X5      $ r   )r  )r  r  r   s     r   r   r    s    )D)D*r   Tr  )r  r  )r   r  r  s   ```r   r  5TransactionStrategy._execute_transaction_with_retries  s;      [[00D  $ 1 
 
 	
 
s   )535c           	      	  #    [        U R                  5      S:  a  [        S5      eSU l        U R	                  5       u  p4U R
                  (       d  UR                  U5      I S h  vN   [        [        SS5      /U[        SS5      /5      nU Vs/ s H%  n[        UR                  ;  d  M  UR                  PM'     nnUR                  U5      n[        R                  " 5       nUR                  U5      I S h  vN   / n	 UR!                  US5      I S h  vN   [-        U R.                  5       HY  u  p[        UR                  ;   a%  U	R'                  XR                  [           45        M>   UR!                  US5      I S h  vN nM[     S n UR!                  US5      I S h  vN nSU l        SU l        Uc  [5        S	5      eU	 H  u  pUR7                  X5        M     [        U5      [        U R.                  5      :w  aK  [9        S
R;                  U R.                   Vs/ s H  oUR                  S   PM     sn[        U5      5      5      eU(       d  [        U	5      S:  a%  U R=                  UU R.                  U5      I S h  vN   / n[?        UU R.                  5       H  u  nn[A        U[B        5      (       dg  UR                  S   nUU RD                  RF                  RH                  ;   a4  U RD                  RF                  RH                  U   " U40 UR                  D6nUR'                  U5        M     [K        S[        R                  " 5       U-
  URL                  URN                  [Q        URR                  5      S9I S h  vN   U$  GN(s  snf  GN GN! ["         a/  n
U R%                  U
SS5        U	R'                  U
5         S n
A
GNS n
A
fU R(                   a  nU R%                  USS5        XKl        e S nAff = f GN! U R0                   a<  nU R%                  XS-   UR                  5        U	R'                  U5         S nAGM  S nAfU R(                   a+  nU R%                  XS-   UR                  5        XKl        e S nAf["         a<  n
U R%                  XS-   UR                  5        U	R'                  U
5         S n
A
GM  S n
A
ff = f GN.! [2         a    U	(       a  U	S   ee f = fs  snf  GNq GNs7f)Nr!  zDAll keys involved in a cluster transaction must map to the same slotTr   MULTIr  r  FzWatched variable changed.zeUnexpected response length for cluster pipeline EXEC. Command stack was {} but response had length {}r  r  )*ro  r  rZ   r  r  r  re  r   r  r6   r}   r~  r  r  r  r  r  rb   rt  r   r  rF  r  rh  r  r\   rf   insertr]   formatr  r  r  r  rg  r  r   r1   r   r   r  r   )r   r  r  r  rF  cr  packed_commandsr  r  r  cluster_errorr  rz  r  
slot_errorr  datar  r   r   s                        r   r   (TransactionStrategy._execute_transaction  se     t##$q(+V  !%!P!P!R
~~11*===Q()Q'(

 %*LEq^188-KFAFFEL$228< ^^%
,,_===	++J@@@ $D$7$78JA/q.."@AB%(77
CHHA 9" 	'66z6JJH   899 DAOOA!  x=C 3 344&CCI6(,(;(;<(;1VVAY(;<c(mD  S[1_))##   (D$7$78FAsa++"xx{4::#<#<#O#OO

11DD\R ZZA KKN 9 (&!^^-
:%??"Z]]+
 	
 	
 Q > M 	> A 	$$Q73MM!%% 	$$]Aw?'1$	 I00 .,,ZQMMM*---- ,,]E7<<P/9,$ %,,QAw||DMM!$$% K 	Qi	* ="	
s0  A#S<%N&+S<N-N=<S<9N:S<N NN AS</PPP	S<S %S&S *A8S<"S1:AS<S6D S<S9S<S<N 
P%$O	S<P"O<<PS<PS0QS<S!&RS0SS<SS<S S..	S<9S<c                   #    / U l          U R                  (       a   U R                  (       aE  U R                  R                  S5      I S h  vN   U R                  R	                  5       I S h  vN   U R
                  R                  U R                  5      I S h  vN   U R                  (       a?  U R
                  (       a.  U R                  S sol        U R
                  R                  U5        S U l        S U l        SU l        SU l        [        5       U l        SU l        g  N N N! U R                   a7    U R                  (       a#  U R                  R                  5       I S h  vN     N[        R                   a6    U R                  (       a#  U R                  R                  5       I S h  vN    e f = f! U R                  (       a?  U R
                  (       a.  U R                  S sol        U R
                  R                  U5        S U l        S U l        SU l        SU l        [        5       U l        SU l        f = f7f)Nr  F)rh  r  r  rJ  rK  r  re  r  rT  r  CancelledErrorrm  r  r@  r  r  rL  s     r   r  TransactionStrategy.resetb  s     +	$ ++~~ #::GG	RRR"::HHJJJ 00EE44   ++0F0F00 9
8 &&..z:+/D(%)D""DN).D&#&5D #DOG SJ -- H33"::EEGGG--  33"::EEGGG" ++0F0F00 9
8 &&..z:+/D(%)D""DN).D&#&5D #DOs   IF; /D) D#!D) .D%/D) 3(F; D'F;  BI#D) %D) 'F; )>F8'E*(F8-F; /AF80F31F88F; ;BH??Ic                     U R                   (       a  [        S5      eU R                  (       a  [        S5      eSU l         g )Nz"Cannot issue nested calls to MULTIz:Commands without an initial WATCH have already been issuedT)r  ra   rh  r   s    r   r'  TransactionStrategy.multi  s:    %%ABBL  &*"r   c                 |   #    U R                   (       a  [        S5      eU R                  " S/UQ76 I S h  vN $  N7f)Nr  r  )r  ra   r  r-  s     r   r  TransactionStrategy.watch  s6     %%ABB))':E::::s   3<:<c                 d   #    U R                   (       a  U R                  S5      I S h  vN $ g N7f)Nr  T)r  r  r   s    r   r2  TransactionStrategy.unwatch  s(     >>--i888 9s   &0.0c                 @   #    U R                  5       I S h  vN   g  N7fr   r
  r   s    r   rk  TransactionStrategy.discard  r  r  c                 0   #    U R                   " S/UQ76 $ 7f)Nr  )r  r-  s     r   r5  TransactionStrategy.unlink  s     ##H5u55s   )	rh  r  r  r  r  r  r  r  r   r=  )1r$  r%  r&  r'  r  r  r  rU   r_   r  rY   OSErrorrW   rc   r  r  r	  r   rx   r*   r  r   rj   ri   r
   r  r  r  r  r  r  r  r  r   rg   r9  r,  r   r  r  r   r  r'  r  r2  rk  r5  r3  r  r  s   @r   r  r  
  s   "")9!55$j1	
_ 
 
D	{J&	'D6U4+;%<  PU ,"<4+,"<8;"<	"<H
!F  &' R&Ww
23W	W GKS"S?CS	cS	
+,	
>B	
u+,u>Bun.$`*;6 6r   r  c            	           \ rS rSrSrSS jrS\4S jr SS\\	   S\
S	\
S\4S
 jjrS\SS4S jrS rSS jrSS\SS4S jjrSS jrSS jrS\SS4S jrS\\\\4      4S jrSrg)_ClusterNodePoolAdapteri  u  Thin adapter exposing the :class:`ConnectionPoolInterface` that
:class:`PubSub` requires, backed by a :class:`ClusterNode`'s own
connection pool.

Connections are acquired from the node via
:meth:`ClusterNode.acquire_connection` and returned via
:meth:`ClusterNode.release`.  :meth:`PubSub.aclose` already
disconnects the connection *before* calling :meth:`release`, so the
connection is returned to the node's free-queue in a disconnected
state — guaranteeing that a subscribed socket is never silently
reused for regular commands.

Methods that do not apply to this adapter (the underlying node's
lifecycle is managed by the cluster, not by individual PubSub
instances) are implemented as no-ops so the adapter remains a valid
:class:`ConnectionPoolInterface`.
r~   Nc                 2    Xl         UR                  U l        g r   _noder   rb  s     r   r	   _ClusterNodePoolAdapter.__init__  s    
!%!7!7r   c                 6    U R                   R                  5       $ r   )r   r  r   s    r   r  #_ClusterNodePoolAdapter.get_encoder  s    zz%%''r   r   r  r  c                    #    U R                   R                  5       n UR                  5       I S h  vN   U$  N! [         a6    UR	                  5       I S h  vN    U R                   R                  U5        e f = f7fr   )r   ra  connectr  rT  rm  )r   r   r  r  rF  s        r   get_connection&_ClusterNodePoolAdapter.get_connection  st      ZZ224
		$$&&&  ' 	
 '')))JJz*	s1   A=: 8: A=: A:A"A::A=rF  c                    #    U R                   R                  U5      I S h  vN   U R                   R                  U5        g  N 7fr   )r   re  rm  rL  s     r   rm  _ClusterNodePoolAdapter.release  s6      jj--j999

:& 	:s   AA!Ac                 :    U R                   R                  SS 5      $ )Nr   )r   r   r   s    r   get_protocol$_ClusterNodePoolAdapter.get_protocol  s    %%))*d;;r   c                     g r   r   r   s    r   r  _ClusterNodePoolAdapter.reset      r   inuse_connectionsc                    #    g 7fr   r   )r   r0  s     r   rT  "_ClusterNodePoolAdapter.disconnect       r  c                    #    g 7fr   r   r   s    r   r  _ClusterNodePoolAdapter.aclose  r3  r  c                     g r   r   r  s     r   r  !_ClusterNodePoolAdapter.set_retry  r/  r   r  c                    #    g 7fr   r   )r   r  s     r   r  (_ClusterNodePoolAdapter.re_auth_callback  r3  r  c                     / $ r   r   r   s    r   get_connection_count,_ClusterNodePoolAdapter.get_connection_count  s    	r   r  r!  r   r   )T)r   r2   r~   N)r$  r%  r&  r'  r(  r	  r"   r  r   r  r
   r)   r&  rm  r+  r  r,  rT  r  r  r3   r  r   r   r  r  r;  r3  r   r   r   r  r    s    $8(W ( -1$SM9<IL	 '(: 't '<$ $ N t d5d+;&< r   r  dispatcher_refzweakref.ref[EventDispatcher]listener
event_typer~   c                 @    U " 5       nUb  UR                  X!/05        g g r   )unregister_listeners)r=  r>  r?  
dispatchers       r    _unregister_slots_cache_listenerrC    s*      !J''Z(@A r   c                   4    \ rS rSrSrS	S jrS\SS4S jrSrg)
ClusterPubSubSlotsCacheListeneri
  aX  
Async listener that forwards AsyncAfterSlotsCacheRefreshEvent to a
ClusterPubSub.

Holds a weak reference to the pubsub so it does not keep the instance
alive. Deterministic cleanup of the dispatcher's strong reference to this
listener is performed by a ``weakref.finalize`` attached to the owning
ClusterPubSub in ``ClusterPubSub.__init__``.
r~   Nc                 :    [         R                  " U5      U l        g r   )weakrefref_pubsub_ref)r   r  s     r   r	  (ClusterPubSubSlotsCacheListener.__init__  s    9@V9Lr   eventc                    #    U R                  5       nUc  g  UR                  5       I S h  vN   g  N! [         a5  n[        R	                  SU[        U5      R                  U5         S nAg S nAff = f7f)Nz2pubsub %r raised during slots-cache change: %s: %s)rI  on_slots_changedr  rr  r  r  r$  )r   rK  r  r  s       r   listen&ClusterPubSubSlotsCacheListener.listen  sn     !!#> 
	))+++ 	 DQ  	 	s6   A53 13 A53 
A2+A-(A5-A22A5)rI  )r  r  r~   N)	r$  r%  r&  r'  r(  r	  r-  rN  r3  r   r   r   rE  rE  
  s     M& T r   rE  c                   z  ^  \ rS rSrSr     S-SSS\S   S\\   S	\\   S
\\   S\\	   S\
SS4U 4S jjjr   S.SSS\S   S\\   S	\\   SS4
S jjrS\S   4S jrS/S jrSSS\4S jrS\S\\   4S jr S0S\S\\\   \\\\
4      4   4S jjrS\\SS4   4S jr   S1S\S\S\S   S\\\\
4      4S jjrS\\-  S\SS4S jrS\
SS4S jrS/S jrS \
S!\\   S"\\   S#SSS4
S$ jr S/S% jr!\"S2S& j5       r#S\S'   4S( jr$S/U 4S) jjr%SSS\S   S\\   S	\\   SS4
S* jr&S\
S\
S\
4U 4S+ jjr'S,r(U =r)$ )3r  i+  a  
Async cluster implementation for pub/sub.

IMPORTANT: before using ClusterPubSub, read about the known limitations
with pubsub in Cluster mode and learn how to workaround them:
https://redis.readthedocs.io/en/stable/clustering.html#known-pubsub-limitations
Nredis_clusterrz   r  rx   r   r   push_handler_funcr   r}   r~   c                   > SU l         U R                  XX45        U R                   b  [        U R                   5      nOSnXl        0 U l        0 U l        [        R                  " 5       U l        [        5       U l
        U R                  5       U l        Uc  [        5       U l        OX`l        [        T
U ]<  " SUUR                   UU R                  S.UD6  UR"                  R                  n	[%        U 5      U l        U	R)                  [*        U R&                  /05        [,        R.                  " U [0        [,        R2                  " U	5      U R&                  [*        5        g)a  
When a pubsub instance is created without specifying a node, a single
node will be transparently chosen for the pubsub connection on the
first command execution. The node will be determined by:
 1. Hashing the channel name in the request to find its keyslot
 2. Selecting a node that handles the keyslot: If read_from_replicas is
    set to true or load_balancing_strategy is set, a replica can be selected.

:param redis_cluster: RedisCluster instance
:param node: ClusterNode to connect to
:param host: Host of the node to connect to
:param port: Port of the node to connect to
:param push_handler_func: Optional push handler function
:param event_dispatcher: Optional event dispatcher
:param kwargs: Additional keyword arguments
Nconnection_poolr   rR  r   r   )r  set_pubsub_noder  clusternode_pubsub_mapping_shard_channel_to_noder  r/   _shard_state_lockr@  _reconcile_tasks_pubsubs_generatorrT   r   r  r	  r   r   rE  _slots_cache_listenerregister_listenersrR   rG  finalizerC  rH  )r   rQ  r  r   r   rR  r   r}   rU  nm_dispatcherr   s             r   r	  ClusterPubSub.__init__4  s5   4 	]$= 99 5dii@O"O$68  79# 07||~365"&"9"9";#%4%6D"%5" 	
+!))/!33		

 	
 &33EE%DT%J"((-0J0J/KL	

 	,KK&&&,	
r   rW  c                     Ub*  U R                  XUR                  UR                  5        UnO=Ub'  Ub$  UR                  X4S9nU R                  XX45        UnOUc  Ub  [	        S5      eSnXPl        g)a  
The pubsub node will be set according to the passed node, host and port
When none of the node, host, or port are specified - the node is set
to None and will be determined by the keyslot of the channel in the
first command to be executed.
RedisClusterException will be thrown if the passed node does not exist
in the cluster.
If host is passed without port, or vice versa, a DataError will be
thrown.
Nr  zSpecify both host and port)_raise_on_invalid_noder   r   ra  r[   r  )r   rW  r  r   r   pubsub_nodes         r   rV  ClusterPubSub.set_pubsub_node  s    " ''tyy$))LK$"2###9D''tBK!1899 K	r   c                     U R                   $ )z
Get the node that is being used as the pubsub connection.

:return: The ClusterNode being used for pubsub, or None if not yet determined
)r  r   s    r   get_pubsub_nodeClusterPubSub.get_pubsub_node  s     yyr   c                 :  #    [        [        5      nU R                  R                  5        H.  u  p#X1[	        U R
                  R                  U5      5         U'   M0     UR                  5        H&  nU R                  X@R                  5      I S h  vN   M(     g  N	7fr   )
r   r  shard_channelsr  rM   r   r  r   _resubscribe
ssubscribe)r   by_slotkvsubscriptionss        r   _resubscribe_shard_channels)ClusterPubSub._resubscribe_shard_channels  sy     
 +6d*;''--/DA;<HT\\00345a8 0$^^-M##M??CCC .Cs   BBB
Bc                 B    U R                   UR                     $ ! [         ay    [        [	        U5      U R
                  R                  U R                  U R                  S9n[        [        R                  U5      Ul        X R                   UR                  '   Us $ f = f)z3Get or create a PubSub instance for the given node.rT  )rX  r  KeyErrorr'   r  rW  r   rR  r   r   r  rq  )r   r  r  s      r   _get_node_pubsubClusterPubSub._get_node_pubsub  s    	++DII66 	 7 =,,"&"8"8!%!7!7	F 2<9962F. 39$$TYY/M	s    B BBr  c                 Z    U R                   R                  5        H  u  p#X1L d  M  Us  $    g r   )rX  r  )r   r  r  	candidates       r   _find_node_name_for_pubsub(ClusterPubSub._find_node_name_for_pubsub  s.    #77==?OD"  @ r   r  c                    #    [        [        U R                  5      5       H8  n[        U R                  5      nUR                  SUS9I Sh  vN nUc  M5  X44s  $    g N7f)z7Generate messages from shard channels across all nodes.Fr  r  Nr  )r  ro  rX  rZ  r\  get_message)r   r  r  r  r=  s        r   _sharded_message_generator(ClusterPubSub._sharded_message_generator  sl      s43345A$112F #..*/ /  G "& 6 s   A
A AA 	A c              #   ~   #     [        U R                  R                  5       5      nU(       d  gU Sh  vN   M7   N7f)z>Generator that yields PubSub instances in round-robin fashion.N)r   rX  r   )r   current_nodess     r   r\   ClusterPubSub._pubsubs_generator  s:      !9!9!@!@!BCM $$$	  %s   2=;=r  r  c                   #    U(       aH  U R                   R                  UR                  5      nU(       a  UR                  SUS9I Sh  vN nOSnOU R	                  US9I Sh  vN u  pEUc  g[        US   5      S:X  a  U R                   ISh  vN   US   U R                  ;   a\  U R                  R                  US   5        U R                  R                  US   S5        U R                  R                  US   S5        UbZ  UR                  (       dI  U R                  U5      nUb5   UR                  5       I Sh  vN   U R                   R                  US5        SSS5      ISh  vN   [        US   5      S;   a  U R                   (       d  U(       a  gU$  GNV GN> GN Nf! [         a     Npf = f NL! , ISh  vN  (       d  f       Na= f7f)	z
Get a message from shard channels.

:param ignore_subscribe_messages: Whether to ignore subscribe messages
:param timeout: Timeout for message retrieval
:param target_node: Specific node to get message from
:return: Message dictionary or None
Fr|  N)r  r  sunsubscribechannel)rl  r  )rX  r   r  r}  r~  rr   rZ  "pending_unsubscribe_shard_channelsrt  rj  r   rY  
subscribedry  r  r  r  )r   r  r  r  r  r=  r  s          r   get_sharded_message!ClusterPubSub.get_sharded_message  s     --11+2B2BCF !' 2 2.3W !3 !  $($C$CG$C$TTOF? (N:
 ---9%)P)PP;;BB79CUV''++GI,>E//33GI4FM %f.?.?::6BD'!"(--/11 0044T4@1 .-6 (,JJ--1Ja U .* 2( ! !- .---s   AG	F
G$F"%,GF%GBF<.F*F(F*F<"G-F:.2G"G%G(F**
F74F<6F77F<:G<GGGGr~  c           
        #    [        X5      nU R                   ISh  vN   UR                  5        GH\  u  pEU R                  R	                  U5      nU(       d  M*  [        [        U R                  US05      5      5      nU R                  R                  U5      nU(       a-  XR                  :w  a  U R                  UUUU5      I Sh  vN   M  U R                  U5      n	U(       a#  U	R                  [        XE5      5      I Sh  vN   OU	R                  U5      I Sh  vN   U R                  R!                  U	R                  5        UR                  U R                  U'   U R"                  R%                  U R                  US05      5        GM_     SSS5      ISh  vN   g GN N N N N! , ISh  vN  (       d  f       g= f7f)z
Subscribe to shard channels.

:param args: Channel names or ``Subscription`` objects
:param kwargs: Channel names with handlers
N)rI   rZ  r  rW  rs  rZ  iter_normalize_keysrY  r   r  _migrate_shard_channelru  rl  rl   rj  r   r  difference_update)
r   r~  r}   
s_channels	s_channelhandlerr  normalized_keyold_namer  s
             r   rl  ClusterPubSub.ssubscribe/  sy     0=
 )))&0&6&6&8"	||55i@ "&d4+?+?D@Q+R&S!T66::>JII 5 55& 	   ..t4 ++L,LMMM ++I666##**6+@+@A>Bii++N;77II(()T):;7 '9 *)) N63 *)))s|   G
F%G
B(F0
F(;F0F*F0 F,!A3F0G
F. G
(F0*F0,F0.G
0G6F97GG
c           
        #    U(       a  [        US   USS 5      nO#[        U R                  R                  5       5      nU R                   ISh  vN   U H  n[        [        U R                  US05      5      5      nU R                  R                  U5      nU(       a  X@R                  ;   a  U R                  U   nOWU R                  R                  U5      nU(       a  UR                  U R                  ;  a  M  U R                  UR                     nUR                  U5      I Sh  vN   U R                  R!                  UR                  5        GM     SSS5      ISh  vN   g GN NB N
! , ISh  vN  (       d  f       g= f7f)zs
Unsubscribe from shard channels.

:param args: Channel names to unsubscribe from. If empty, unsubscribe from all.
r   r!  N)rH   r   rj  r  rZ  rZ  r  r  rY  r   rX  rW  rs  r  r  r  r   )r   r~  r  r  r  r  r  s          r   r  ClusterPubSub.sunsubscribe^  s4     Qab2D++0023D
 )))!	!%d4+?+?D@Q+R&S!T 2266~FD$<$<<!55d;F<<99)DD499D4L4L#L !55dii@F)))44477>>== " *)) 5 *)))s[   AFE+FCE2+E.,.E2F%E0&F.E20F2F	8E;9F	Fc           
        #    / nSnSnU R                    ISh  vN   [        U R                  R                  5       5       Hj  u  pE U R                  R                  U5      nU R                  R                  U5      nXvR                  :X  a  MM   U R                  XEXv5      I Sh  vN   SnMl     [        U R&                  R                  5       5       HM  u  pU
R(                  (       a  M   U
R+                  5       I Sh  vN   U R&                  R/                  U	S5        MO     SSS5      ISh  vN   U(       a  [        [1        U5       SU< 35      eUb
  U(       d  Uegg GNI! [         a    UR                  U5         GMC  f = f N! [        [        [        4 a<  n[        R!                  SU[#        U5      R$                  U5        Uc  Un SnAGM  SnAff = f N! [,         a     Nf = f N! , ISh  vN  (       d  f       N= f7f)a3  
Reconcile per-node shard subscriptions against the cluster's current
slot ownership map. For each tracked shard channel whose owning node
has changed (e.g. after CLUSTER SETSLOT / failover), sunsubscribe on
the old node's pubsub and ssubscribe on the new owner's pubsub,
preserving any registered handler.
FNTz+shard channel %r migration deferred: %s: %szI shard channel(s) left unreconciled; slot(s) not covered by the cluster: )rZ  r   rj  r  rW  rs  rc   r   rY  r   r  r  rY   rd   r  rr  warningr  r$  rX  r  r  r  r   ro  )r   	uncoveredmade_progressfirst_migrate_errorr  r  new_noder  r  r  r  s              r    reinitialize_shard_subscriptions.ClusterPubSub.reinitialize_shard_subscriptions  s     	7;)))$()<)<)B)B)D$E 
#||==gFH  66::7C}},55(   %)M' %FL !%T%=%=%C%C%E F((($mmo-- ,,00t< !GO *)\  &y>" #77@mE  *= &% 4A*m * +  $$W- (w?  NNEQ((	 +2./+* .$ W *)))s   HE"H)G2E%",G2F
$F%F
+;G2+G >G?G  G2#H.G0/4H%F G2FG2F

G0GG2GG2G  
G-*G2,G--G20H2H	8G;9H	Hr  r  r  r  c                 Z  #    U(       a8  X0R                   ;   a)  U R                   U   n UR                  U5      I S h  vN   U R                  U5      nU(       a#  UR                  [        X5      5      I S h  vN   OUR                  U5      I S h  vN   U R                  R                  UR                  5        [        [!        U R#                  US 05      5      5      nUR$                  U R&                  U'   U R(                  R+                  U R#                  US 05      5        g  N! [        [        [        4 ae    U R
                  R                  US9cG   UR                  5       I S h  vN    O! [         a     Of = fU R                   R                  US 5         GN_f = f GN. GN7f)Nr_  )rX  r  rY   rd   r  rW  ra  r  r  r   ru  rl  rl   rj  r   rZ  r  r  r  rY  r  r  )r   r  r  r  r  
old_pubsub
new_pubsubr  s           r   r  $ClusterPubSub._migrate_shard_channel  st     $<$<<11(;JA --g6666 **84
''W(FGGG''000"":#<#<=d4#7#7$#HIJ6>mm##N3//AA  '41	
G 7#\7; A$ <<((8(<D(//111$ ,,004@/A8 H0s   &F+D( D&D( 5F+7F%8F+F(BF+&D( (/F"E2+E.,E21F"2
E?<F">E??F"F+!F""F+(F+c                 .  #    U R                   (       d  g [        R                  " U R                  5       5      nU R                  R                  U5        UR                  U R                  R                  5        UR                  U R                  5        g 7fr   )	rj  r  r  r  r[  ri  rj  rk  _log_reconcile_task_exception)r   rl  s     r   rM  ClusterPubSub.on_slots_changed  sp      """"4#H#H#JK!!$'t44<<= 	tAABs   BBc                     U R                  5       (       a  g U R                  5       nUb  [        R                  SXS9  g g )Nz,shard subscription reconciliation failed: %rrp  )	cancelledr  rr  r  )rl  r[  s     r   r  +ClusterPubSub._log_reconcile_task_exception  s>    >>nn?LL>   r   r)   c                     U R                   $ )a1  
Get the Redis connection of the pubsub connected node.

Returns the pubsub's dedicated connection (acquired from its own
connection pool), not from the ClusterNode's connection pool.
This avoids the connection pool resource leak that would occur
if we called node.acquire_connection() without releasing.
)rF  r   s    r   get_redis_connection"ClusterPubSub.get_redis_connection  s     r   c                   >#    U R                   (       aL  [        U R                   5      nU H  nUR                  5         M     [        R                  " USS06I Sh  vN   U R
                   ISh  vN   U R                   R                  5         U R                  R                  5        H  nUR                  5       I Sh  vN   M     U R                  R                  5         [        U 5      R                  U 5      U l        [        TU ]%  5       I Sh  vN   U R                  R                  5         SSS5      ISh  vN   g N N N N5 N! , ISh  vN  (       d  f       g= f7f)z#
Disconnect the pubsub connection.
rX  TN)r[  r   cancelr  r  rZ  clearrX  r   r  r  r\  r  rY  )r   tasksrl  r  r   s       r   r  ClusterPubSub.aclose)  s       ../E ..%@4@@@ )))!!'')2299;mmo%% <
 $$**, '+4j&C&C'D# '.""" ''--/7 *)) A * && #/ *)))s   AE'EE'1E2E'5AE EAEE	E2E'=E>E'E'E	EE'E$EE$ E'c                 b    Ub  UR                  UR                  S9c  [        SU SU S35      eg)zT
Raise a RedisClusterException if the node is None or doesn't exist in
the cluster.
Nr_  zNode :z doesn't exist in the cluster)ra  r  r`   )r   rQ  r  r   r   s        r   rc  $ClusterPubSub._raise_on_invalid_node^  sF     <=11DII1FN'vQtf$AB  Or   c                   >#    U(       a  US   R                  5       OSnUS;   aa  [        U5      S:  aR  US   nU R                  R                  U5      nU(       a+  U R	                  U5      nUR
                  " U0 UD6I Sh  vN $ U R                  c  U R                  c  [        U5      S:  ap  US   nU R                  R                  U5      nU R                  R                  R                  UU R                  R                  U R                  R                  5      nOU R                  R                  5       nXPl        [        U5      U l        [         TU ]  " U0 UD6I Sh  vN $  N N7f)z|
Execute a command on the appropriate cluster node.

Taken code from redis-py and tweaked to make it work within a cluster.
r   r  )
SSUBSCRIBESUNSUBSCRIBESPUBLISHr!  N)r  ro  rW  rs  ru  r  rF  rU  rm  r   r|  r   r   rZ  r  r  r  )	r   r~  r}   rz  r  r  r  rp  r   s	           r   r  ClusterPubSub.execute_commandn  s<     &*$q'--/r@@4y1}q'||55g>!2248F!'!7!7!H!HHH ??"##+t9q= #1gG<<//8D<<55HH77<<D  <<779D 	'>t'D$ W,d=f===- I, >s%   A?E,E(C E,#E*$E,*E,)
r   r\  r[  rY  rZ  r]  rW  rU  r  rX  )NNNNNr"  r   )        )Fr  N)rl  zasyncio.Taskr~   N)*r$  r%  r&  r'  r(  r   r  r  r   rT   r
   r	  rV  rg  rq  r'   ru  ry  r.  r   r   r~  r   r\  r,  r  rh   rl   rk   rl  r  r  r  rM  staticmethodr  r  r  rc  r  r3  r  r  s   @r   r  r  +  s    )-""046:M
%M
 }%M
 sm	M

 smM
 $H-M
 #?3M
 M
 
M
 M
d )-""     }%   sm	  
 sm   
  D-!8 	D] v ( HSM   #	x$sCx.!99	:%IfdD.@$A % +0/3	D#'D D m,	D
 
$sCx.	!DL-,-8E-	-^  BI&V1
1
 (#1
 3-	1

  1
 
1
fC&  h/C&D 30j% }% sm	
 sm 
 (>3 (># (># (> (>r   r  )r  r>  loggingrX  r<  r  r  r0  rG  abcr   r   r   r   	itertoolsr   typesr   typingr	   r
   r   r   r   r   r   r   r   r   r   r   r   r   r   r   r  r   redis._defaultsr   r   r   r   r   r    redis._parsersr!   r"   redis._parsers.commandsr#   r$   r%   redis._parsers.helpersr&   redis.asyncio.clientr'   r(   redis.asyncio.connectionr)   r*   r+   r,   r-   redis.asyncio.lockr/   $redis.asyncio.observability.recorderr0   r1   redis.asyncio.retryr2   redis.auth.tokenr3   redis.backoffr4   r5   redis.clientr6   r7   r8   redis.clusterr9   r:   r;   r<   r=   r>   r?   r@   rA   rB   rC   rD   rE   redis.commandsrF   rG   redis.commands.helpersrH   rI   redis.commands.policiesrJ   rK   	redis.crcrL   rM   redis.credentialsrN   redis.driver_inforO   rP   redis.eventrQ   rR   rS   rT   redis.exceptionsrU   rV   rW   rX   rY   rZ   r[   r\   r]   r^   r_   r`   ra   rb   rc   rd   re   rf   redis.typingrg   rh   ri   rj   rk   rl   redis.utilsrm   rn   ro   rp   rq   rr   rs   r   rt   ru   rv   	getLoggerr$  rr  r  rw   rz   rx   r   r  rz  replacer  setattrr  rI  re  r  r  r  r-  rC  rE  r  r   r   r   <module>r     sw            # #       &   8 R R 9 :  $ & + A D D    D K R 8 0 =     (    77JJK			8	$C](;T#}BT=U
]="68Q ]@)l l^	w wtf@m%9;T f@R )Gooc3'--/G""OW&<W&EF )	@ 	@a aHC(( C(L`8' `8FF6* F6RL5 L^
B2
B)
B V
B 
	
B&A Bk	>F k	>r   