U
    7ÄT_@4  ã                   @   s®   d dl 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mZ ddl	m
Z
 ddlmZ zd dlZW n ek
r|   dZY nX dd„ ZG dd	„ d	ƒZG d
d„ dƒZeZeZdS )é    N)ÚFutureÚThreadPoolExecutoré   )ÚCurrentThreadExecutor)ÚLocalc              	   C   sZ   | D ]P}z&|  ¡ |   |¡kr,| |   |¡¡ W q tk
rR   | |   |¡¡ Y qX qd S ©N)ÚgetÚsetÚLookupError)ÚcontextZcvar© r   ú0/tmp/pip-unpacked-wheel-_lltyzrh/asgiref/sync.pyÚ_restore_context   s    r   c                   @   sD   e Zd ZdZi Zeƒ Zddd„Zdd„ Zdd„ Z	d	d
„ Z
dd„ ZdS )ÚAsyncToSyncaè  
    Utility class which turns an awaitable that only works on the thread with
    the event loop into a synchronous callable that works in a subthread.

    If the call stack contains an async loop, the code runs there.
    Otherwise, the code runs in a new loop in a new thread.

    Either way, this thread then pauses and waits to run any thread_sensitive
    code called from further down the call stack using SyncToAsync, before
    finally exiting once the async task returns.
    Fc                 C   sn   || _ z| j j| _W n tk
r(   Y nX |r6d | _n4zt ¡ | _W n$ tk
rh   ttj	dd ƒ| _Y nX d S )NÚmain_event_loop)
Ú	awaitableÚ__self__ÚAttributeErrorr   ÚasyncioÚget_event_loopÚRuntimeErrorÚgetattrÚSyncToAsyncÚthreadlocal)Úselfr   Zforce_new_loopr   r   r   Ú__init__1   s      ÿzAsyncToSync.__init__c              	   O   sL  zt  ¡ }W n tk
r    Y nX | ¡ r2tdƒ‚td k	rFt ¡ g}nd }tƒ }t ¡ }t	| j
dƒrn| j
j}nd }tƒ }|| j
_zˆ|  ||||t ¡ |¡}	| jrª| j ¡ sät  ¡ }
tdd�}| | j|
|	¡}|rÚ| |¡ | ¡  n"| j | jj|	¡ |�r| |¡ W 5 t	| j
dƒ�r| j
`|�r,|| j
_td k	�rBt|d ƒ X | ¡ S )NznYou cannot use AsyncToSync in the same thread as an async event loop - just await the async function directly.Úcurrentr   r   ©Úmax_workers)r   r   r   Z
is_runningÚcontextvarsÚcopy_contextr   Ú	threadingÚcurrent_threadÚhasattrÚ	executorsr   r   r   Ú	main_wrapÚsysÚexc_infor   Znew_event_loopr   ZsubmitÚ_run_event_loopZrun_until_futureÚresultZcall_soon_threadsafeZcreate_task)r   ÚargsÚkwargsZ
event_loopr   Úcall_resultÚsource_threadZold_current_executorZcurrent_executorr   ÚloopZloop_executorZloop_futurer   r   r   Ú__call__D   sf    ÿ
     ÿ
  ÿ

 ÿ
zAsyncToSync.__call__c                 C   sÔ   t  |¡ z| 	|¡ W 5 zœtjdkr2t  |¡}nt j |¡}|D ]}| ¡  qB| 	t j
|ddiŽ¡ |D ]0}| ¡ rxqj| ¡ dk	rj| d| ¡ |dœ¡ qjt|dƒr´| 	| ¡ ¡ W 5 | ¡  t  | j¡ X X dS )zP
        Runs the given event loop (designed to be called in a thread).
        )é   é   r   Zreturn_exceptionsTNz(unhandled exception during loop shutdown)ÚmessageÚ	exceptionÚtaskÚshutdown_asyncgens)r   Zset_event_loopÚcloser   r&   Úversion_infoZ	all_tasksÚTaskÚcancelZrun_until_completeZgatherZ	cancelledr3   Zcall_exception_handlerr#   r5   )r   r.   ÚcoroZtasksr4   r   r   r   r(   �   s0    


ýÿ
zAsyncToSync._run_event_loopc                 C   s   t  | j|¡}t  || j¡S ©z*
        Include self for methods
        )Ú	functoolsÚpartialr/   Úupdate_wrapperr   )r   ÚparentÚobjtypeÚfuncr   r   r   Ú__get__°   s    zAsyncToSync.__get__c           
   
   Ã   sÒ   |dk	rt |d ƒ t ¡ }|| j|< zˆzL|d r`z|d ‚W qr   | j||ŽI dH }Y qrX n| j||ŽI dH }W n, tk
r  }	 z| |	¡ W 5 d}	~	X Y nX | 	|¡ W 5 | j|= |dk	rÌt ¡ |d< X dS )zs
        Wraps the awaitable with something that puts the result into the
        result/exception future.
        Nr   r   )
r   r   Úget_current_taskÚ
launch_mapr   r    r   Ú	ExceptionZset_exceptionZ
set_result)
r   r*   r+   r,   r-   r'   r   Úcurrent_taskr)   Úer   r   r   r%   ·   s"    
zAsyncToSync.main_wrapN)F)Ú__name__Ú
__module__Ú__qualname__Ú__doc__rD   r   r$   r   r/   r(   rB   r%   r   r   r   r   r      s   
I#r   c                   @   s€   e Zd ZdZdejkr8e ¡ Ze 	e
eejd ƒd�¡ i Ze ¡ Ze
dd�Zddd„Zdd	„ Zd
d„ Zdd„ Zedd„ ƒZdS )r   aK  
    Utility class which turns a synchronous callable into an awaitable that
    runs in a threadpool. It also sets a threadlocal inside the thread so
    calls to AsyncToSync can escape it.

    If thread_sensitive is passed, the code will run in the same thread as any
    outer code. This is needed for underlying Python code that is not
    threadsafe (for example, code which handles SQLite database connections).

    If the outermost program is async (i.e. SyncToAsync is outermost), then
    this will be a dedicated single sub-thread that all sync code runs in,
    one after the other. If the outermost program is sync (i.e. AsyncToSync is
    outermost), this will just be the main thread. This is achieved by idling
    with a CurrentThreadExecutor while AsyncToSync is blocking its sync parent,
    rather than just blocking.
    ZASGI_THREADSr   r   Fc                 C   sH   || _ t | |¡ || _tjj| _z|j| _W n tk
rB   Y nX d S r   )	rA   r<   r>   Ú_thread_sensitiver   Z
coroutinesZ_is_coroutiner   r   )r   rA   Zthread_sensitiver   r   r   r   ú   s    
zSyncToAsync.__init__c           
   	   Ï   sÀ   t  ¡ }| jr,ttjdƒr$tjj}q0| j}nd }td k	rft 	¡ }t
j| jf|ž|Ž}|j}|f}i }n| j}| |t
j| j||  ¡ t ¡ |f|ž|Ž¡}t j|d d�I d H }	td k	r¼t|ƒ |	S )Nr   )Útimeout)r   r   rL   r#   r   r$   r   Úsingle_thread_executorr   r    r<   r=   rA   ÚrunZrun_in_executorÚthread_handlerrC   r&   r'   Úwait_forr   )
r   r*   r+   r.   Úexecutorr   ÚchildrA   ÚfutureÚretr   r   r   r/     s>    
ûúùþzSyncToAsync.__call__c                 C   s   t  | j|¡S r;   )r<   r=   r/   )r   r?   r@   r   r   r   rB   /  s    zSyncToAsync.__get__c           	      O   sŒ   || j _t ¡ }tj |¡|kr&d}n|| j|< d}zD|d rhz|d ‚W qv   |||Ž Y W ¢S X n|||ŽW ¢S W 5 |r†| j|= X dS )zE
        Wraps the sync application with exception handling.
        FTr   N)r   r   r!   r"   r   rD   r   )	r   r.   Zsource_taskr'   rA   r*   r+   r"   Z
parent_setr   r   r   rP   5  s    
zSyncToAsync.thread_handlerc                   C   s@   z$t tdƒrt ¡ W S tj ¡ W S W n tk
r:   Y dS X dS )zs
        Cross-version implementation of asyncio.current_task()

        Returns None if there is no task.
        rF   N)r#   r   rF   r8   r   r   r   r   r   rC   U  s    

zSyncToAsync.get_current_taskN)F)rH   rI   rJ   rK   ÚosÚenvironr   r   r.   Zset_default_executorr   ÚintrD   r!   Úlocalr   rN   r   r/   rB   rP   ÚstaticmethodrC   r   r   r   r   r   Ø   s   
ÿ


+ r   )r   Zasyncio.coroutinesr<   rV   r&   r!   Úconcurrent.futuresr   r   Zcurrent_thread_executorr   rY   r   r   ÚImportErrorr   r   r   Zsync_to_asyncZasync_to_syncr   r   r   r   Ú<module>   s&   
 < 