
    ng\jZ                        S SK r S SKrS SKrS SKJrJrJrJr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  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'J(r(J)r)  S SK*J+r+  S SK,J-r-  S SK.J/r/  S SK0J1r1J2r2J3r3  S SK4J5r5  \Rl                  " \75      r8\5 " S S\\5      5       r9S\4S jr: " S S\\5      r; " S S5      r<g)    N)AnyCallableListLiteralOptional)HealthCheckHealthCheckPolicy)BackgroundScheduler)	NoBackoff)PubSubWorkerThread)CoreCommandsRedisModuleCommands)MaintNotificationsConfig)CircuitBreaker)State)DefaultCommandExecutor)DEFAULT_GRACE_PERIODDatabaseConfigInitialHealthCheckMultiDbConfig)Database	DatabasesSyncDatabase)InitialHealthCheckFailedErrorNoValidDatabaseExceptionUnhealthyDatabaseException)FailureDetector)GeoFailoverReason)Retry)ChannelTPubSubHandlerSubscription)experimentalc                   0   \ rS rSrSrS\4S j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
4   4S jrS rS	\S\4S jrS\\\4   4S jr S  r!S!\"S"\#S#\#4S$ jr$S% r%S&r&g
)(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        4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        [=        5       U l        [@        RB                  " 5       U l"        X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_databasesr0   _auto_fallback_intervalr/   _event_dispatcherr+   _command_retryupdate_supported_errorsConnectionRefusedErrorr   r-   r.   command_executorinitializedr
   _bg_scheduler	threadingLock_hc_lock_config)selfr'   s     O/home/edenadmin/noVNC/venv/lib/python3.13/site-packages/redis/multidb/client.py__init__MultiDBClient.__init__+   s    **, '' ((*%% 	
 '-&B&B#&&,,. 	!
 ++ ,,.)) 	 ''/ ,,.)) 	
 	--doo>'-'D'D$!'!8!8$22335K4MN 6"55oo--"55$66!00!33#'#?#?	!
 !02!(    c                 F     U R                  5         g ! [         a     g f = fN)close	ExceptionrK   s    rL   __del__MultiDBClient.__del__U   s$    	JJL 	 		    
  c                    U R                   R                  U R                  5        U R                   R                  U R                  U R
                  5        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)zD
Perform initialization of databases to define their initial state.
FTz4Initial connection failed - no active database foundN)rF   run_coro_sync_perform_initial_health_checkrun_recurring_coror6   _check_databases_healthr1   circuiton_state_changed!_on_circuit_state_change_callbackstateCBStateCLOSEDrD   _active_databaser   rE   )rK   is_active_db_founddatabaseweights       rL   
initializeMultiDBClient.initialize^   s     	(()K)KL 	--''((	

 # $H--d.T.TU %%7@R@R :B%%6%)" !0 "*F   rO   returnc                     U R                   $ )z5
Returns a sorted (by weight) list of all databases.
)r1   rT   s    rL   get_databasesMultiDBClient.get_databases   s     rO   re   Nc                    SnU R                    H  u  p4X1:X  d  M  Sn  O   U(       d  [        S5      eU R                  R                  U R                  U5        UR
                  R                  [        R                  :X  aB  U R                   R                  S5      S   u  pTU[        R                  4U R                  l        g[        S5      e)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)r1   
ValueErrorrF   rY   _check_db_healthr]   r`   ra   rb   	get_top_nr   MANUALrD   active_databaser   )rK   re   existsexisting_db_highest_weighted_dbs         rL   set_active_database!MultiDBClient.set_active_database   s     "ooNK& .
 NOO(()>)>I!!W^^3%)__%>%>q%A!%D"!((5D!!1 &?
 	
rO   skip_initial_health_checkc                    [        S[        5       S9UR                  S'   SUR                  ;  a  [        SS9UR                  S'   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                  R                  U R                   U5        U R$                  R'                  S
5      S   u  pgU R$                  R)                  XUR                  5        U R+                  XV5        g! ["         a    U(       d  e  Nkf = 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.
r   )retriesbackoffretrymaint_notifications_configF)enabled)connection_poolN)clientr]   rf   health_check_urlro    )r   r   client_kwargsr   from_urlrJ   client_class	from_pool	set_retryr]   default_circuit_breakerr   rf   r   rF   rY   rq   r   r1   rr   add_change_active_database)rK   r'   r{   r   r]   re   rx   highest_weights           rL   add_databaseMultiDBClient.add_database   s    ).a(MW% (v/C/CC(7   !=> ??\\..77#)#7#7F &&uQ	'LM\\..88 & 0 0 9 F \\..F1E1EFF ~~% **, 	 ==#44	
	,,T-B-BHM
 /3oo.G.G.J1.M+Hoo6$$XC * 	, -	s   -&G* *G?>G?new_databasehighest_weight_databasec                     UR                   UR                   :  aK  UR                  R                  [        R                  :X  a"  U[
        R                  4U R                  l        g g g rQ   )	rf   r]   r`   ra   rb   r   	AUTOMATICrD   rt   )rK   r   r   s      rL   r   %MultiDBClient._change_active_database   sZ     "9"@"@@$$**gnn< !++5D!!1 = ArO   c                    U R                   R                  U5      nU R                   R                  S5      S   u  p4XB::  aK  UR                  R                  [
        R                  :X  a"  U[        R                  4U R                  l
        ggg)z,
Removes a database from the database list.
ro   r   N)r1   removerr   r]   r`   ra   rb   r   rs   rD   rt   )rK   re   rf   rx   r   s        rL   remove_databaseMultiDBClient.remove_database   s~     ''1.2oo.G.G.J1.M+ $#++11W^^C $!((5D!!1 D %rO   rf   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        g)z,
Updates a database from the database list.
NTrn   ro   r   )r1   rp   rr   update_weightrf   r   )rK   re   rf   ru   rv   rw   rx   r   s           rL   update_database_weight$MultiDBClient.update_database_weight   sz     "ooNK& .
 NOO.2oo.G.G.J1.M+%%h7 $$XCrO   failure_detectorc                 :    U R                   R                  U5        g)z.
Adds a new failure detector to the database.
N)r;   append)rK   r   s     rL   add_failure_detector"MultiDBClient.add_failure_detector  s     	&&'78rO   healthcheckc                     U R                      U R                  R                  U5        SSS5        g! , (       d  f       g= f)z*
Adds a new health check to the database.
N)rI   r4   r   )rK   r   s     rL   add_health_checkMultiDBClient.add_health_check  s)     ]]&&{3 ]]s	   2
A c                 |    U R                   (       d  U R                  5         U R                  R                  " U0 UD6$ )z2
Executes a single command and return its result.
)rE   rg   rD   execute_commandrK   argsoptionss      rL   r   MultiDBClient.execute_command  s3     OO$$44dFgFFrO   c                     [        U 5      $ )z*
Enters into pipeline mode of the client.
)PipelinerT   s    rL   pipelineMultiDBClient.pipeline#  s     ~rO   funcr   c                     U R                   (       d  U R                  5         U R                  R                  " U/UQUQ76 $ )z#
Executes callable as transaction.
)rE   rg   rD   execute_transaction)rK   r   watchesr   s       rL   transactionMultiDBClient.transaction)  s8     OO$$88RR'RRrO   c                 \    U R                   (       d  U R                  5         [        U 40 UD6$ )z
Return a Publish/Subscribe object. With this object, you can
subscribe to channels and listen for messages that get published to
them.
)rE   rg   PubSub)rK   kwargss     rL   pubsubMultiDBClient.pubsub2  s'     OOd%f%%rO   c                   #    U R                      [        U R                  5      nSSS5        U R                  R	                  W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$ ! , (       d  f       N= f N7f)z?
Runs health checks on the given database until first failure.
N)
rI   listr4   r9   executer]   r`   ra   OPENrb   )rK   re   r2   
is_healthys       rL   rq   MultiDBClient._check_db_health=  s      ]] !4!45M   44<<]HUU
%%5)0  &H,,22gnnD%,^^H" ] Vs(   DC1'DDB$D1
C?;D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)	_hc_tasksr1   asynciocreate_taskrq   r   gatherzipitems
isinstancer   re   ra   r   r]   r`   loggerdebugoriginal_exception)	rK   
task_to_dbre   rw   taskresultsresult
db_resultsunhealthy_dbs	            rL   r\   %MultiDBClient._check_databases_healthP  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: )r\   rJ   initial_health_check_policyr   ALL_AVAILABLEvaluesMAJORITY_AVAILABLEsumlenONE_AVAILABLEr   )rK   r   r   s      rL   rZ   +MultiDBClient._perform_initial_health_checkr  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r]   	old_state	new_statec                    U[         R                  :X  a1  U R                  R                  U R                  UR
                  5        g U[         R                  :X  a\  U[         R                  :X  aH  [        R                  SUR
                   S35        U R                  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.)ra   	HALF_OPENrF   run_coro_fire_and_forgetrq   re   rb   r   r   warningrun_oncer   _half_open_circuitinfo)rK   r]   r   r   s       rL   r_   /MultiDBClient._on_circuit_state_change_callback  s     )))77%%w'7'7 &9+DNNG,,--Z[ ''$&8' &9+FKK)G$4$4#55IJK ,G&rO   c                 n   U R                   (       aJ   U R                   R                  U R                  R                  5        U R                   R                  5         U R                  R                  (       a/  U R                  R                  R                  R                  5         gg! [         a     Nqf = f)z*
Closes the client and all its resources.
N)	rF   rY   r9   rR   rS   stoprD   rt   r   rT   s    rL   rR   MultiDBClient.close  s     ""001J1J1P1PQ ##%  00!!1188>>@ 1	  s   /B' '
B43B4)r?   rF   rA   rJ   r1   r@   r=   r;   rI   r   r6   r9   r4   rD   rE   )T)'__name__
__module____qualname____firstlineno____doc__r   rM   rU   rg   r   rk   r   ry   r   boolr   r   r   r   floatr   r   r   r   r   r   r   r   r   r   rq   dictr\   rZ   r   ra   r_   rR   __static_attributes__r   rO   rL   r%   r%   $   s,   
(} (T$ Ly 
L 
T 
: IM6D$6DAE6Dp
(
CO
  D| DU D&9_ 94K 4GS*t); < S	&|  & tHdN/C  D0L%L29LFML*ArO   r%   r]   c                 .    [         R                  U l        g rQ   )ra   r   r`   )r]   s    rL   r   r     s    %%GMrO   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\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.
F_is_async_clientr   c                     / U l         Xl        g rQ   )_command_stack_client)rK   r   s     rL   rM   Pipeline.__init__  s     rO   ri   c                     U $ rQ   r   rT   s    rL   	__enter__Pipeline.__enter__      rO   c                 $    U R                  5         g rQ   reset)rK   exc_type	exc_value	tracebacks       rL   __exit__Pipeline.__exit__      

rO   c                 F     U R                  5         g ! [         a     g f = frQ   r  rS   rT   s    rL   rU   Pipeline.__del__  s"    	JJL 		rW   c                 ,    [        U R                  5      $ rQ   )r   r   rT   s    rL   __len__Pipeline.__len__  s    4&&''rO   c                     g)z1Pipeline instances should always evaluate to TrueTr   rT   s    rL   __bool__Pipeline.__bool__  s    rO   Nc                     / U l         g rQ   )r   rT   s    rL   r  Pipeline.reset  s
     rO   c                 $    U R                  5         g)zClose the pipelineNr  rT   s    rL   rR   Pipeline.close  s    

rO   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      rL   pipeline_execute_command!Pipeline.pipeline_execute_command  s     	""D?3rO   c                 &    U R                   " U0 UD6$ )zAdds a command to the stack)r  rK   r   r   s      rL   r   Pipeline.execute_command  s    ,,d=f==rO   c                 (   U R                   R                  (       d  U R                   R                  5          U R                   R                  R	                  [        U R                  5      5      U R                  5         $ ! U R                  5         f = f)z0Execute all the commands in the current pipeline)r   rE   rg   rD   execute_pipelinetupler   r  rT   s    rL   r   Pipeline.execute  s_    ||''LL##%	<<00AAd))* JJLDJJLs   7A? ?B)r   r   )ri   r   ri   N)r   r   r   r   r   r   r   __annotations__r%   rM   r   r  rU   intr  r   r  r  rR   r  r   r   r   r   r   r   rO   rL   r   r     so     (-gen,} ( ($ !>
c 
rO   r   c                   2   \ rS rSrSrS\4S jrS S jrS!S jrS!S	 jr	S!S
 jr
\S\4S j5       rS rS\\-  S\SS4S jrS rS\\-  S\SS4S jrS rS\\-  S\SS4S jrS r S"S\S\4S jjr S"S\S\4S jjr    S#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   rD   r   )rK   r   r   s      rL   rM   PubSub.__init__   s$     %%,,6v6rO   ri   c                     U $ rQ   r   rT   s    rL   r   PubSub.__enter__  r   rO   Nc                 F     U R                  5         g ! [         a     g f = frQ   r  rT   s    rL   rU   PubSub.__del__  s$    	 JJL 		rW   c                 L    U R                   R                  R                  S5      $ )Nr  r   rD   execute_pubsub_methodrT   s    rL   r  PubSub.reset  s    ||,,BB7KKrO   c                 $    U R                  5         g rQ   r  rT   s    rL   rR   PubSub.close  r	  rO   c                 V    U R                   R                  R                  R                  $ rQ   )r   rD   active_pubsub
subscribedrT   s    rL   r3  PubSub.subscribed  s    ||,,::EEErO   c                 P    U R                   R                  R                  " S/UQ76 $ )Nr   r,  rK   r   s     rL   r   PubSub.execute_command!  s*    ||,,BB
 $
 	
rO   r   r   c                 V    U R                   R                  R                  " S/UQ70 UD6$ )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()``.

psubscriber,  r  s      rL   r9  PubSub.psubscribe&  4     ||,,BB

#)
 	
rO   c                 P    U R                   R                  R                  " S/UQ76 $ )zR
Unsubscribe from the supplied patterns. If empty, unsubscribe from
all patterns.
punsubscriber,  r6  s     rL   r=  PubSub.punsubscribe4  ,    
 ||,,BB
!
 	
rO   c                 V    U R                   R                  R                  " S/UQ70 UD6$ )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()``.
	subscriber,  r  s      rL   rA  PubSub.subscribe=  s4     ||,,BB

"(
 	
rO   c                 P    U R                   R                  R                  " S/UQ76 $ )zQ
Unsubscribe from the supplied channels. If empty, unsubscribe from
all channels
unsubscriber,  r6  s     rL   rD  PubSub.unsubscribeK  s%    
 ||,,BB=XSWXXrO   c                 V    U R                   R                  R                  " S/UQ70 UD6$ )aJ  
Subscribes the client to the specified shard 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_sharded_message()``.

ssubscriber,  r  s      rL   rG  PubSub.ssubscribeR  r;  rO   c                 P    U R                   R                  R                  " S/UQ76 $ )z]
Unsubscribe from the supplied shard_channels. If empty, unsubscribe from
all shard_channels
sunsubscriber,  r6  s     rL   rJ  PubSub.sunsubscribe`  r?  rO   ignore_subscribe_messagestimeoutc                 L    U R                   R                  R                  SUUS9$ )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rL  rM  r,  rK   rL  rM  s      rL   rO  PubSub.get_messagei  s0     ||,,BB&? C 
 	
rO   c                 L    U R                   R                  R                  SUUS9$ )z
Get the next message if one is available in a sharded channel, 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_sharded_messagerP  r,  rQ  s      rL   rT  PubSub.get_sharded_messagey  s0     ||,,BB!&? C 
 	
rO   
sleep_timedaemonexception_handlersharded_pubsubr   c                 P    U R                   R                  R                  UUUU US9$ )N)rW  rX  r   rY  )r   rD   execute_pubsub_run)rK   rV  rW  rX  rY  s        rL   run_in_threadPubSub.run_in_thread  s6     ||,,??/) @ 
 	
rO   )r   )ri   r   r!  )F        )r^  FNF)r   r   r   r   r   r%   rM   r   rU   r  rR   propertyr   r3  r   r    r"   r!   r9  r=  rA  rD  rG  rJ  r   rO  rT  r   r   r\  r   r   rO   rL   r   r     sF   	7} 	7L FD F F


,
8E
	


,
8E
	
Y
,
8E
	

 IL
)-
@E
" IL
)-
@E
$  04$

 
 $H-	

 
 

 
rO   r   )=r   loggingrG   typingr   r   r   r   r   !redis.asyncio.multidb.healthcheckr   r	   redis.backgroundr
   redis.backoffr   redis.clientr   redis.commandsr   r   redis.maint_notificationsr   redis.multidb.circuitr   r   ra   redis.multidb.command_executorr   redis.multidb.configr   r   r   r   redis.multidb.databaser   r   r   redis.multidb.exceptionr   r   r   redis.multidb.failure_detectorr   redis.observability.attributesr   redis.retryr   redis.typingr    r!   r"   redis.utilsr#   	getLoggerr   r   r%   r   r   r   r   rO   rL   <module>rs     s       9 9 L 0 # + < > 0 2 A  E D 
 ; <  > > $			8	$ JA' JA JAZ& &B"L BJ[
 [
rO   