o
    Z�¿j[/  ã                   @   sb  U 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	 ddl
mZmZ ddlmZmZ ddlmZmZmZ dd	lmZ d
gZejdkr^e	ƒ Ze	ed< dZde jde jfdd„Znde jde jfdd„Zd ZdZdZdZ dZ!dZ"dZ#dZ$dZ%dZ&dZ'd Z(dZ)dZ*	 e +¡ de,dede,fdd„ƒZ-e +¡ de,de,dedede,f
dd „ƒZ.G d!d
„ d
ƒZ/dS )"é    N)Úsuppress)ÚAnyÚOptional)ÚWeakKeyDictionaryé   )ÚffiÚlib)Ú	CurlECodeÚCurlMOpt)ÚDEFAULT_CACERTÚCurlÚ	CurlError)ÚCurlCffiWarningÚ	AsyncCurlÚwin32Ú
_selectorsa  
    Proactor event loop does not implement add_reader family of methods required.
    Registering an additional selector thread for add_reader support.
    To avoid this warning use:
        asyncio.set_event_loop_policy(WindowsSelectorEventLoopPolicy())
    Úasyncio_loopÚreturnc                    sv   ˆ t v rt ˆ  S tˆ ttdtdƒƒƒsˆ S tjttdd� ddl	m
} |ˆ ƒ ‰t ˆ < ˆ j‰‡ ‡‡fdd„}|ˆ _ˆS )	zåGet selector-compatible loop

        Returns an object with ``add_reader`` family of methods,
        either the loop itself or a SelectorThread instance.

        Workaround Windows proactor removal of *reader methods.
        ÚProactorEventLoopNé   ©Ú
stacklevelr   )ÚAddThreadSelectorEventLoopc                      s   ˆˆ _ t ˆ d ¡ ˆ  ¡  d S ©N)Úcloser   Úpop© ©r   Ú
loop_closeÚselector_loopr   úH/home/dinkstrade/pdmp/venv/lib/python3.10/site-packages/curl_cffi/aio.pyÚ_close_selector_and_loop7   s   z.get_selector.<locals>._close_selector_and_loop)r   Ú
isinstanceÚgetattrÚasyncioÚtypeÚwarningsÚwarnÚPROACTOR_WARNINGr   Ú_asyncio_selectorr   r   )r   r   r!   r   r   r    Úget_selector   s   ÿÿr*   Úloopc                 C   s   | S r   r   )r+   r   r   r    r*   C   s   r   é   é   éÿÿÿÿÚ
timeout_msÚclientpc                 C   s>   t  |¡}|jr|j ¡  d|_|j |d |jtt¡|_dS )zD
    see: https://curl.se/libcurl/c/CURLMOPT_TIMERFUNCTION.html
    Niè  r   )	r   Úfrom_handleÚ_timerÚcancelr+   Ú
call_laterÚprocess_dataÚCURL_SOCKET_TIMEOUTÚCURL_POLL_NONE)Úcurlmr/   r0   Ú
async_curlr   r   r    Útimer_functionu   s   

ür:   ÚsockfdÚwhatÚdatac                 C   s’   t  |¡}|j}||jv r| |¡ | |¡ |t@ r*| ||j|t	¡ |j 
|¡ |t@ r=| ||j|t¡ |j 
|¡ |tkrG|j |¡ dS )z[This callback is called when libcurl decides it's time to interact with certain
    socketsr   )r   r1   r+   Ú_sockfdsÚremove_readerÚremove_writerÚCURL_POLL_INÚ
add_readerr5   ÚCURL_CSELECT_INÚaddÚCURL_POLL_OUTÚ
add_writerÚCURL_CSELECT_OUTÚCURL_POLL_REMOVEÚremove)Úcurlr;   r<   r0   r=   r9   r+   r   r   r    Úsocket_function�   s   



rK   c                   @   sÀ   e Zd ZdZd%defdd„Zdd„ Zd	d
„ Zdd„ Zde	fdd„Z
dededefdd„Zdedefdd„Zde	fdd„Zde	fdd„Zde	fdd„Zde	fdd„Zded efd!d"„Zd#d$„ ZdS )&r   zhWrapper around curl_multi handle to provide asyncio support. It uses the libcurl
    socket_action APIs.Ú NÚcacertc                 C   sf   t  ¡ | _|pt| _i | _i | _tƒ | _t	|dur|nt
 ¡ ƒ| _| j |  ¡ ¡| _d| _|  ¡  dS )z—
        Parameters:
            cacert: CA cert path to use, by default, certs from ``certifi`` are used.
            loop: EventLoop to use.
        N)r   Úcurl_multi_initÚ_curlmr   Ú_cacertÚ_curl2futureÚ
_curl2curlÚsetr>   r*   r$   Úget_running_loopr+   Úcreate_taskÚ_force_timeoutÚ_timeout_checkerr2   Ú_setup)ÚselfrM   r+   r   r   r    Ú__init__¯   s   

ÿzAsyncCurl.__init__c                 C   sP   |   tjtj¡ |   tjtj¡ t | ¡| _	|   tj
| j	¡ |   tj| j	¡ d S r   )Úsetoptr
   ÚTIMERFUNCTIONr   r:   ÚSOCKETFUNCTIONrK   r   Ú
new_handleÚ_self_handleÚ
SOCKETDATAÚ	TIMERDATA©rY   r   r   r    rX   Á   s
   zAsyncCurl._setupc                 Ã   sÎ   �| j  ¡  ttjƒ� | j I dH  W d  ƒ n1 sw   Y  | j ¡ D ]\}}t | j	|j
¡ | ¡ s?| ¡ s?| d¡ q&t | j	¡ d| _	| jD ]}| j |¡ | j |¡ qL| jre| j ¡  dS dS )z?Close and cleanup running timers, readers, writers and handles.N)rW   r3   r   r$   ÚCancelledErrorrQ   Úitemsr   Úcurl_multi_remove_handlerO   Ú_curlÚdoneÚ	cancelledÚ
set_resultÚcurl_multi_cleanupr>   r+   r?   r@   r2   )rY   rJ   Úfuturer;   r   r   r    r   É   s$   €
ÿ
€
ÿzAsyncCurl.closec                 Ã   s,   �	 | j sdS |  tt¡ t d¡I dH  q)zpThis coroutine is used to safeguard from any missing signals from curl, and
        put everything back on trackTgš™™™™™¹?N)rO   Úsocket_actionr6   r7   r$   Úsleeprb   r   r   r    rV   ä   s   €üzAsyncCurl._force_timeoutrJ   c                 C   sF   |  ¡  t | j|j¡}|  |¡ | j ¡ }|| j|< || j	|j< |S )znAdd a curl handle to be managed by curl_multi. This is the equivalent of
        `perform` in the async world.)
Ú_ensure_cacertr   Úcurl_multi_add_handlerO   rf   Ú_check_errorr+   Úcreate_futurerQ   rR   )rY   rJ   Úerrcoderk   r   r   r    Ú
add_handleí   s   


zAsyncCurl.add_handler;   Ú
ev_bitmaskr   c                 C   s.   t  d¡}t | j|||¡}|  |¡ |d S )zYwrapper for curl_multi_socket_action,
        returns the number of running curl handles.úint *r   )r   Únewr   Úcurl_multi_socket_actionrO   rp   )rY   r;   rt   Úrunning_handlerr   r   r   r    rl   ù   s   

ÿ
zAsyncCurl.socket_actionc                 C   sè   | j stjdtdd� dS |  ||¡ t d¡}	 zHt | j |¡}|tj	kr)W dS |j
tkr\| j|j }|jj}| ¡ }|durG|  ||¡ n|dkrQ|  |¡ n|  || |d¡¡ ntd	ƒ W n tyr   tjd
tdd� Y nw q)z8Call curl_multi_info_read to read data for given socket.z0Curlm already closed! quitting from process_datar   r   Nru   Tr   ÚperformzNOT DONEzLUnexpected curl multi state in process_data, please open an issue on GitHub
)rO   r&   r'   r   rl   r   rv   r   Úcurl_multi_info_readÚNULLÚmsgÚCURLMSG_DONErR   Úeasy_handler=   ÚresultÚ_get_callback_exceptionÚset_exceptionri   Ú
_get_errorÚprintÚ	Exception)rY   r;   rt   Úmsg_in_queueÚcurl_msgrJ   ÚretcodeÚcallback_exceptionr   r   r    r5     sB   ý


€
üÿîzAsyncCurl.process_datac                 C   s8   t  | j|j¡}|  |¡ | j |jd ¡ | j |d ¡S r   )r   re   rO   rf   rp   rR   r   rQ   )rY   rJ   rr   r   r   r    Ú_pop_future*  s   
zAsyncCurl._pop_futurec                 C   s6   |   |¡}|r| ¡ s| ¡ s| ¡  dS dS dS dS )z&Cancel a future for given curl handle.N)r‰   rg   rh   r3   ©rY   rJ   rk   r   r   r    Úremove_handle0  s   
ÿzAsyncCurl.remove_handlec                 C   s8   |   |¡}|r| ¡ s| ¡ s| d¡ dS dS dS dS )z,Mark a future as done for given curl handle.N)r‰   rg   rh   ri   rŠ   r   r   r    ri   6  ó   
ÿzAsyncCurl.set_resultc                 C   s8   |   |¡}|r| ¡ s| ¡ s| |¡ dS dS dS dS )z2Raise exception of a future for given curl handle.N)r‰   rg   rh   r�   )rY   rJ   Ú	exceptionrk   r   r   r    r�   <  rŒ   zAsyncCurl.set_exceptionrr   Úargsc                 G   sH   |t jkrd S t |¡}d dd„ |D ƒ¡}td|› d|› d|› d�ƒ‚)Nú c                 S   s   g | ]}t |ƒ‘qS r   )Ústr)Ú.0Úar   r   r    Ú
<listcomp>F  s    z*AsyncCurl._check_error.<locals>.<listcomp>z
Failed in z
, multi: (z) z„. See https://curl.se/libcurl/c/libcurl-errors.html first for more details. Please open an issue on GitHub to help debug this error.)r	   ÚOKr   Úcurl_multi_strerrorÚjoinr   )rY   rr   rŽ   ÚerrmsgÚactionr   r   r    rp   B  s   

ÿzAsyncCurl._check_errorc                 C   sB   |t jt jt jt jt jt jfv rt d|¡}n|}t	 
| j||¡S )z!Wrapper around curl_multi_setopt.zlong*)r
   Ú
PIPELININGÚMAXCONNECTSÚMAX_HOST_CONNECTIONSÚMAX_PIPELINE_LENGTHÚMAX_TOTAL_CONNECTIONSÚMAX_CONCURRENT_STREAMSr   rv   r   Úcurl_multi_setoptrO   )rY   ÚoptionÚvalueÚc_valuer   r   r    r[   M  s   úzAsyncCurl.setopt)rL   N)Ú__name__Ú
__module__Ú__qualname__Ú__doc__r�   rZ   rX   r   rV   r   rs   Úintrl   r5   r‰   r‹   ri   r�   r   rp   r[   r   r   r   r    r   «   s    	
')0r$   Úsysr&   Ú
contextlibr   Útypingr   r   Úweakrefr   Ú_wrapperr   r   Úconstr	   r
   rJ   r   r   r   Úutilsr   Ú__all__Úplatformr   Ú__annotations__r(   ÚAbstractEventLoopr*   r7   rA   rE   ÚCURL_POLL_INOUTrH   r6   ÚCURL_SOCKET_BADrC   rG   ÚCURL_CSELECT_ERRr}   ÚCURLPIPE_NOTHINGÚCURLPIPE_HTTP1ÚCURLPIPE_MULTIPLEXÚ
def_externr§   r:   rK   r   r   r   r   r    Ú<module>   sP   
 
ÿþ* 