
    ցjy&              
           S r SSKrSSKJr  SSKJr  SSKJr  SSK	J
r
  SS	/0r/ S
Qr " S S5      r\SS j5       rS rS rSSSS\SSSSS.	S jrS rS rSS.S jrg)z,
Thin wrappers around `concurrent.futures`.
    N)contextmanagerlength_hint   )tqdm)TqdmWarningzgithub.com/	casperdcl)
thread_mapprocess_mapinterpreter_mapc                   R    \ rS rSrSrSSKJr  SSKJr	  S r
SS jrS rS	 rS
 rSrg)_InterpreterLock   z3Reentrant lock backed by a cross-interpreter queue.r   )	get_ident)	monotonicc                 P    SSK Jn  Xl        U" 5       U l        S U l        SU l        g )Nr   )RLock)	threadingr   _queue_lock_owner_depth)selfqueuer   s      Q/mnt/workspace/venv-train/lib/python3.13/site-packages/tqdm/contrib/concurrent.py__init___InterpreterLock.__init__   s!    #W
    c                    SSK Jn  U R                  5       nUS:X  a  U R                  R	                  U5      nOU R                  R	                  X5      nU(       d  gU R
                  (       a  U =R
                  S-  sl        g U(       d  U R                  R                  5         OZUS:X  a  U R                  R                  5         O9[        SX R                  5       U-
  -
  5      nU R                  R                  US9   U R                  5       U l        SU l        g! U a    U R                  R                  5          gf = f)Nr   )EmptyF   T)timeout)r   r    _timer   acquirer   r   
get_nowaitgetmaxreleaser   r   )r   blockingr#   r    startacquired	remainings          r   r%   _InterpreterLock.acquire   s    

b=zz))(3Hzz))(<H;;KK1K
	&&(B!7jjlU.B#CD		2 nn&  	JJ 	s   !D & D 8D  E ?E c                    U R                   U R                  5       :w  a  [        S5      eU =R                  S-  sl        U R                  (       d"  S U l         U R                  R                  S 5        U R                  R                  5         g )Nzcannot release un-acquired lockr"   )r   r   RuntimeErrorr   r   putr   r)   r   s    r   r)   _InterpreterLock.release6   s]    ;;$..**@AAq{{DKKKOOD!

r   c                 &    U R                  5         U $ N)r%   r2   s    r   	__enter___InterpreterLock.__enter__?   s    r   c                 $    U R                  5         g r5   )r)   )r   excs     r   __exit___InterpreterLock.__exit__C   s    r   )r   r   r   r   N)Tr!   )__name__
__module____qualname____firstlineno____doc__r   r   timer   r$   r   r%   r)   r6   r:   __static_attributes__ r   r   r   r      s$    =#'6r   r    c              #      #    [        U SS5      nUc  U=(       d    U R                  5       n[        X!U5      nU R                  U5        Uv   Uc  U ?gU R                  U5        g7f)z>get (create if necessary) and then restore `tqdm_class`'s lockr   N)getattrget_lockset_lockr   )
tqdm_class	lock_namelockold_locks       r   ensure_lockrM   G   sf      z7D1H|0:..04D)D
JH%s   A#A%c           	          S[         R                  < SU R                  < SU R                  R	                  S5      < SU< S3	n[
        U44$ )zGReturn an initializer which bootstraps the parent import path and lock.zimport sys
sys.path[:] = z
from concurrent import interpreters
from importlib import import_module
from tqdm.contrib.concurrent import _InterpreterLock
tqdm_class = import_module(z)
for name in .z:
    tqdm_class = getattr(tqdm_class, name)
tqdm_class.monitor_interval = 0
tqdm_class.set_lock(_InterpreterLock(interpreters.Queue(z))))syspathr=   r>   splitexec)rI   lock_queue_idcodes      r   _get_interpreter_initrV   V   sk    	 %& '1&;&;%> ?!..44S9< =C DQBSSV		X 	 $=r   c                 .   ^ [        U4S jU  5       5      $ )z min(map(length_hint, iterables))c              3   P   >#    U  H  n[        US 5      =mS:  d  M  Tv   M     g7f)r!   r   Nr   ).0itns     r   	<genexpr>_min_map_len.<locals>.<genexpr>h   s&     H9Rk"b.A)Aa(Gqq9s   &	&)min)	iterablesr[   s    @r   _min_map_lenr`   f   s    H9HHHr   r"   g        )	max_workersr#   	chunksizerJ   rI   	smoothingr   _initializer	_initargsc       	         R  ^^ UR                  5       nSU;  a  [        U5      US'   0 nSU;   a  UR                  S5      US'   0 nS H  nUU;   d  M  UR                  U5      UU'   M!     SnUS   (       aC  SU;  a=   SSKJn  U=(       d    [        S	U" 5       =(       d    S
S-   5      nUS   U:  a  UUS'   Sn[        XeUS9 nU	c  UR                  n	U4n
U " SX)U
S.UD6 nU" SSU0UD6 mUb  STl
        UR                  mUU4S jnUUl        [        UR                  " U/UQ7X4S.UD65      sSSS5        sSSS5        sSSS5        $ ! [
         a	    SSKJn   Nf = f! , (       d  f       O= f SSS5        O! , (       d  f       O= fSSS5        g! , (       d  f       g= f)z
Implementation of `thread_map`, `process_map` and `interpreter_map`.

Parameters
----------
max_workers  : int
timeout  : int
buffersize  : int
    Requires Python>=3.14.
thread_name_prefix  : str
max_tasks_per_child  : int
mp_context  : str
total
buffersize)thread_name_prefixmax_tasks_per_child
mp_contextNminitersr   )process_cpu_count)	cpu_count    r"      T)rJ   rK   )ra   initializerinitargsrc   c                  B   > T" U 0 UD6nUR                  U4S j5        U$ )Nc                 $   > TR                  5       $ r5   )update)_pbars    r   <lambda>4_executor_map.<locals>.patchsubmit.<locals>.<lambda>   s    DKKMr   )add_done_callback)argskwargsfut	orisubmitrw   s      r   patchsubmit"_executor_map.<locals>.patchsubmit   s&    #T4V4C))*ABJr   )r#   rb   rC   )copyr`   poposrm   ImportErrorrn   r^   rM   rH   dynamic_miniterssubmitlistmap)PoolExecutorfnra   r#   rb   rJ   rI   rc   r   rd   re   r_   tqdm_kwargsr|   
map_kwargspool_kwargskr   rn   	rough_maxlkexr   r~   rw   s                          @@r   _executor_mapr   k   s   $ Ff&y1wJv#)::l#;
< KH;#ZZ]KN I g:V3	%9  B3rIK,<1+A#B	'?Y&!*F:#	Z5	AR%..LI )kV_ )'),.:i:6:d#/,0D)II	 (	BFFX"X,3XLVX Y ;:) ) 
B	A  	%$	% ;::) ) ) 
B	A	AsU   =E F!E>,AE#1	E>:	FE E #
E1-E>5	F>
F	F
F&c                 ,    SSK Jn  [        X0/UQ70 UD6$ )a`  
Equivalent of `list(map(fn, *iterables))`
driven by `concurrent.futures.ThreadPoolExecutor`.

Parameters
----------
max_workers  : int, optional
    Maximum number of workers to spawn; passed to `concurrent.futures.ThreadPoolExecutor`.
thread_name_prefix  : str, optional
    Passed to `concurrent.futures.ThreadPoolExecutor` [default: ''].
timeout  : int or float, optional
    Seconds to wait before raising `TimeoutError` if `__next__` is called and the
    result isn't available. [default: None].
buffersize  : int, optional
    Requires Python>=3.14 [default: None].
tqdm_class  : optional
    `tqdm` class to use for bars [default: tqdm.auto.tqdm].
smoothing  : float, optional
    Passed to `tqdm_class`; the [default: 0] is average (due to erratic update frequency).
lock_name  : str, optional
    Member of `tqdm_class.get_lock()` to use [default: ''].
r   )ThreadPoolExecutor)concurrent.futuresr   r   )r   r_   r   r   s       r   r
   r
      s    . 6+K)K{KKr   c                     SSK Jn  SSKJn  UR	                  5       nUR                  S5        UR                  S[        5      n[        XeR                  5      u  px[        X@/UQ7[        U5      XxS.UD6$ )aB  
Equivalent of `list(map(fn, *iterables))`
driven by `concurrent.futures.InterpreterPoolExecutor` (Python 3.14+).

Parameters
----------
Same as `thread_map`.

Notes
-----
`fn`, its arguments, and its return values must be pickleable.
Worker progress bars using the same `tqdm_class` share a cross-interpreter write lock.
r   )interpreters)InterpreterPoolExecutorNrI   )r   rd   re   )
concurrentr   r   r   create_queuer1   r'   	tqdm_autorV   idr   r   )	r   r_   r   r   r   
lock_queuerI   rq   rr   s	            r   r   r      sz     (:**,JNN4y9J1*mmLKE&/E7G
7S E8CE Er   mp_lock)rJ   c                    SSK Jn  U(       a,  SU;  a&  [        U5      nUS:  a  SSKJn  U" SU-  [
        SS9  [        X@/UQ7S	U0UD6$ )
a  
Equivalent of `list(map(fn, *iterables))`
driven by `concurrent.futures.ProcessPoolExecutor`.

Parameters
----------
max_workers  : int, optional
    Maximum number of workers to spawn; passed to `concurrent.futures.ProcessPoolExecutor`.
timeout  : int or float, optional
    Seconds to wait before raising `TimeoutError` if `__next__` is called and the
    result isn't available. [default: None].
chunksize  : int, optional
    Approximate size of chunks sent to worker processes; passed to
    `concurrent.futures.ProcessPoolExecutor.map`. [default: 1].
buffersize  : int, optional
    Requires Python>=3.14 [default: None].
max_tasks_per_child  : int, optional
    Maximum number of tasks a worker process can complete before being replaced
    with a new process; passed to `concurrent.futures.ProcessPoolExecutor`.
mp_context  : multiprocessing.BaseContext, optional
    Multiprocessing context to use, e.g. `multiprocessing.get_context('fork')`.
lock_name  : str, optional
    Member of `tqdm_class.get_lock()` to use [default: mp_lock].
tqdm_class  : optional
    `tqdm` class to use for bars [default: tqdm.auto.tqdm].
smoothing  : float, optional
    Passed to `tqdm_class`; the [default: 0] is average (due to erratic update frequency).
r   )ProcessPoolExecutorrb   i  )warnzIterable length %d > 1000 but `chunksize` is not set. This may seriously degrade multiprocess performance. Set `chunksize=1` or more.r   )
stacklevelrJ   )r   r   r`   warningsr   r   r   )r   rJ   r_   r   r   shortest_iterable_lenr   s          r   r   r      sd    : 7[3 !-Y 7 4'% /1FG , ,a9a	aU`aar   )rD   N)r@   rP   
contextlibr   operatorr   autor   r   stdr   
__author____all__r   rM   rV   r`   r   r
   r   r   rC   r   r   <module>r      s     %   $ k]+

:5 5p & & I /3DAY[Ct$RV9YxL6E2 +4 (br   