
Â�Â­Â�Â­Â�Â­Â�Â­Â�Â­Â�Â­Â�Â­Â�Â­Â�Â­Â�Â­Â�Â­Â�Â­Â�Â­Â�Â­Â�Â­Â�Â­Â�Â­Â�Â­
<!DOCTYPE html>
<html>
U
    'â7`2P  ã                   @   s8  d dl Z d dlZd dlZd dlmZmZmZmZmZm	Z	m
Z
 ddlmZ ddlmZmZmZ ddlmZ zd dlmZ W n  ek
r˜   d dlmZ Y nX dZe
d	ƒZG d
d„ deƒZG dd„ dee ƒZG dd„ dƒZG dd„ dƒZG dd„ deƒZG dd„ deƒZeƒ ZG dd„ dee ƒZ G dd„ de e ƒZ!dS )é    N)Ú	AwaitableÚCallableÚGenericÚListÚOptionalÚTupleÚTypeVaré   )ÚBaseProtocol)ÚBaseTimerContextÚset_exceptionÚ
set_result)Úinternal_logger)ÚDeque)ÚEMPTY_PAYLOADÚ	EofStreamÚStreamReaderÚ	DataQueueÚFlowControlDataQueueÚ_Tc                   @   s   e Zd ZdZdS )r   zeof stream indication.N)Ú__name__Ú
__module__Ú__qualname__Ú__doc__© r   r   úB/opt/alt/python38/lib64/python3.8/site-packages/aiohttp/streams.pyr      s   r   c                   @   sD   e Zd Zeg ee f ddœdd„Zddœdd„Zedœd	d
„ZdS )ÚAsyncStreamIteratorN)Ú	read_funcÚreturnc                 C   s
   || _ d S ©N)r   )Úselfr   r   r   r   Ú__init__   s    zAsyncStreamIterator.__init__zAsyncStreamIterator[_T]©r   c                 C   s   | S r   r   ©r    r   r   r   Ú	__aiter__"   s    zAsyncStreamIterator.__aiter__c                 Ã   s<   z|   ¡ I d H }W n tk
r*   t‚Y nX |dkr8t‚|S ©Nó    )r   r   ÚStopAsyncIteration©r    Úrvr   r   r   Ú	__anext__%   s    
zAsyncStreamIterator.__anext__)	r   r   r   r   r   r   r!   r$   r*   r   r   r   r   r      s   r   c                   @   s@   e Zd Zdddœdd„Zd dœdd„Zeeef dœd	d
„ZdS )ÚChunkTupleAsyncStreamIteratorr   N)Ústreamr   c                 C   s
   || _ d S r   )Ú_stream)r    r,   r   r   r   r!   0   s    z&ChunkTupleAsyncStreamIterator.__init__r"   c                 C   s   | S r   r   r#   r   r   r   r$   3   s    z'ChunkTupleAsyncStreamIterator.__aiter__c                 Ã   s    | j  ¡ I d H }|dkrt‚|S )N©r&   F)r-   Ú	readchunkr'   r(   r   r   r   r*   6   s    z'ChunkTupleAsyncStreamIterator.__anext__)	r   r   r   r!   r$   r   ÚbytesÚboolr*   r   r   r   r   r+   /   s   r+   c                   @   sR   e Zd Zee dœdd„Zeee dœdd„Zee dœdd„Ze	dœd	d
„Z
dS )ÚAsyncStreamReaderMixinr"   c                 C   s
   t | jƒS r   )r   Úreadliner#   r   r   r   r$   >   s    z AsyncStreamReaderMixin.__aiter__©Únr   c                    s   t ‡ ‡fdd„ƒS )zzReturns an asynchronous iterator that yields chunks of size n.

        Python-3.5 available for Python 3.5+ only
        c                      s
   ˆ  ˆ ¡S r   )Úreadr   ©r5   r    r   r   Ú<lambda>F   r&   z5AsyncStreamReaderMixin.iter_chunked.<locals>.<lambda>)r   ©r    r5   r   r7   r   Úiter_chunkedA   s    z#AsyncStreamReaderMixin.iter_chunkedc                 C   s
   t | jƒS )z¡Returns an asynchronous iterator that yields all the available
        data as soon as it is received

        Python-3.5 available for Python 3.5+ only
        )r   Úreadanyr#   r   r   r   Úiter_anyH   s    zAsyncStreamReaderMixin.iter_anyc                 C   s   t | ƒS )a  Returns an asynchronous iterator that yields chunks of data
        as they are received by the server. The yielded objects are tuples
        of (bytes, bool) as returned by the StreamReader.readchunk method.

        Python-3.5 available for Python 3.5+ only
        )r+   r#   r   r   r   Úiter_chunksP   s    z"AsyncStreamReaderMixin.iter_chunksN)r   r   r   r   r0   r$   Úintr:   r<   r+   r=   r   r   r   r   r2   =   s   r2   c                   @   s¨  e Zd ZdZdZdddœeeee ee	j
 ddœdd„Zedœd	d
„Zeeef dœdd„Zee dœdd„Zeddœdd„Zeg df ddœdd„Zddœdd„Zedœdd„Zedœdd„Zddœdd„Zeddœdd„Zd<eedd œd!d"„Zddœd#d$„Zddœd%d&„Zedd'œd(d)„Zedœd*d+„Z d=eed-œd.d/„Z!edœd0d1„Z"eeef dœd2d3„Z#eed-œd4d5„Z$d>eed-œd6d7„Z%eed-œd8d9„Z&eed-œd:d;„Z'dS )?r   a*  An enhancement of asyncio.StreamReader.

    Supports asynchronous iteration by line, chunk or as available::

        async for line in reader:
            ...
        async for chunk in reader.iter_chunked(1024):
            ...
        async for slice in reader.iter_any():
            ...

    r   N)ÚtimerÚloop)ÚprotocolÚlimitr?   r@   r   c                C   sv   || _ || _|d | _|d kr&t ¡ }|| _d| _d| _d | _t	 
¡ | _d| _d| _d | _d | _d | _|| _g | _d S )Né   r   F)Ú	_protocolÚ
_low_waterÚ_high_waterÚasyncioZget_event_loopÚ_loopÚ_sizeÚ_cursorÚ_http_chunk_splitsÚcollectionsÚdequeÚ_bufferÚ_buffer_offsetÚ_eofÚ_waiterÚ_eof_waiterÚ
_exceptionÚ_timerÚ_eof_callbacks)r    rA   rB   r?   r@   r   r   r   r!   j   s"    

zStreamReader.__init__r"   c                 C   sŠ   | j jg}| jr | d| j ¡ | jr0| d¡ | jdkrP| d| j| jf ¡ | jrf| d| j ¡ | jr|| d| j ¡ dd 	|¡ S )	Nz%d bytesÚeofi   zlow=%d high=%dzw=%rze=%rz<%s>ú )
Ú	__class__r   rI   ÚappendrP   rE   rF   rQ   rS   Újoin)r    Úinfor   r   r   Ú__repr__„   s    


zStreamReader.__repr__c                 C   s   | j | jfS r   )rE   rF   r#   r   r   r   Úget_read_buffer_limits’   s    z#StreamReader.get_read_buffer_limitsc                 C   s   | j S r   ©rS   r#   r   r   r   Ú	exception•   s    zStreamReader.exception©Úexcr   c                 C   sP   || _ | j ¡  | j}|d k	r.d | _t||ƒ | j}|d k	rLd | _t||ƒ d S r   )rS   rU   ÚclearrQ   r   rR   ©r    ra   Úwaiterr   r   r   r   ˜   s    

zStreamReader.set_exception©Úcallbackr   c                 C   sB   | j r2z
|ƒ  W q> tk
r.   t d¡ Y q>X n| j |¡ d S ©NúException in eof callback)rP   Ú	Exceptionr   r_   rU   rY   ©r    rf   r   r   r   Úon_eof¦   s    
zStreamReader.on_eofc              	   C   s†   d| _ | j}|d k	r$d | _t|d ƒ | j}|d k	rBd | _t|d ƒ | jD ].}z
|ƒ  W qH tk
rt   t d¡ Y qHX qH| j ¡  d S )NTrh   )	rP   rQ   r   rR   rU   ri   r   r_   rb   )r    rd   Úcbr   r   r   Úfeed_eof¯   s    



zStreamReader.feed_eofc                 C   s   | j S )z&Return True if  'feed_eof' was called.©rP   r#   r   r   r   Úis_eofÄ   s    zStreamReader.is_eofc                 C   s   | j o| j S )z=Return True if the buffer is empty and 'feed_eof' was called.©rP   rN   r#   r   r   r   Úat_eofÈ   s    zStreamReader.at_eofc                 Ã   sB   | j r
d S | jd kst‚| j ¡ | _z| jI d H  W 5 d | _X d S r   )rP   rR   ÚAssertionErrorrH   Úcreate_futurer#   r   r   r   Úwait_eofÌ   s    zStreamReader.wait_eof)Údatar   c                 C   sx   t jdtdd� |sdS | jr>| jd | jd… | jd< d| _|  jt|ƒ7  _|  jt|ƒ8  _| j |¡ d| _	dS )zDrollback reading some data from stream, inserting it to buffer head.zJunread_data() is deprecated and will be removed in future releases (#3260)rC   )Ú
stacklevelNr   )
ÚwarningsÚwarnÚDeprecationWarningrO   rN   rI   ÚlenrJ   Ú
appendleftÚ_eof_counter)r    ru   r   r   r   Úunread_data×   s    üzStreamReader.unread_data©ru   Úsizer   c                 C   s†   | j rtdƒ‚|sd S |  jt|ƒ7  _| j |¡ |  jt|ƒ7  _| j}|d k	rdd | _t|d ƒ | j| j	kr‚| j
js‚| j
 ¡  d S )Nzfeed_data after feed_eof)rP   rr   rI   rz   rN   rY   Útotal_bytesrQ   r   rF   rD   Ú_reading_pausedÚpause_reading©r    ru   r   rd   r   r   r   Ú	feed_dataë   s    
zStreamReader.feed_datac                 C   s"   | j d kr| jrtdƒ‚g | _ d S )Nz?Called begin_http_chunk_receiving whensome data was already fed)rK   r€   ÚRuntimeErrorr#   r   r   r   Úbegin_http_chunk_receivingý   s    
ÿz'StreamReader.begin_http_chunk_receivingc                 C   sd   | j d krtdƒ‚| j r"| j d nd}| j|kr4d S | j  | j¡ | j}|d k	r`d | _t|d ƒ d S )NzFCalled end_chunk_receiving without calling begin_chunk_receiving firstéÿÿÿÿr   )rK   r…   r€   rY   rQ   r   )r    Úposrd   r   r   r   Úend_http_chunk_receiving  s    
ÿ

z%StreamReader.end_http_chunk_receiving)Ú	func_namer   c              	   Ã   sf   | j d k	rtd| ƒ‚| j ¡  }| _ z2| jrL| j� |I d H  W 5 Q R X n
|I d H  W 5 d | _ X d S )NzH%s() called while another coroutine is already waiting for incoming data)rQ   r…   rH   rs   rT   )r    rŠ   rd   r   r   r   Ú_wait#  s    
ÿÿzStreamReader._waitc                 Ã   s¶   | j d k	r| j ‚g }d}d}|r¬| jrŽ|rŽ| j}| jd  d|¡d }|  |rV|| nd¡}| |¡ |t|ƒ7 }|rzd}|| jkr tdƒ‚q | j	r–q¬|r|  
d¡I d H  qd	 |¡S )
Nr   Tó   
r	   r‡   FzLine is too longr3   r&   )rS   rN   rO   ÚfindÚ_read_nowait_chunkrY   rz   rF   Ú
ValueErrorrP   r‹   rZ   )r    ÚlineZ	line_sizeZ
not_enoughÚoffsetZicharru   r   r   r   r3   8  s*    




zStreamReader.readliner‡   r4   c                 Ã   s¬   | j d k	r| j ‚| jrF| jsFt| ddƒd | _| jdkrFtjddd� |sNdS |dk r„g }|  ¡ I d H }|snqz| |¡ qZd 	|¡S | js¢| js¢|  
d	¡I d H  q„|  |¡S )
Nr|   r   r	   é   zEMultiple access to StreamReader in eof state, might be infinite loop.T)Ú
stack_infor&   r6   )rS   rP   rN   Úgetattrr|   r   Úwarningr;   rY   rZ   r‹   Ú_read_nowait)r    r5   ÚblocksÚblockr   r   r   r6   V  s*    

ý
zStreamReader.readc                 Ã   s8   | j d k	r| j ‚| js.| js.|  d¡I d H  q|  d¡S )Nr;   r‡   )rS   rN   rP   r‹   r–   r#   r   r   r   r;   €  s
    
zStreamReader.readanyc                 Ã   sŽ   | j dk	r| j ‚| jrZ| j d¡}|| jkr0dS || jkrN|  || j ¡dfS t d¡ q| jrn|  d¡dfS | j	rxdS |  
d	¡I dH  q dS )
zþReturns a tuple of (data, end_of_http_chunk). When chunked transfer
        encoding is used, end_of_http_chunk is a boolean indicating if the end
        of the data corresponds to the end of a HTTP chunk , otherwise it is
        always False.
        Nr   ©r&   TTzESkipping HTTP chunk end due to data consumption beyond chunk boundaryr‡   Fr.   r/   )rS   rK   ÚpoprJ   r–   r   r•   rN   rŽ   rP   r‹   )r    rˆ   r   r   r   r/   Œ  s     


ÿzStreamReader.readchunkc                 Ã   sp   | j d k	r| j ‚g }|dkrf|  |¡I d H }|sNd |¡}t |t|ƒ| ¡‚| |¡ |t|ƒ8 }qd |¡S )Nr   r&   )rS   r6   rZ   rG   ÚIncompleteReadErrorrz   rY   )r    r5   r—   r˜   Úpartialr   r   r   Úreadexactly¬  s    


zStreamReader.readexactlyc                 C   s2   | j d k	r| j ‚| jr(| j ¡ s(tdƒ‚|  |¡S )Nz9Called while some coroutine is waiting for incoming data.)rS   rQ   Údoner…   r–   r9   r   r   r   Úread_nowait»  s    
ÿzStreamReader.read_nowaitc                 C   sÞ   | j d }| j}|dkrHt|ƒ| |krH|||| … }|  j|7  _n,|rj| j  ¡  ||d … }d| _n
| j  ¡ }|  jt|ƒ8  _|  jt|ƒ7  _| j}|r¼|d | jk r¼| d¡ qž| j| jk rÚ| j	j
rÚ| j	 ¡  |S )Nr   r‡   )rN   rO   rz   ÚpopleftrI   rJ   rK   rš   rE   rD   r�   Úresume_reading)r    r5   Zfirst_bufferr‘   ru   Zchunk_splitsr   r   r   rŽ   Ê  s$    



zStreamReader._read_nowait_chunkc                 C   sP   g }| j r>|  |¡}| |¡ |dkr|t|ƒ8 }|dkrq>q|rLd |¡S dS )z8 Read not more than n bytes, or whole buffer if n == -1 r‡   r   r&   )rN   rŽ   rY   rz   rZ   )r    r5   ÚchunksÚchunkr   r   r   r–   å  s    

zStreamReader._read_nowait)r   )r‡   )r‡   )(r   r   r   r   r€   r
   r>   r   r   rG   ÚAbstractEventLoopr!   Ústrr\   r   r]   ÚBaseExceptionr_   r   r   rk   rm   r1   ro   rq   rt   r0   r}   r„   r†   r‰   r‹   r3   r6   r;   r/   r�   rŸ   rŽ   r–   r   r   r   r   r   Z   sB   úù	* r   c                   @   sô   e Zd Zee dœdd„Zeddœdd„Zeg df ddœd	d
„Zddœdd„Z	e
dœdd„Ze
dœdd„Zddœdd„Zd%eeddœdd„Zedœdd„Zd&eedœdd„Zedœdd„Zeee
f dœdd „Zeedœd!d"„Zedœd#d$„ZdS )'ÚEmptyStreamReaderr"   c                 C   s   d S r   r   r#   r   r   r   r_   õ  s    zEmptyStreamReader.exceptionNr`   c                 C   s   d S r   r   )r    ra   r   r   r   r   ø  s    zEmptyStreamReader.set_exceptionre   c                 C   s.   z
|ƒ  W n t k
r(   t d¡ Y nX d S rg   )ri   r   r_   rj   r   r   r   rk   û  s    
zEmptyStreamReader.on_eofc                 C   s   d S r   r   r#   r   r   r   rm     s    zEmptyStreamReader.feed_eofc                 C   s   dS ©NTr   r#   r   r   r   ro     s    zEmptyStreamReader.is_eofc                 C   s   dS r¨   r   r#   r   r   r   rq     s    zEmptyStreamReader.at_eofc                 Ã   s   d S r   r   r#   r   r   r   rt   
  s    zEmptyStreamReader.wait_eofr   )ru   r5   r   c                 C   s   d S r   r   )r    ru   r5   r   r   r   r„     s    zEmptyStreamReader.feed_datac                 Ã   s   dS r%   r   r#   r   r   r   r3     s    zEmptyStreamReader.readliner‡   r4   c                 Ã   s   dS r%   r   r9   r   r   r   r6     s    zEmptyStreamReader.readc                 Ã   s   dS r%   r   r#   r   r   r   r;     s    zEmptyStreamReader.readanyc                 Ã   s   dS )Nr™   r   r#   r   r   r   r/     s    zEmptyStreamReader.readchunkc                 Ã   s   t  d|¡‚d S r%   )rG   r›   r9   r   r   r   r�     s    zEmptyStreamReader.readexactlyc                 C   s   dS r%   r   r#   r   r   r   rŸ     s    zEmptyStreamReader.read_nowait)r   )r‡   )r   r   r   r   r¦   r_   r   r   rk   rm   r1   ro   rq   rt   r0   r>   r„   r3   r6   r;   r   r/   r�   rŸ   r   r   r   r   r§   ô  s   r§   c                   @   s°   e Zd ZdZejddœdd„Zedœdd„Ze	dœd	d
„Z
e	dœdd„Zee dœdd„Zeddœdd„Zdeeddœdd„Zddœdd„Zedœdd„Zee dœdd„ZdS )r   z>DataQueue is a general-purpose blocking queue with one reader.N)r@   r   c                 C   s,   || _ d| _d | _d | _d| _t ¡ | _d S )NFr   )rH   rP   rQ   rS   rI   rL   rM   rN   )r    r@   r   r   r   r!   )  s    zDataQueue.__init__r"   c                 C   s
   t | jƒS r   )rz   rN   r#   r   r   r   Ú__len__1  s    zDataQueue.__len__c                 C   s   | j S r   rn   r#   r   r   r   ro   4  s    zDataQueue.is_eofc                 C   s   | j o| j S r   rp   r#   r   r   r   rq   7  s    zDataQueue.at_eofc                 C   s   | j S r   r^   r#   r   r   r   r_   :  s    zDataQueue.exceptionr`   c                 C   s.   d| _ || _| j}|d k	r*d | _t||ƒ d S r¨   )rP   rS   rQ   r   rc   r   r   r   r   =  s    zDataQueue.set_exceptionr   r~   c                 C   s@   |  j |7  _ | j ||f¡ | j}|d k	r<d | _t|d ƒ d S r   )rI   rN   rY   rQ   r   rƒ   r   r   r   r„   F  s    zDataQueue.feed_datac                 C   s(   d| _ | j}|d k	r$d | _t|d ƒ d S r¨   )rP   rQ   r   )r    rd   r   r   r   rm   O  s
    zDataQueue.feed_eofc              	   Ã   s˜   | j sX| jsX| jrt‚| j ¡ | _z| jI d H  W n$ tjtjfk
rV   d | _‚ Y nX | j r~| j  	¡ \}}|  j
|8  _
|S | jd k	r�| j‚nt‚d S r   )rN   rP   rQ   rr   rH   rs   rG   ÚCancelledErrorÚTimeoutErrorr    rI   rS   r   ©r    ru   r   r   r   r   r6   W  s    

zDataQueue.readc                 C   s
   t | jƒS r   )r   r6   r#   r   r   r   r$   k  s    zDataQueue.__aiter__)r   )r   r   r   r   rG   r¤   r!   r>   r©   r1   ro   rq   r   r¦   r_   r   r   r„   rm   r6   r   r$   r   r   r   r   r   &  s   		r   c                       sX   e Zd ZdZeeejddœ‡ fdd„Zde	eddœ‡ fdd	„Z
e	d
œ‡ fdd„Z‡  ZS )r   zgFlowControlDataQueue resumes and pauses an underlying stream.

    It is a destination for parsed data.N)rA   rB   r@   r   c                   s"   t ƒ j|d� || _|d | _d S )N)r@   rC   )Úsuperr!   rD   Ú_limit)r    rA   rB   r@   ©rX   r   r   r!   t  s    zFlowControlDataQueue.__init__r   r~   c                    s0   t ƒ  ||¡ | j| jkr,| jjs,| j ¡  d S r   )r­   r„   rI   r®   rD   r�   r‚   r¬   r¯   r   r   r„   |  s    zFlowControlDataQueue.feed_datar"   c                 ƒ   s:   ztƒ  ¡ I d H W ¢S | j | jk r4| jjr4| j ¡  X d S r   )rI   r®   rD   r�   r¡   r­   r6   r#   r¯   r   r   r6   ‚  s    zFlowControlDataQueue.read)r   )r   r   r   r   r
   r>   rG   r¤   r!   r   r„   r6   Ú__classcell__r   r   r¯   r   r   o  s     þr   )"rG   rL   rw   Útypingr   r   r   r   r   r   r   Zbase_protocolr
   Zhelpersr   r   r   Úlogr   r   ÚImportErrorZtyping_extensionsÚ__all__r   ri   r   r   r+   r2   r   r§   r   r   r   r   r   r   r   Ú<module>   s0   $   /I