o
    .j,-                     @  s|   d dl mZ d dl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 G d	d
 d
ZG dd dZG dd de	ZdS )    )annotationsN   )HTTPHeaderDict)AsyncSocketSSLAsyncSocket   )BaseBackendResponsePromise)BytesQueueBufferc                   @  s   e Zd Z		d2d3d
dZed4ddZd5ddZd4ddZd4ddZd4ddZ	d6ddZ
d6ddZd7d8d!d"Zd9d%d&Zd7d:d(d)Zd;d*d+Z	,d<d=d.d/Zd>d0d1ZdS )?AsyncDirectStreamAccessN	stream_idintreadtyping.Callable[[int | None, int | None, bool, bool], typing.Awaitable[tuple[list[bytes], bool, HTTPHeaderDict | None]]] | NonewriteBtyping.Callable[[bytes, int, bool], typing.Awaitable[None]] | NonereturnNonec                 C  s$   || _ || _|| _t | _d| _d S NF)
_stream_id_read_writer
   _buffer_eot)selfr   r   r    r   Q/home/thesage/.local/lib/python3.10/site-packages/urllib3/backend/_async/_base.py__init__   s
   
z AsyncDirectStreamAccess.__init__boolc                 C  s   | j d u o	| jd u S N)r   r   r   r   r   r   closed   s   zAsyncDirectStreamAccess.closedb	bytearrayc                   sP   | j d u r
td| t|I d H }t|dkrdS ||d t|< t|S )Nz!read operation on a closed streamr   )r   OSErrorrecvlen)r   r"   tempr   r   r   readinto"   s   
z AsyncDirectStreamAccess.readintoc                 C  
   | j d uS r   )r   r    r   r   r   readable.      
z AsyncDirectStreamAccess.readablec                 C  r)   r   )r   r    r   r   r   writable1   r+   z AsyncDirectStreamAccess.writablec                 C     dS r   r   r    r   r   r   seekable4      z AsyncDirectStreamAccess.seekablec                 C  r-   Nr   r    r   r   r   fileno7   r/   zAsyncDirectStreamAccess.filenoc                 C  r-   r0   r   r    r   r   r   name:   r/   zAsyncDirectStreamAccess.namer   !_AsyncDirectStreamAccess__bufsize_AsyncDirectStreamAccess__flagsbytesc                   s   |  |I d H \}}}|S r   )recv_extended)r   r4   r5   data_r   r   r   r%   =   s   zAsyncDirectStreamAccess.recv
int | None)tuple[bytes, bool, HTTPHeaderDict | None]c                   s   | j d u r
tdd }| js*| js*|  || j|d udI d H \}| _}| j| | jrA| j|d ur:|dkr:|nt| j}nd}| joI| j }|rOd | _ |||fS )Nzstream closed errorFr       )r   r$   r   r   r   put_manygetr&   )r   r4   trailerschunksr8   eotr   r   r   r7   A   s.   

z%AsyncDirectStreamAccess.recv_extended_AsyncDirectStreamAccess__datac                   s.   | j d u r
td|  || jdI d H  d S Nstream write not permittedFr   r$   r   )r   rB   r5   r   r   r   sendalld   s   
zAsyncDirectStreamAccess.sendallc                   s2   | j d u r
td|  || jdI d H  t|S rC   )r   r$   r   r&   )r   rB   r   r   r   r   j   s
   
zAsyncDirectStreamAccess.writeF&_AsyncDirectStreamAccess__close_streamc                   s.   | j d u r
td|  || j|I d H  d S )NrD   rE   )r   rB   rG   r   r   r   sendall_extendedr   s   
z(AsyncDirectStreamAccess.sendall_extendedc                   sX   | j d ur|  d| jdI d H  d | _ | jd ur*| d | jddI d H  d | _d S d S )Nr<   TF)r   r   r   r    r   r   r   closez   s   


zAsyncDirectStreamAccess.close)NN)r   r   r   r   r   r   r   r   r   r   )r"   r#   r   r   )r   r   )r   )r4   r   r5   r   r   r6   )r4   r:   r   r;   )rB   r6   r5   r   r   r   )rB   r6   r   r   )F)rB   r6   rG   r   r   r   r   r   )__name__
__module____qualname__r   propertyr!   r(   r*   r,   r.   r2   r3   r%   r7   rF   r   rH   rI   r   r   r   r   r      s&    






#
	r   c                   @  s   e Zd ZU dZded< ddddddd2ddZed3ddZed4d d!Zej	d5d$d!Zed6d%d&Z
d7d(d)Zd8d9d,d-Zd:d.d/Zd:d0d1ZdS );AsyncLowLevelResponsezImplemented for backward compatibility purposes. It is there to impose http.client like
    basic response object. So that we don't have to change urllib3 tested behaviors.styping.Callable[[int | None, int | None], typing.Awaitable[tuple[list[bytes], bool, HTTPHeaderDict | None]]] | None(_AsyncLowLevelResponse__internal_read_stN)	authorityportr   dsastream_abortmethodstrstatusr   versionreasonheadersr   bodyrS   
str | NonerT   r:   r   rU   AsyncDirectStreamAccess | NonerV   5typing.Callable[[int], typing.Awaitable[None]] | Noner   r   c                C  s   || _ || _|| _|| _|| _|| _| jd u}|du | _| j| _|| _|| _	d| _
| jdko5d| jdk| _d | _d | _d| _| jsR| jd}|rOt|nd | _d| _|	| _t | _d | _|
| _|| _d | _d S )NFr      chunkedztransfer-encodingzcontent-length)rY   rZ   r[   msg_methodrR   r!   r   rS   rT   
debuglevelr>   rb   
chunk_leftlength
will_closer   data_in_countr   r
   %_AsyncLowLevelResponse__buffer_excess_AsyncLowLevelResponse__promise_dsa_stream_abortr?   )r   rW   rY   rZ   r[   r\   r]   rS   rT   r   rU   rV   has_bodycontent_lengthr   r   r   r      s8   


zAsyncLowLevelResponse.__init__typing.NoReturnc                 C  s   t d)Nzurllib3-future no longer expose a filepointer-like in responses. It was a remnant from the http.client era. We no longer support it.)RuntimeErrorr    r   r   r   fp   s   zAsyncLowLevelResponse.fpResponsePromise | Nonec                 C     | j S r   )rk   r    r   r   r   from_promise      z"AsyncLowLevelResponse.from_promisevaluer	   c                 C  s   |j | jkr
td|| _d S )NzCTrying to assign a ResponsePromise to an unrelated LowLevelResponse)r   r   
ValueErrorrk   )r   rw   r   r   r   ru      s
   
c                 C  rt   )z'Original HTTP verb used in the request.)rd   r    r   r   r   rW      s   zAsyncLowLevelResponse.methodr   c                 C  rt   )z:Here we do not create a fp sock like http.client Response.)r!   r    r   r   r   isclosed   rv   zAsyncLowLevelResponse.isclosed_AsyncLowLevelResponse__sizer6   c                   s&  | j du s| jd u rtd|dkrdS t| j}|d uo%|dko%||k}| jdu rG|sG| || jI d H \}| _| _| j| t| j}|rY| j	|d urV|dkrV|n|nd}t|}||8 }| jrs|dkrsd | _
d| _ d | _| jr~|rz|nd | _n| jd ur|  j|8  _|  j|7  _|S )NTzI/O operation on closed file.r   r<   F)r!   rR   rx   r&   rj   r   r   r?   r=   r>   rm   _sockrb   rf   rg   ri   )r   rz   buf_capacitydata_ready_to_gor@   r8   size_inr   r   r   r      sD   


zAsyncLowLevelResponse.readc                   sV   | j d ur'| jdu r)| jd ur|  | jI d H  d| _d | _ d| _d | _d S d S d S )NFT)rm   r   r   r!   rl   r    r   r   r   abort  s   



zAsyncLowLevelResponse.abortc                 C  s   d | _ d| _d | _d S )NT)rR   r!   rl   r    r   r   r   rI   &  s   
zAsyncLowLevelResponse.close)rW   rX   rY   r   rZ   r   r[   rX   r\   r   r]   rQ   rS   r^   rT   r:   r   r:   rU   r_   rV   r`   r   r   )r   rp   )r   rs   )rw   r	   r   r   )r   rX   rJ   r   )rz   r:   r   r6   rK   )rL   rM   rN   __doc____annotations__r   rO   rr   ru   setterrW   ry   r   r   rI   r   r   r   r   rP      s*   
 @

1
rP   c                   @  s   e Zd ZU ded< d(ddZd(ddZd)d
dZd(ddZ	d*dddd+ddZddd,ddZ	d(ddZ
dd d-d$d%Zd(d&d'ZdS ).AsyncBaseBackendz#AsyncSocket | SSLAsyncSocket | Nonesockr   r   c                      t )z+Upgrade conn from svn ver to max supported.NotImplementedErrorr    r   r   r   _upgrade/     zAsyncBaseBackend._upgradec                   r   )z>Emit proper CONNECT request to the http (server) intermediary.r   r    r   r   r   _tunnel3  r   zAsyncBaseBackend._tunnelAsyncSocket | Nonec                   r   )zRun protocol initialization from there. Return None to ensure that the child
        class correctly create the socket / connection.r   r    r   r   r   	_new_conn7     zAsyncBaseBackend._new_connc                   r   )zhShould be called after _new_conn proceed as expected.
        Expect protocol handshake to be done here.r   r    r   r   r   
_post_conn<  r   zAsyncBaseBackend._post_connNF)encode_chunkedexpect_body_afterwardmessage_bodybytes | Noner   r   r   rs   c                  r   )z6This method conclude the request context construction.r   )r   r   r   r   r   r   r   
endheadersA     zAsyncBaseBackend.endheaders)promiser   rP   c                  r   )zFetch the HTTP response. You SHOULD not retrieve the body in that method, it SHOULD be done
        in the LowLevelResponse, so it enable stream capabilities and remain efficient.
        r   )r   r   r   r   r   getresponseK  s   zAsyncBaseBackend.getresponsec                   r   )z9End the connection, do some reinit, closing of fd, etc...r   r    r   r   r   rI   S  r   zAsyncBaseBackend.close)rA   r8   <bytes | typing.IO[typing.Any] | typing.Iterable[bytes] | strrA   c                  r   )zThe send() method SHOULD be invoked after calling endheaders() if and only if the request
        context specify explicitly that a body is going to be sent.r   )r   r8   rA   r   r   r   sendW  r   zAsyncBaseBackend.sendc                   r   r   r   r    r   r   r   pinga  s   zAsyncBaseBackend.pingrK   )r   r   r   )r   r   r   r   r   r   r   rs   )r   rs   r   rP   )r8   r   rA   r   r   rs   )rL   rM   rN   r   r   r   r   r   r   r   rI   r   r   r   r   r   r   r   ,  s"   
 





r   )
__future__r   typing_collectionsr   contrib.ssar   r   _baser   r	   util.responser
   r   rP   r   r   r   r   r   <module>   s    x *