
    ng\jYT                        S SK r S SKrS SKJrJrJrJrJrJrJ	r	  S SK
Jr  S SKJrJrJrJr  S SKJrJrJr  S SKJ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!J"r"  S SK#J$r$  S SK#J%r&  S SK'J(r(J)r)J*r*  S SK+J,r,  S SK-J.r.J/r/J0r0J1r1J2r2  S SK3J4r4  \Rj                  " \65      r7\4 " S S\"\!5      5       r8S\$4S jr9 " S S\"\!5      r: " S S5      r;g)    N)Any	AwaitableCallableListLiteralOptionalUnion)DefaultCommandExecutor)DEFAULT_GRACE_PERIODDatabaseConfigInitialHealthCheckMultiDbConfig)AsyncDatabaseDatabase	Databases)AsyncFailureDetector)HealthCheckHealthCheckPolicy)Retry)BackgroundScheduler)	NoBackoff)AsyncCoreCommandsAsyncRedisModuleCommands)CircuitBreaker)State)InitialHealthCheckFailedErrorNoValidDatabaseExceptionUnhealthyDatabaseException)GeoFailoverReason)ChannelT
EncodableTKeyTPubSubHandlerSubscription)experimentalc                   r   \ rS rSrSrS\4S jrS.S jrS rS r	S	 r
S\4S
 jrS\SS4S jr S/S\S\4S jjrS\S\4S jrS\4S jrS\S\4S jrS\4S jrS\4S jrS rS rSSSS.S\S/\\\\   4   4   S \ S!\!\"   S"\S#\!\   4
S$ jjr#S% r$S\%\&\4   4S& jr'S' r(S\S\4S( jr)S)\*S*\+S+\+4S, jr,S-r-g)0MultiDBClient!   z~
Client that operates on multiple logical Redis databases.
Should be used in Client-side geographic failover database setups.
configc                    UR                  5       U l        UR                  (       d  UR                  5       OUR                  U l        UR
                  U l        UR                  R                  5       U l	        UR                  (       d  UR                  5       OUR                  U l        UR                  c  UR                  5       OUR                  U l        U R                  R!                  U R                  5        UR"                  U l        UR&                  U l        UR*                  U l        U R,                  R/                  [0        /5        [3        U R                  U R                  U R,                  U R                  UR4                  UR6                  U R(                  U R$                  S9U l        SU l        [<        R>                  " 5       U l         [C        5       U l"        Xl#        S U l$        / U l%        S U l&        g )N)failure_detectors	databasescommand_retryfailover_strategyfailover_attemptsfailover_delayevent_dispatcherauto_fallback_intervalF)'r,   
_databaseshealth_checksdefault_health_checks_health_checkshealth_check_interval_health_check_intervalhealth_check_policyvalue_health_check_policyr+   default_failure_detectors_failure_detectorsr.   default_failover_strategy_failover_strategyset_databasesr2   _auto_fallback_intervalr1   _event_dispatcherr-   _command_retryupdate_supported_errorsConnectionRefusedErrorr
   r/   r0   command_executorinitializedasyncioLock_hc_lockr   _bg_scheduler_config_recurring_hc_task	_hc_tasks_half_open_state_task)selfr)   s     W/home/edenadmin/noVNC/venv/lib/python3.13/site-packages/redis/asyncio/multidb/client.py__init__MultiDBClient.__init__(   s    **, '' ((*%% 	
 '-&B&B#&&,,. 	!
 ++ ,,.)) 	 ''/ ,,.)) 	
 	--doo>'-'D'D$!'!8!8$22335K4LM 6"55oo--"55$66!00!33#'#?#?	!
 !02"&%)"    returnc                 d   #    U R                   (       d  U R                  5       I S h  vN   U $  N7fN)rG   
initializerP   s    rQ   
__aenter__MultiDBClient.__aenter__U   s(     //### $s   %0.0c                   #    U R                   (       a  U R                   R                  5         U R                  (       a  U R                  R                  5         U R                   H  nUR                  5         M     U R                  R                  5       I S h  vN   U R                  R                  (       a7  U R                  R                  R                  R                  5       I S h  vN   g g  NW N7frW   )
rM   cancelrO   rN   r;   closerF   active_databaseclientaclose)rP   hc_tasks     rQ   ra   MultiDBClient.acloseZ   s     ""##**,%%&&--/~~GNN & ''--///   00''77>>EEGGG 1 	0 Hs%   BC5C1AC5*C3+C53C5c                 @   #    U R                  5       I S h  vN   g  N7frW   ra   rP   exc_type	exc_value	tracebacks       rQ   	__aexit__MultiDBClient.__aexit__j        kkm   c                   #    U R                  5       I Sh  vN   [        R                  " U R                  R	                  U R
                  U R                  5      5      U l        SnU R                   Ho  u  p#UR                  R                  U R                  5        UR                  R                  [        R                  :X  d  MT  U(       a  M]  X R                  l        SnMq     U(       d  [#        S5      eSU l        g N7f)zD
Perform initialization of databases to define their initial state.
NFTz4Initial connection failed - no active database found)_perform_initial_health_checkrH   create_taskrK   run_recurring_asyncr8   _check_databases_healthrM   r3   circuiton_state_changed!_on_circuit_state_change_callbackstateCBStateCLOSEDrF   _active_databaser   rG   )rP   is_active_db_founddatabaseweights       rQ   rX   MultiDBClient.initializem   s      00222 #*"5"522++,,#
 # $H--d.T.TU %%7@R@R :B%%6%)" !0 "*F   9 	3s   DC?B+DD1Dc                     U R                   $ )z5
Returns a sorted (by weight) list of all databases.
)r3   rY   s    rQ   get_databasesMultiDBClient.get_databases   s     rT   r{   Nc                   #    SnU R                    H  u  p4X1:X  d  M  Sn  O   U(       d  [        S5      eU R                  U5      I Sh  vN   UR                  R                  [
        R                  :X  aS  U R                   R                  S5      S   u  pTU R                  R                  U[        R                  5      I Sh  vN   g[        S5      e N N7f)z<
Promote one of the existing databases to become an active.
NT/Given database is not a member of database list   r   z1Cannot set active database, database is unhealthy)r3   
ValueError_check_db_healthrs   rv   rw   rx   	get_top_nrF   set_active_databaser   MANUALr   )rP   r{   existsexisting_db_highest_weighted_dbs         rQ   r   !MultiDBClient.set_active_database   s      "ooNK& .
 NOO##H---!!W^^3%)__%>%>q%A!%D"'';;+22   &?
 	
 	.s)   C,C	C
A9CCCCskip_initial_health_checkc                   #    UR                   R                  S[        S[        5       S905        UR                  (       a<  U R
                  R                  R                  " UR                  40 UR                   D6nOUR                  (       aY  UR                  R                  [        S[        5       S95        U R
                  R                  R                  UR                  S9nO&U R
                  R                  " S0 UR                   D6nUR                  c  UR                  5       OUR                  n[        UUUR                  UR                  S9n U R                  U5      I Sh  vN   U R                   R#                  S5      S   u  pgU R                   R%                  XUR                  5        U R'                  XV5      I Sh  vN   g Nc! [         a    U(       d  e  Nuf = f N7f)	z
Adds a new database to the database list.

Args:
    config: DatabaseConfig object that contains the database configuration.
    skip_initial_health_check: If True, adds the database even if it is unhealthy.
retryr   )retriesbackoff)connection_poolN)r`   rs   r|   health_check_urlr    )client_kwargsupdater   r   from_urlrL   client_class	from_pool	set_retryrs   default_circuit_breakerr   r|   r   r   r   r3   r   add_change_active_database)rP   r)   r   r`   rs   r{   r   highest_weights           rQ   add_databaseMultiDBClient.add_database   s     	##WeAy{.S$TU??\\..77#)#7#7F &&uQ	'LM\\..88 & 0 0 9 F \\..F1E1EFF ~~% **, 	 ==#44	
	''111
 /3oo.G.G.J1.M+Hoo6**8III 2) 	, -	 	JsI   EG+G +G,G 0AG+	G)
G+G G&#G+%G&&G+new_databasehighest_weight_databasec                    #    UR                   UR                   :  a\  UR                  R                  [        R                  :X  a3  U R
                  R                  U[        R                  5      I S h  vN   g g g  N7frW   )	r|   rs   rv   rw   rx   rF   r   r   	AUTOMATIC)rP   r   r   s      rQ   r   %MultiDBClient._change_active_database   sn      "9"@"@@$$**gnn<'';;/99   = As   A0A<2A:3A<c                 H  #    U R                   R                  U5      nU R                   R                  S5      S   u  p4XB::  a\  UR                  R                  [
        R                  :X  a3  U R                  R                  U[        R                  5      I Sh  vN   ggg N7f)z,
Removes a database from the database list.
r   r   N)r3   remover   rs   rv   rw   rx   rF   r   r   r   )rP   r{   r|   r   r   s        rQ   remove_databaseMultiDBClient.remove_database   s      ''1.2oo.G.G.J1.M+ $#++11W^^C'';;#%6%=%=   D %s   BB"B B"r|   c                 (  #    SnU R                    H  u  pEXA:X  d  M  Sn  O   U(       d  [        S5      eU R                   R                  S5      S   u  pgU R                   R                  X5        X!l        U R                  X5      I Sh  vN   g N7f)z,
Updates a database from the database list.
NTr   r   r   )r3   r   r   update_weightr|   r   )rP   r{   r|   r   r   r   r   r   s           rQ   update_database_weight$MultiDBClient.update_database_weight   s      "ooNK& .
 NOO.2oo.G.G.J1.M+%%h7 **8IIIs   BA-B
BBfailure_detectorc                 :    U R                   R                  U5        g)z.
Adds a new failure detector to the database.
N)r=   append)rP   r   s     rQ   add_failure_detector"MultiDBClient.add_failure_detector  s     	&&'78rT   healthcheckc                    #    U R                    ISh  vN   U R                  R                  U5        SSS5      ISh  vN   g N0 N! , ISh  vN  (       d  f       g= f7f)z*
Adds a new health check to the database.
N)rJ   r6   r   )rP   r   s     rQ   add_health_checkMultiDBClient.add_health_check  s3      ===&&{3 !=====sA   A"AA"AA"AA"A"AAAA"c                    #    U R                   (       d  U R                  5       I Sh  vN   U R                  R                  " U0 UD6I Sh  vN $  N( N7f)z2
Executes a single command and return its result.
N)rG   rX   rF   execute_commandrP   argsoptionss      rQ   r   MultiDBClient.execute_command  sG      //###**::DLGLLL $Ls!   %AA#AAAAc                     [        U 5      $ )z*
Enters into pipeline mode of the client.
)PipelinerY   s    rQ   pipelineMultiDBClient.pipeline&  s     ~rT   F
shard_hintvalue_from_callablewatch_delayfuncr   watchesr   r   r   c                   #    U R                   (       d  U R                  5       I Sh  vN   U R                  R                  " U/UQ7UUUS.6I Sh  vN $  N. N7f)z#
Executes callable as transaction.
Nr   )rG   rX   rF   execute_transaction)rP   r   r   r   r   r   s         rQ   transactionMultiDBClient.transaction,  sc      //###**>>

 " 3#
 
 	
 $
s!   %AA)AAAAc                 x   #    U R                   (       d  U R                  5       I Sh  vN   [        U 40 UD6$  N7f)z
Return a Publish/Subscribe object. With this object, you can
subscribe to channels and listen for messages that get published to
them.
N)rG   rX   PubSub)rP   kwargss     rQ   pubsubMultiDBClient.pubsubB  s5      //###d%f%% $s   %:8:c                   #    0 n/ U l         U R                   HI  u  p#[        R                  " U R	                  U5      5      nX!U'   U R                   R                  U5        MK     [        R                  " U R                   SS06I Sh  vN n[        U R                   U5       VVs0 s H
  u  pFX   U_M     nnnUR                  5        Hi  u  p&[        U[        5      (       d  M  UR                  n[        R                  UR                  l        [         R#                  SUR$                  S9  SXx'   Mk     U$  Ns  snnf 7f)zS
Runs health checks as a recurring task.
Runs health checks against all databases.
return_exceptionsTNz%Health check failed, due to exception)exc_infoF)rN   r3   rH   rp   r   r   gatherzipitems
isinstancer   r{   rw   OPENrs   rv   loggerdebugoriginal_exception)	rP   
task_to_dbr{   r   taskresultsresult
db_resultsunhealthy_dbs	            rQ   rr   %MultiDBClient._check_databases_healthM  s#    
 46
??KH&&t'<'<X'FGD'tNN!!$' +
  O$OO :=T^^W9U
9UJf$9U 	 
 !+ 0 0 2H&"<==%-4\\$$*;#66  
 ,1
( !3 ' P
s+   BED9	E&D;7+E&AE;Ec                 &  #    U R                  5       I Sh  vN nSnU R                  R                  [        R                  :X  a  SUR                  5       ;  nOU R                  R                  [        R                  :X  a)  [        UR                  5       5      [        U5      S-  :  nO;U R                  R                  [        R                  :X  a  SUR                  5       ;   nU(       d"  [        SU R                  R                   35      eg N7f)zZ
Runs initial health check and evaluate healthiness based on initial_health_check_policy.
NTF   z:Initial health check failed. Initial health check policy: )rr   rL   initial_health_check_policyr   ALL_AVAILABLEvaluesMAJORITY_AVAILABLEsumlenONE_AVAILABLEr   )rP   r   
is_healthys      rQ   ro   +MultiDBClient._perform_initial_health_checko  s      4466
<<337I7W7WWgnn&66JLL44!445 W^^-.W1AAJLL448J8X8XX!11J/LT\\MuMuLvw   7s   DDC9Dc                   #    U R                   R                  U R                  U5      I Sh  vN nU(       dI  UR                  R                  [
        R                  :w  a  [
        R                  UR                  l        U$ U(       aG  UR                  R                  [
        R                  :w  a  [
        R                  UR                  l        U$  N7f)z?
Runs health checks on the given database until first failure.
N)r;   executer6   rs   rv   rw   r   rx   )rP   r{   r   s      rQ   r   MultiDBClient._check_db_health  s     
  44<<
 

 %%5)0  &H,,22gnnD%,^^H"
s   *CCB%Crs   	old_state	new_statec                 &   [         R                  " 5       nU[        R                  :X  a5  [         R                  " U R                  UR                  5      5      U l        g U[        R                  :X  aR  U[        R                  :X  a>  [        R                  SUR                   S35        UR                  [        [        U5        U[        R                  :w  a9  U[        R                  :X  a$  [        R                  SUR                   S35        g g g )Nz	Database z- is unreachable. Failover has been initiated.z is reachable again.)rH   get_running_looprw   	HALF_OPENrp   r   r{   rO   rx   r   r   warning
call_laterr   _half_open_circuitinfo)rP   rs   r   r   loops        rQ   ru   /MultiDBClient._on_circuit_state_change_callback  s     '')))))0)<)<%%g&6&67*D& &9+DNNG,,--Z[ OO02DgN&9+FKK)G$4$4#55IJK ,G&rT   )rA   rK   rC   rL   r3   rB   r?   r=   rO   rJ   rN   r8   r;   r6   rM   rF   rG   )rP   r'   rU   r'   )T).__name__
__module____qualname____firstlineno____doc__r   rR   rZ   ra   rj   rX   r   r   r   r   r   boolr   r   r   floatr   r   r   r   r   r   r   r   r	   r   r   r"   r   strr   r   dictr   rr   ro   r   r   rw   ru   __static_attributes__r   rT   rQ   r'   r'   !   sz   
+*} +*Z
H " Hy 
- 
D 
8 IM/J$/JAE/Jb	)	DQ	m J] JE J&95I 94+ 4M %)$)'+

|U3	#+>%??@
 
 SM	

 "
 e_
,	& tHdN/C  D0}  $L%L29LFMLrT   r'   rs   c                 .    [         R                  U l        g rW   )rw   r   rv   )rs   s    rQ   r   r     s    %%GMrT   c                       \ rS rSr% SrSr\S   \S'   S\4S jr	SS jr
S	 rS
 rS rS\4S jrS\4S jrSS jrSS jrSS jrS rS\\   4S jrSrg)r   i  z?
Pipeline implementation for multiple logical Redis databases.
T_is_async_clientr`   c                     / U l         Xl        g rW   )_command_stack_client)rP   r`   s     rQ   rR   Pipeline.__init__  s     rT   rU   c                    #    U $ 7frW   r   rY   s    rQ   rZ   Pipeline.__aenter__  
        c                    #    U R                  5       I S h  vN   U R                  R                  XU5      I S h  vN   g  N) N7frW   )resetr  rj   rf   s       rQ   rj   Pipeline.__aexit__  s6     jjlll$$X)DDD 	Ds   AA #AAAAc                 >    U R                  5       R                  5       $ rW   )_async_self	__await__rY   s    rQ   r  Pipeline.__await__  s    !++--rT   c                    #    U $ 7frW   r   rY   s    rQ   r  Pipeline._async_self  r  r  c                 ,    [        U R                  5      $ rW   )r   r
  rY   s    rQ   __len__Pipeline.__len__  s    4&&''rT   c                     g)z1Pipeline instances should always evaluate to TrueTr   rY   s    rQ   __bool__Pipeline.__bool__  s    rT   Nc                    #    / U l         g 7frW   )r
  rY   s    rQ   r  Pipeline.reset  s      s   	c                 @   #    U R                  5       I Sh  vN   g N7f)zClose the pipelineN)r  rY   s    rQ   ra   Pipeline.aclose  s     jjlrm   c                 >    U R                   R                  X45        U $ )a:  
Stage a command to be executed when execute() is next called

Returns the current Pipeline object back so commands can be
chained together, such as:

pipe = pipe.set('foo', 'bar').incr('baz').decr('bang')

At some other point, you can then run: pipe.execute(),
which will execute all commands queued in the pipe.
)r
  r   r   s      rQ   pipeline_execute_command!Pipeline.pipeline_execute_command  s     	""D?3rT   c                 &    U R                   " U0 UD6$ )zAdds a command to the stack)r%  rP   r   r   s      rQ   r   Pipeline.execute_command  s    ,,d=f==rT   c                 ~  #    U R                   R                  (       d"  U R                   R                  5       I Sh  vN    U R                   R                  R	                  [        U R                  5      5      I Sh  vN U R                  5       I Sh  vN   $  N] N N	! U R                  5       I Sh  vN    f = f7f)z0Execute all the commands in the current pipelineN)r  rG   rX   rF   execute_pipelinetupler
  r  rY   s    rQ   r   Pipeline.execute  s     ||'',,))+++	66GGd))*  **, , $**,sW   9B=BB=;B <B=B  B=BB=B B=B:3B64B::B=)r  r
  )rP   r   rU   r   rU   N)rU   r   )r   r   r   r   r   r  r   __annotations__r'   rR   rZ   rj   r  r  intr  r  r  r  ra   r%  r   r   r   r   r  r   rT   rQ   r   r     su     '+gdm*} E.( ($ !>
tCy 
rT   r   c                       \ rS rSrSrS\4S jrSS jrSS jrS	 r	\
S\4S
 j5       rS\4S jrS\\-  S\SS4S jrS\4S jrS\\-  S\SS4S jrS r SS\S\\   4S jjrSSS.S\SS4S jjrSrg)r   i  z*
PubSub object for multi database client.
r`   c                 \    Xl         U R                   R                  R                  " S0 UD6  g)zInitialize the PubSub object for a multi-database client.

Args:
    client: MultiDBClient instance to use for pub/sub operations
    **kwargs: Additional keyword arguments to pass to the underlying pubsub implementation
Nr   )r  rF   r   )rP   r`   r   s      rQ   rR   PubSub.__init__  s$     %%,,6v6rT   rU   c                    #    U $ 7frW   r   rY   s    rQ   rZ   PubSub.__aenter__  r  r  Nc                 @   #    U R                  5       I S h  vN   g  N7frW   re   rf   s       rQ   rj   PubSub.__aexit__  rl   rm   c                 h   #    U R                   R                  R                  S5      I S h  vN $  N7f)Nra   r  rF   execute_pubsub_methodrY   s    rQ   ra   PubSub.aclose  s&     \\22HHRRRR   )202c                 V    U R                   R                  R                  R                  $ rW   )r  rF   active_pubsub
subscribedrY   s    rQ   r?  PubSub.subscribed  s    ||,,::EEErT   r   c                 l   #    U R                   R                  R                  " S/UQ76 I S h  vN $  N7f)Nr   r9  rP   r   s     rQ   r   PubSub.execute_command  s7     \\22HH
 $
 
 	
 
   +424r   c                 r   #    U R                   R                  R                  " S/UQ70 UD6I Sh  vN $  N7f)a  
Subscribe to channel patterns. Patterns supplied as keyword arguments
expect a pattern name as the key and a callable as the value. A
pattern's callable will be invoked automatically when a message is
received on that pattern rather than producing a message via
``listen()``.

psubscribeNr9  r(  s      rQ   rF  PubSub.psubscribe  sA      \\22HH

#)
 
 	
 
   .757c                 l   #    U R                   R                  R                  " S/UQ76 I Sh  vN $  N7f)zR
Unsubscribe from the supplied patterns. If empty, unsubscribe from
all patterns.
punsubscribeNr9  rB  s     rQ   rJ  PubSub.punsubscribe(  s9     
 \\22HH
!
 
 	
 
rD  c                 r   #    U R                   R                  R                  " S/UQ70 UD6I Sh  vN $  N7f)a"  
Subscribe to channels. Channels supplied as keyword arguments expect
a channel name as the key and a callable as the value. A channel's
callable will be invoked automatically when a message is received on
that channel rather than producing a message via ``listen()`` or
``get_message()``.
	subscribeNr9  r(  s      rQ   rM  PubSub.subscribe1  sA      \\22HH

"(
 
 	
 
rH  c                 l   #    U R                   R                  R                  " S/UQ76 I Sh  vN $  N7f)zQ
Unsubscribe from the supplied channels. If empty, unsubscribe from
all channels
unsubscribeNr9  rB  s     rQ   rP  PubSub.unsubscribe?  s9     
 \\22HH
 
 
 	
 
rD  ignore_subscribe_messagestimeoutc                 h   #    U R                   R                  R                  SUUS9I Sh  vN $  N7f)z
Get the next message if one is available, otherwise None.

If timeout is specified, the system will wait for `timeout` seconds
before returning. Timeout should be specified as a floating point
number or None to wait indefinitely.
get_message)rR  rS  Nr9  )rP   rR  rS  s      rQ   rU  PubSub.get_messageH  s=      \\22HH&? I 
 
 	
 
r<  g      ?)exception_handlerpoll_timeoutrX  c                f   #    U R                   R                  R                  X!U S9I Sh  vN $  N7f)a`  Process pub/sub messages using registered callbacks.

This is the equivalent of :py:meth:`redis.PubSub.run_in_thread` in
redis-py, but it is a coroutine. To launch it as a separate task, use
``asyncio.create_task``:

    >>> task = asyncio.create_task(pubsub.run())

To shut it down, use asyncio cancellation:

    >>> task.cancel()
    >>> await task
)
sleep_timerW  r   N)r  rF   execute_pubsub_run)rP   rW  rX  s      rQ   run
PubSub.runX  s:     & \\22EE#QU F 
 
 	
 
s   (1/1)r  )rU   r   r.  )Fg        )r   r   r   r   r   r'   rR   rZ   rj   ra   propertyr  r?  r!   r   r    r$   r#   rF  rJ  rM  rP  r   r  rU  r\  r  r   rT   rQ   r   r     s    	7} 	7S FD F F
: 


,
8E
	

 

,
8E
	

 SV
)-
@H
& !	
 	

 

 
rT   r   )<rH   loggingtypingr   r   r   r   r   r   r	   &redis.asyncio.multidb.command_executorr
   redis.asyncio.multidb.configr   r   r   r   redis.asyncio.multidb.databaser   r   r   &redis.asyncio.multidb.failure_detectorr   !redis.asyncio.multidb.healthcheckr   r   redis.asyncio.retryr   redis.backgroundr   redis.backoffr   redis.commandsr   r   redis.multidb.circuitr   r   rw   redis.multidb.exceptionr   r   r   redis.observability.attributesr   redis.typingr    r!   r"   r#   r$   redis.utilsr%   	getLoggerr   r   r'   r   r   r   r   rT   rQ   <module>rp     s      K K K I  N M G L % 0 # F 0 2 
 = P P $			8	$ IL,.? IL ILX& &C'): CLu
 u
rT   