Ë
    æÿæi¡3  ã                   óp  — d dl mZmZmZ d dlZd dlZd dlZd dlZd dl	Z	d dl
mZ ddlmZmZ ddlmZmZmZ 	 d dlZd dlZe�)e�'d dlmZmZmZmZmZ d dlmZ d d	lmZ d d
l m!Z! 	 d dl m"Z# d„ Z% G d„ d«      Z&d„ Z'd„ Z( G d„ d«      Z)d„ Z* G d„ dee«      Z+y# e$ r dZdZY Œcw xY w# e$ r	 d dl$m"Z# Y ŒIw xY w)é    )Úabsolute_importÚdivisionÚprint_functionN)Úuuid4é   )Ú*_retrieve_traceback_capturing_wrapped_callÚ_TracebackCapturingWrapper)ÚAutoBatchingMixinÚParallelBackendBaseÚparallel_config)ÚClientÚas_completedÚ
get_clientÚrejoinÚsecede)Úsizeof)Úfuncname)Úthread_state)ÚTimeoutErrorc                 óN   — 	 t        j                  | «       y# t        $ r Y yw xY w)NTF)ÚweakrefÚrefÚ	TypeError)Úobjs    úa/Volumes/fast/ai/experiments/voice-extract-mac/.venv/lib/python3.12/site-packages/joblib/_dask.pyÚis_weakrefabler   +   s(   € ðÜ�‰�CÔØøÜò Ùðús   ‚ ˜	$£$c                   ó.   — e Zd ZdZd„ Zd„ Zd„ Zd„ Zd„ Zy)Ú_WeakKeyDictionarya¬  A variant of weakref.WeakKeyDictionary for unhashable objects.

    This datastructure is used to store futures for broadcasted data objects
    such as large numpy arrays or pandas dataframes that are not hashable and
    therefore cannot be used as keys of traditional python dicts.

    Furthermore using a dict with id(array) as key is not safe because the
    Python is likely to reuse id of recently collected arrays.
    c                 ó   — i | _         y ©N©Ú_data©Úselfs    r   Ú__init__z_WeakKeyDictionary.__init__>   s	   € Øˆ�
ó    c                 ód   — | j                   t        |«         \  }} |«       |urt        |«      ‚|S r    )r"   ÚidÚKeyError)r$   r   r   Úvals       r   Ú__getitem__z_WeakKeyDictionary.__getitem__A   s1   € Ø—:‘:œb ›gÑ&‰ˆˆSÙ‹5˜Ñä˜3“-ÐØˆ
r&   c                 óæ   ‡ ‡— t        |«      Š	 ‰ j                  ‰   \  }} |«       |urt        |«      ‚	 ||f‰ j                  ‰<   y # t        $ r ˆˆ fd„}t        j                  ||«      }Y Œ9w xY w)Nc                 ó    •— ‰j                   ‰= y r    r!   )Ú_Úkeyr$   s    €€r   Ú
on_destroyz2_WeakKeyDictionary.__setitem__.<locals>.on_destroyS   s   ø€ Ø—J‘J˜s‘Or&   )r(   r"   r)   r   r   )r$   r   Úvaluer   r.   r0   r/   s   `     @r   Ú__setitem__z_WeakKeyDictionary.__setitem__H   sv   ù€ Ü�‹gˆð	/Ø—Z‘Z ‘_‰FˆC�Ù‹u˜CÑä˜s“mÐ#ð  ð ˜u˜*ˆ�
‰
�3Šøô ò 	/õ$ô —+‘+˜c :Ó.ŠCð	/ús   �&A Á%A0Á/A0c                 ó,   — t        | j                  «      S r    )Úlenr"   r#   s    r   Ú__len__z_WeakKeyDictionary.__len__Y   s   € Ü�4—:‘:‹Ðr&   c                 ó8   — | j                   j                  «        y r    )r"   Úclearr#   s    r   r7   z_WeakKeyDictionary.clear\   s   € Ø�
‰
×ÑÕr&   N)	Ú__name__Ú
__module__Ú__qualname__Ú__doc__r%   r+   r2   r5   r7   © r&   r   r   r   3   s    „ ñòòò%ò"ór&   r   c                 ó|   — 	 t        | t        «      r| d   d   } t        | «      S # t        $ r Y t        | «      S w xY w)Nr   )Ú
isinstanceÚlistÚ	Exceptionr   )Úxs    r   Ú	_funcnamerB   `   sH   € ðÜ�aœÔØ�!‘�Q‘ˆAô �A‹;Ðøô ò ØÜ�A‹;Ððús   ‚% ¥	;º;c                 ó’   — | D ���ch c]  \  }}}|’Œ
 }}}}t        |«      dk(  rd}nd}t        | «      |t        | «      fS c c}}}w )z8Summarize of list of (func, args, kwargs) function callsr   FT)r4   rB   )ÚtasksÚfuncÚargsÚkwargsÚunique_funcsÚmixeds         r   Ú_make_tasks_summaryrJ   i   sO   € á38Õ9±5Ñ/˜T 4¨’D°5€LÒ9ä
ˆ<Ó˜AÒØ‰àˆÜˆu‹:�uœi¨Ó.Ð.Ð.ùô :s   ‡Ac                   ó$   — e Zd ZdZd„ Zdd„Zd„ Zy)ÚBatchz6dask-compatible wrapper that executes a batch of tasksc                 ó@   — t        |«      \  | _        | _        | _        y r    )rJ   Ú
_num_tasksÚ_mixedrB   )r$   rD   s     r   r%   zBatch.__init__w   s   € ô 8KÈ5Ó7QÑ4ˆŒ˜œ d¥nr&   Nc           	      ó’   — g }t        d¬«      5  |D ]  \  }}}|j                   ||i |¤Ž«       Œ |cd d d «       S # 1 sw Y   y xY w)NÚdask)Úbackend)r   Úappend)r$   rD   ÚresultsrE   rF   rG   s         r   Ú__call__zBatch.__call__|   sE   € ØˆÜ VÖ,Û&+Ñ"��d˜FØ—‘™t TÐ4¨VÑ4Õ5ð ',à÷ -×,Ò,ús	   �$=½Ac                 ób   — d| j                   › d| j                  › d�}| j                  rd|z   }|S )NÚ	batch_of_r.   Ú_callsÚmixed_)rB   rN   rO   )r$   Údescrs     r   Ú__repr__zBatch.__repr__ƒ   s6   € Ø˜DŸN™NÐ+¨1¨T¯_©_Ð,=¸VÐDˆØ�;Š;Ø˜uÑ$ˆEØˆr&   r    )r8   r9   r:   r;   r%   rU   r[   r<   r&   r   rL   rL   t   s   „ Ù@òRó
ór&   rL   c                   ó   — y r    r<   r<   r&   r   Ú_joblib_probe_taskr]   Š   s   € àr&   c                   ó¦   ‡ — e Zd ZdZdZdZdZ	 	 	 	 	 dˆ fd„	Zd„ Zd„ Z	d„ Z
dd	„Zd
„ Zd„ Zd„ Zd„ Zdd„Zd„ Zdd„Zej(                  d„ «       Zˆ xZS )ÚDaskDistributedBackendgš™™™™™É?g      ð?Téÿÿÿÿc                 ó¶  •— t         ‰| �  «        t        €d}t        |«      ‚|€|rt	        ||d¬«      }n	 t        «       }|| _        |�7t        |t        t        f«      s!t        dt        |«      j                  z  «      ‚|�jt        |«      dkD  r\t        |«      | _        | j                  j                  |d¬«      }	t!        ||	«      D �
�ci c]  \  }
}t#        |
«      |“Œ c}}
| _        ng | _        i | _        || _        || _        t+        g |j,                  dd¬	«      | _        i | _        i | _        y # t        $ r}d}t        |«      |‚d }~ww xY wc c}}
w )
Nz{You are trying to use 'dask' as a joblib parallel backend but dask is not installed. Please install dask to fix this error.F)ÚloopÚset_as_defaultz¢To use Joblib with Dask first create a Dask Client

    from dask.distributed import Client
    client = Client()
or
    client = Client('scheduler-address:8786')z&scatter must be a list/tuple, got `%s`r   T)Ú	broadcast)rb   Úwith_resultsÚraise_errors)Úsuperr%   ÚdistributedÚ
ValueErrorr   r   Úclientr>   r?   Útupler   Útyper8   r4   Ú_scatterÚscatterÚzipr(   Údata_futuresÚwait_for_workers_timeoutÚsubmit_kwargsr   rb   Úwaiting_futuresÚ_resultsÚ
_callbacks)r$   Úscheduler_hostrn   rj   rb   rq   rr   ÚmsgÚeÚ	scatteredrA   ÚfÚ	__class__s               €r   r%   zDaskDistributedBackend.__init__•   sf  ø€ ô 	‰ÑÔäÐð%ð ô
 ˜S“/Ð!àˆ>ÙÜ °TÈ%ÔP‘ð1Ü'›\�Fð ˆŒàÐ¤z°'¼DÄ%¸=Ô'IÜØ8¼4À»=×;QÑ;QÑQóð ð Ð¤3 w£<°!Ò#3ä  ›MˆDŒMØŸ™×+Ñ+¨G¸tÐ+ÓDˆIÜ69¸'À9Ô6MÔ NÑ6M©d¨a°¤ A£¨¡Ð6MÒ NˆDÕàˆDŒMØ "ˆDÔØ(@ˆÔ%Ø*ˆÔÜ+Ø�V—[‘[¨tÀ%ô 
ˆÔð ˆŒØˆ�øôA "ò 	1ðHð ô % S›/¨qÐ0ûð	1üó, !Os   ·
D6 ÃEÄ6	EÄ?EÅEc              ƒ   ó¢  K  — | j                   r�| j                  2 3 d {  –—† \  }}| j                  j                  |«      }| j                  j                  |«      }|j
                  dk(  r|\  }}}|j                  |«       Œi|j                  |«        ||«       Œƒy y 7 Œ€6 t        j                  d«      ƒ d {  –—†7   | j                   rŒ¿Œ1­w)NÚerrorç{®Gáz„?)
Ú	_continuers   rt   Úpopru   ÚstatusÚset_exceptionÚ
set_resultÚasyncioÚsleep)r$   ÚfutureÚresultÚ	cf_futureÚcallbackÚtypÚexcÚtbs           r   Ú_collectzDaskDistributedBackend._collectÐ   s°   è ø€ Ø�nŠnØ(,×(<Ò(<÷ %‘n�f˜fØ ŸM™M×-Ñ-¨fÓ5�	ØŸ?™?×.Ñ.¨vÓ6�Ø—=‘= GÒ+Ø#)‘L�C˜˜bØ×+Ñ+¨CÕ0à×(Ñ(¨Ô0Ù˜VÕ$øð ð%øÐ(<ô —-‘- Ó%×%Ñ%ð �n‹nús2   ‚C›B"ŸB  B"£A=CÂ B"Â"CÂ;B>Â<Cc                 ó   — t         dfS )Nr<   )r_   r#   s    r   Ú
__reduce__z!DaskDistributedBackend.__reduce__Ý   s   € Ü&¨Ð+Ð+r&   c                 ó2   — t        | j                  ¬«      dfS )N)rj   r`   )r_   rj   r#   s    r   Úget_nested_backendz)DaskDistributedBackend.get_nested_backendà   s   € Ü%¨T¯[©[Ô9¸2Ð=Ð=r&   c                 ó2   — || _         | j                  |«      S r    )ÚparallelÚeffective_n_jobs)r$   Ún_jobsr“   Úbackend_argss       r   Ú	configurez DaskDistributedBackend.configureã   s   € Ø ˆŒØ×$Ñ$ VÓ,Ð,r&   c                 óŽ   — d| _         | j                  j                  j                  | j                  «       t        «       | _        y )NT)r   rj   rb   Úadd_callbackr�   r   Úcall_data_futuresr#   s    r   Ú
start_callz!DaskDistributedBackend.start_callç   s0   € ØˆŒØ�‰×Ñ×%Ñ% d§m¡mÔ4Ü!3Ó!5ˆÕr&   c                 óp   — d| _         t        j                  d«       | j                  j	                  «        y )NFr~   )r   Útimer…   rš   r7   r#   s    r   Ú	stop_callz DaskDistributedBackend.stop_callì   s+   € ð ˆŒô 	�
‰
�4ÔØ×Ñ×$Ñ$Õ&r&   c           	      ó   — t        | j                  j                  «       j                  «       «      }|dk7  s| j                  s|S 	 | j                  j                  t        «      j                  | j                  ¬«       t        | j                  j                  «       j                  «       «      S # t        $ rD}dj                  | j                  t        dd| j                  z  «      «      }t        |«      |‚d }~ww xY w)Nr   )Útimeoutz÷DaskDistributedBackend has no worker after {} seconds. Make sure that workers are started and can properly connect to the scheduler and increase the joblib/dask connection timeout with:

parallel_config(backend='dask', wait_for_workers_timeout={})é
   é   )Úsumrj   ÚncoresÚvaluesrq   Úsubmitr]   r‡   Ú_TimeoutErrorÚformatÚmaxr   )r$   r•   r”   rx   Ú	error_msgs        r   r”   z'DaskDistributedBackend.effective_n_jobsö   sã   € Ü˜tŸ{™{×1Ñ1Ó3×:Ñ:Ó<Ó=ÐØ˜qÒ ¨×(EÒ(EØ#Ð#ð
	1Ø�K‰K×ÑÔ1Ó2×9Ñ9Ø×5Ñ5ð :ô ô �4—;‘;×%Ñ%Ó'×.Ñ.Ó0Ó1Ð1øô ò 	1ðO÷
 ‰fØ×-Ñ-Ü�B˜˜D×9Ñ9Ñ9Ó:óð ô ˜yÓ)¨qÐ0ûð	1ús   Á9B0 Â0	C=Â9?C8Ã8C=c           
   ƒ   ót  ‡ ‡‡K  — t        «       Št        ‰ dd «      Šˆˆˆ fd„}g }|j                  D ]r  \  }}}t         ||«      ƒ d {  –—† «      }t        t	        |j                  «        ||j                  «       «      ƒ d {  –—† «      «      }|j                  |||f«       Œt t        |«      |fS 7 Œj7 Œ1­w)Nrš   c              “   óð  •K  — g }| D ]Ö  }t        |«      }|‰v r|j                  ‰|   «       Œ'‰	j                  j                  |d «      }|€m‰�k	 ‰|   ƒ d {  –—† }|€[t        |«      rPt        |«      dkD  rB‰	j                  j                  |dd¬«      }t        j                  |«      }|‰|<   |ƒ d {  –—† }|�|j                  |«       ŒÆ|j                  |«       ŒØ |S 7 ŒŠ# t        $ r Y Œ“w xY w7 Œ>­w)Ng     @�@TF)ÚasynchronousÚhash)r(   rS   rp   Úgetr)   r   r   rj   rn   r„   ÚTask)
rF   ÚoutÚargÚarg_idrz   Ú_coroÚtrš   Úitemgettersr$   s
          €€€r   Úmaybe_to_futuresz>DaskDistributedBackend._to_func_args.<locals>.maybe_to_futures  s  øè ø€ ØˆCÛ�Ü˜C›�Ø˜[Ñ(Ø—J‘J˜{¨6Ñ2Ô3Øà×%Ñ%×)Ñ)¨&°$Ó7�Ø�9Ð!2Ð!>ðØ"3°CÑ"8×8˜ð �yÜ)¨#Ô.´6¸#³;ÀÒ3Dð %)§K¡K×$7Ñ$7Ø #°$¸Uð %8ó %˜Eô !(§¡¨UÓ 3˜AØ56Ð-¨cÑ2à&'§˜Aà�=Ø—J‘J˜q•Mà—J‘J˜s•OðM ðN ˆJð= 9ùÜ#ò Ùðúð. !(úsI   ƒAC6ÁC%ÁC#ÁC%ÁAC6Â5C4Â6-C6Ã#C%Ã%	C1Ã.C6Ã0C1Ã1C6)	ÚdictÚgetattrÚitemsr?   ro   Úkeysr¥   rS   rL   )	r$   rE   r·   rD   rz   rF   rG   rš   r¶   s	   `      @@r   Ú_to_func_argsz$DaskDistributedBackend._to_func_args  s§   úè ø€ Ü“fˆô $ DÐ*=¸tÓDÐö)	ðV ˆØ#Ÿzœz‰OˆAˆt�VÜÑ.¨tÓ4×4Ó5ˆDÜœ#˜fŸk™k›mÑ3CÀFÇMÁMÃOÓ3T×-TÓUÓVˆFØ�L‰L˜!˜T 6Ð*Õ+ð  *ô
 �e“˜eÐ$Ð$ð	 5øØ-Tús$   …AB8Á	B4
Á
:B8ÂB6Â0B8Â6B8c                 óÂ   ‡ ‡— t         j                  j                  «       Š‰j                  ‰_        ˆˆ fd„}‰ j
                  j                  j                  |||«       ‰S )Nc              “   óf  •K  — ‰j                  | «      ƒ d {  –—† \  }}t        |«      › dt        «       j                  › �} ‰j                  j
                  t        |«      f||dœ‰j                  ¤Ž}‰j                  j                  |«       |‰j                  |<   ‰‰j                  |<   y 7 Œ–­w)NÚ-)rD   r/   )r¼   Úreprr   Úhexrj   r¦   r	   rr   rs   Úaddru   rt   )rE   r‰   ÚbatchrD   r/   Údask_futurerˆ   r$   s         €€r   rz   z-DaskDistributedBackend.apply_async.<locals>.fN  s¨   øè ø€ Ø!%×!3Ñ!3°DÓ!9×9‰LˆE�5Ü˜%“[�M ¤5£7§;¡; -Ð0ˆCà,˜$Ÿ+™+×,Ñ,Ü*¨5Ó1ðàØñð ×$Ñ$ñ	ˆKð × Ñ ×$Ñ$ [Ô1Ø+3ˆD�O‰O˜KÑ(Ø)2ˆD�M‰M˜+Ò&ð :ús   ƒB1˜B/™BB1)Ú
concurrentÚfuturesÚFuturer‡   r¯   rj   rb   r™   )r$   rE   r‰   rz   rˆ   s   `   @r   Úapply_asyncz"DaskDistributedBackend.apply_asyncJ  sM   ù€ Ü×&Ñ&×-Ñ-Ó/ˆ	Ø!×(Ñ(ˆ	Œõ	3ð 	�‰×Ñ×%Ñ% a¨¨xÔ8àÐr&   c                 ó   — t        |«      S r    )r   )r$   r±   s     r   Úretrieve_result_callbackz/DaskDistributedBackend.retrieve_result_callback`  s   € Ü9¸#Ó>Ð>r&   c                 ó|  — | j                   j                  5  | j                   j                  j                  «        | j                   j                  j                  «       sI| j                   j                  j                  «        | j                   j                  j                  «       sŒIddd«       y# 1 sw Y   yxY w)z€Tell the client to cancel any task submitted via this instance

        joblib.Parallel will never access those results
        N)rs   ÚlockrÆ   r7   ÚqueueÚemptyr¯   )r$   Úensure_readys     r   Úabort_everythingz'DaskDistributedBackend.abort_everythingc  s�   € ð
 ×!Ñ!×&Ó&Ø× Ñ ×(Ñ(×.Ñ.Ô0Ø×*Ñ*×0Ñ0×6Ñ6Ô8Ø×$Ñ$×*Ñ*×.Ñ.Ô0ð ×*Ñ*×0Ñ0×6Ñ6Õ8÷ '×&Ñ&ús   —BB2Â2B;c              #   ó~   K  — t        t        d«      r
t        «        d–— t        t        d«      rt        «        yy­w)zÙOverride ParallelBackendBase.retrieval_context to avoid deadlocks.

        This removes thread from the worker's thread pool (using 'secede').
        Seceding avoids deadlock in nested parallelism settings.
        Úexecution_stateN)Úhasattrr   r   r   r#   s    r   Úretrieval_contextz(DaskDistributedBackend.retrieval_contextm  s0   è ø€ ô ”<Ð!2Ô3äŒHãä”<Ð!2Ô3Ü�Hð 4ùs   ‚;=)NNNNr¡   )r   Nr    )T)r8   r9   r:   ÚMIN_IDEAL_BATCH_DURATIONÚMAX_IDEAL_BATCH_DURATIONÚsupports_retrieve_callbackÚdefault_n_jobsr%   r�   r�   r‘   r—   r›   rž   r”   r¼   rÈ   rÊ   rÐ   Ú
contextlibÚcontextmanagerrÔ   Ú__classcell__)r{   s   @r   r_   r_   �   sƒ   ø„ Ø"ÐØ"ÐØ!%ÐØ€Nð ØØØØ!#õ9òv&ò,ò>ó-ò6ò
'ò2ò48%ótò,?ó1ð ×Ññó ôr&   r_   ),Ú
__future__r   r   r   r„   Úconcurrent.futuresrÅ   rÙ   r�   r   Úuuidr   Ú_utilsr   r	   r“   r
   r   r   rQ   rh   ÚImportErrorÚdask.distributedr   r   r   r   r   Údask.sizeofr   Ú
dask.utilsr   Údistributed.utilsr   r   r§   Útornado.genr   r   rB   rJ   rL   r]   r_   r<   r&   r   Ú<module>ræ      sÍ   ðß @Ñ @ã Û Û Û Û Ý ÷÷ NÑ MðÛÛð
 Ð˜Ð/÷õ õ #Ý#Ý.ð>õ 	Dò
÷*ñ *òZò/÷ñ ò,	ô
nÐ.Ð0Cõ nøðy ò Ø€DØ‚Kðûð( ò >ß=ð>ús#   ¸B Á%B' Â	B$Â#B$Â'B5Â4B5