grthtrhthjhtyjytjytkergtrhtrjytjerhrfh4:24 29/09/2026a ´i hã@s&dZddlZddlZddlZddlZeedƒr6ed7ZddlmZddlmZddlm Z dd lm Z dd lm Z dd l m Z dd lmZd Zddedœdd„Zd dedœdd„Zeedƒràd!dedœdd„Zd"dedœdd„ZGdd„de jƒZGdd„dee jƒZGdd„dƒZGdd„dƒZdS)#)Ú StreamReaderÚ StreamWriterÚStreamReaderProtocolÚopen_connectionÚ start_serveréNÚAF_UNIX)Úopen_unix_connectionÚstart_unix_serveré)Ú coroutines)Úevents)Ú exceptions)Úformat_helpers)Ú protocols)Úlogger)Úsleepi)ÚloopÚlimitc ‹sx|durt ¡}ntjdtdd�t||d�}t||d�‰|j‡fdd„||fi|¤ŽIdH\}}t|ˆ||ƒ}||fS) aÂA wrapper for create_connection() returning a (reader, writer) pair. The reader returned is a StreamReader instance; the writer is a StreamWriter instance. The arguments are all the usual arguments to create_connection() except protocol_factory; most common are positional host and port, with various optional keyword arguments following. Additional optional keyword arguments are loop (to set the event loop instance to use) and limit (to set the buffer limit passed to the StreamReader). (If you want to customize the StreamReader and/or StreamReaderProtocol classes, just copy the code -- there's really nothing special here except some convenience.) Nú[The loop argument is deprecated since Python 3.8, and scheduled for removal in Python 3.10.é©Ú stacklevel©rr©rcsˆS©N©r©Úprotocolrú'/usr/lib64/python3.9/asyncio/streams.pyÚ5óz!open_connection..) r Úget_event_loopÚwarningsÚwarnÚDeprecationWarningrrÚcreate_connectionr) ÚhostÚportrrÚkwdsÚreaderÚ transportÚ_Úwriterrrrrs þ  ÿÿrc‹sNˆdurt ¡‰ntjdtdd�‡‡‡fdd„}ˆj|||fi|¤ŽIdHS)aÉStart a socket server, call back for each client connected. The first parameter, `client_connected_cb`, takes two parameters: client_reader, client_writer. client_reader is a StreamReader object, while client_writer is a StreamWriter object. This parameter can either be a plain callback function or a coroutine; if it is a coroutine, it will be automatically converted into a Task. The rest of the arguments are all the usual arguments to loop.create_server() except protocol_factory; most common are positional host and port, with various optional keyword arguments following. The return value is the same as loop.create_server(). Additional optional keyword arguments are loop (to set the event loop instance to use) and limit (to set the buffer limit passed to the StreamReader). The return value is the same as loop.create_server(), i.e. a Server object which can be used to stop the service. Nrrrcstˆˆd�}t|ˆˆd�}|S©Nrr©rr©r)r©Úclient_connected_cbrrrrÚfactoryXs  ÿzstart_server..factory)r r!r"r#r$Ú create_server)r1r&r'rrr(r2rr0rr:s þrc‹sv|durt ¡}ntjdtdd�t||d�}t||d�‰|j‡fdd„|fi|¤ŽIdH\}}t|ˆ||ƒ}||fS) z@Similar to `open_connection` but works with UNIX Domain Sockets.NrrrrrcsˆSrrrrrrrpr z&open_unix_connection..) r r!r"r#r$rrZcreate_unix_connectionr)Úpathrrr(r)r*r+r,rrrrds þ   ÿÿrc‹sLˆdurt ¡‰ntjdtdd�‡‡‡fdd„}ˆj||fi|¤ŽIdHS)z=Similar to `start_server` but works with UNIX Domain Sockets.Nrrrcstˆˆd�}t|ˆˆd�}|Sr-r.r/r0rrr2~s  ÿz"start_unix_server..factory)r r!r"r#r$Zcreate_unix_server)r1r4rrr(r2rr0rr ts þr c@sBeZdZdZddd„Zdd„Zdd„Zd d „Zd d „Zd d„Z dS)ÚFlowControlMixina)Reusable flow control logic for StreamWriter.drain(). This implements the protocol methods pause_writing(), resume_writing() and connection_lost(). If the subclass overrides these it must call the super methods. StreamWriter.drain() must wait for _drain_helper() coroutine. NcCs0|durt ¡|_n||_d|_d|_d|_dS©NF)r r!Ú_loopÚ_pausedÚ _drain_waiterÚ_connection_lost)ÚselfrrrrÚ__init__‘s  zFlowControlMixin.__init__cCs d|_|j ¡rt d|¡dS)NTz%r pauses writing)r8r7Ú get_debugrÚdebug©r;rrrÚ pause_writingšs zFlowControlMixin.pause_writingcCsFd|_|j ¡rt d|¡|j}|durBd|_| ¡sB| d¡dS)NFz%r resumes writing)r8r7r=rr>r9ÚdoneÚ set_result©r;ÚwaiterrrrÚresume_writing s  zFlowControlMixin.resume_writingcCsVd|_|jsdS|j}|dur"dSd|_| ¡r4dS|durH| d¡n | |¡dS©NT)r:r8r9rArBÚ set_exception©r;ÚexcrDrrrÚconnection_lost¬s z FlowControlMixin.connection_lostcÃs<|jrtdƒ‚|jsdS|j}|j ¡}||_|IdHdS)NzConnection lost)r:ÚConnectionResetErrorr8r9r7Ú create_futurerCrrrÚ _drain_helper¼s zFlowControlMixin._drain_helpercCst‚dSr)ÚNotImplementedError©r;ÚstreamrrrÚ_get_close_waiterÇsz"FlowControlMixin._get_close_waiter)N) Ú__name__Ú __module__Ú __qualname__Ú__doc__r<r@rErJrMrQrrrrr5‡s   r5csfeZdZdZdZd‡fdd„ Zedd„ƒZdd„Z‡fd d „Z d d „Z d d„Z dd„Z dd„Z ‡ZS)ra=Helper class to adapt between Protocol and StreamReader. (This is a helper class instead of making StreamReader itself a Protocol subclass, because the StreamReader has other potential uses, and to prevent the user of the StreamReader to accidentally call inappropriate methods of the protocol.) Ncsntƒj|d�|dur,t |¡|_|j|_nd|_|dur@||_d|_d|_d|_ ||_ d|_ |j   ¡|_dS)NrF)Úsuperr<ÚweakrefÚrefÚ_stream_reader_wrÚ_source_tracebackÚ_strong_readerÚ_reject_connectionÚ_stream_writerÚ _transportÚ_client_connected_cbÚ _over_sslr7rLÚ_closed)r;Z stream_readerr1r©Ú __class__rrr<Ös  zStreamReaderProtocol.__init__cCs|jdurdS| ¡Sr)rYr?rrrÚ_stream_readerés z#StreamReaderProtocol._stream_readercCs®|jr6ddi}|jr|j|d<|j |¡| ¡dS||_|j}|durT| |¡| d¡du|_ |j durªt ||||jƒ|_ |  ||j ¡}t  |¡r¤|j |¡d|_dS)NÚmessagezpAn open stream was garbage collected prior to establishing network connection; call "stream.close()" explicitly.Zsource_tracebackZ sslcontext)r\rZr7Zcall_exception_handlerÚabortr^rdÚ set_transportÚget_extra_infor`r_rr]r Z iscoroutineZ create_taskr[)r;r*Úcontextr)ÚresrrrÚconnection_madeïs0ÿ    þÿ  z$StreamReaderProtocol.connection_madecsx|j}|dur*|dur | ¡n | |¡|j ¡sV|durJ|j d¡n |j |¡tƒ |¡d|_d|_ d|_ dSr) rdÚfeed_eofrGrarArBrVrJrYr]r^)r;rIr)rbrrrJ s     z$StreamReaderProtocol.connection_lostcCs|j}|dur| |¡dSr)rdÚ feed_data)r;Údatar)rrrÚ data_receivedsz"StreamReaderProtocol.data_receivedcCs$|j}|dur| ¡|jr dSdS)NFT)rdrlr`)r;r)rrrÚ eof_received s z!StreamReaderProtocol.eof_receivedcCs|jSr)rarOrrrrQ+sz&StreamReaderProtocol._get_close_waitercCs"|j}| ¡r| ¡s| ¡dSr)rarAÚ cancelledÚ exception)r;ÚclosedrrrÚ__del__.szStreamReaderProtocol.__del__)NN)rRrSrTrUrZr<ÚpropertyrdrkrJrorprQrtÚ __classcell__rrrbrrËs   rc@sveZdZdZdd„Zdd„Zedd„ƒZdd „Zd d „Z d d „Z dd„Z dd„Z dd„Z dd„Zddd„Zdd„ZdS)ra'Wraps a Transport. This exposes write(), writelines(), [can_]write_eof(), get_extra_info() and close(). It adds drain() which returns an optional Future on which you can wait for flow control. It also adds a transport property which references the Transport directly. cCs4||_||_||_||_|j ¡|_|j d¡dSr)r^Ú _protocolÚ_readerr7rLZ _complete_futrB)r;r*rr)rrrrr<@s  zStreamWriter.__init__cCs@|jjd|j›�g}|jdur0| d|j›�¡d d |¡¡S)Nú transport=zreader=ú<{}>ú )rcrRr^rxÚappendÚformatÚjoin©r;ÚinforrrÚ__repr__Js zStreamWriter.__repr__cCs|jSr©r^r?rrrr*PszStreamWriter.transportcCs|j |¡dSr)r^Úwrite©r;rnrrrrƒTszStreamWriter.writecCs|j |¡dSr)r^Ú writelinesr„rrrr…WszStreamWriter.writelinescCs |j ¡Sr)r^Ú write_eofr?rrrr†ZszStreamWriter.write_eofcCs |j ¡Sr)r^Ú can_write_eofr?rrrr‡]szStreamWriter.can_write_eofcCs |j ¡Sr)r^Úcloser?rrrrˆ`szStreamWriter.closecCs |j ¡Sr)r^Ú is_closingr?rrrr‰cszStreamWriter.is_closingcÃs|j |¡IdHdSr)rwrQr?rrrÚ wait_closedfszStreamWriter.wait_closedNcCs|j ||¡Sr)r^rh)r;ÚnameÚdefaultrrrrhiszStreamWriter.get_extra_infocÃsL|jdur |j ¡}|dur |‚|j ¡r8tdƒIdH|j ¡IdHdS)zyFlush the write buffer. The intended use is to write w.write(data) await w.drain() Nr)rxrrr^r‰rrwrM)r;rIrrrÚdrainls   zStreamWriter.drain)N)rRrSrTrUr<r�rur*rƒr…r†r‡rˆr‰rŠrhr�rrrrr6s    rc@s¢eZdZdZedfdd„Zdd„Zdd„Zdd „Zd d „Z d d „Z dd„Z dd„Z dd„Z dd„Zdd„Zdd„Zd&dd„Zd'dd„Zd d!„Zd"d#„Zd$d%„ZdS)(rNcCsv|dkrtdƒ‚||_|dur*t ¡|_n||_tƒ|_d|_d|_d|_ d|_ d|_ |j  ¡rrt  t d¡¡|_dS)NrzLimit cannot be <= 0Fr )Ú ValueErrorÚ_limitr r!r7Ú bytearrayÚ_bufferÚ_eofÚ_waiterÚ _exceptionr^r8r=rÚ extract_stackÚsysÚ _getframerZ)r;rrrrrr<Šs   ÿzStreamReader.__init__cCs¶dg}|jr"| t|jƒ›d�¡|jr2| d¡|jtkrN| d|j›�¡|jrf| d|j›�¡|jr~| d|j›�¡|jr–| d|j›�¡|j r¦| d¡d   d   |¡¡S) Nrz bytesÚeofzlimit=zwaiter=z exception=ryZpausedrzr{) r‘r|Úlenr’r�Ú_DEFAULT_LIMITr“r”r^r8r}r~rrrrr� s    zStreamReader.__repr__cCs|jSr)r”r?rrrrr²szStreamReader.exceptioncCs0||_|j}|dur,d|_| ¡s,| |¡dSr)r”r“rqrGrHrrrrGµs zStreamReader.set_exceptioncCs*|j}|dur&d|_| ¡s&| d¡dS)z1Wakeup read*() functions waiting for data or EOF.N)r“rqrBrCrrrÚ_wakeup_waiter¾s zStreamReader._wakeup_waitercCs ||_dSrr‚)r;r*rrrrgÆszStreamReader.set_transportcCs*|jr&t|jƒ|jkr&d|_|j ¡dSr6)r8r™r‘r�r^Úresume_readingr?rrrÚ_maybe_resume_transportÊsz$StreamReader._maybe_resume_transportcCsd|_| ¡dSrF)r’r›r?rrrrlÏszStreamReader.feed_eofcCs|jo |j S)z=Return True if the buffer is empty and 'feed_eof' was called.)r’r‘r?rrrÚat_eofÓszStreamReader.at_eofcCsr|sdS|j |¡| ¡|jdurn|jsnt|jƒd|jkrnz|j ¡Wntyfd|_Yn0d|_dS)NrT) r‘Úextendr›r^r8r™r�Z pause_readingrNr„rrrrm×s  ÿþ  zStreamReader.feed_datacÃs^|jdurt|›d�ƒ‚|jr.d|_|j ¡|j ¡|_z|jIdHWd|_nd|_0dS)zpWait until feed_data() or feed_eof() is called. If stream was paused, automatically resume it. NzF() called while another coroutine is already waiting for incoming dataF)r“Ú RuntimeErrorr8r^rœr7rL)r;Ú func_namerrrÚ_wait_for_dataís ÿ  zStreamReader._wait_for_datac Ãsºd}t|ƒ}z| |¡IdH}Wn”tjyL}z|jWYd}~Sd}~0tjy´}zP|j ||j¡r€|jd|j|…=n |j  ¡|  ¡t |j dƒ‚WYd}~n d}~00|S)aÂRead chunk of data from the stream until newline (b' ') is found. On success, return chunk that ends with newline. If only partial line can be read due to EOF, return incomplete line without terminating newline. When EOF was reached while no bytes read, empty bytes object is returned. If limit is reached, ValueError will be raised. In that case, if newline was found, complete line including newline will be removed from internal buffer. Else, internal buffer will be cleared. Limit is compared against part of the line without newline. If stream was paused, this function will automatically resume it if needed. ó Nr) r™Ú readuntilr ÚIncompleteReadErrorÚpartialÚLimitOverrunErrorr‘Ú startswithÚconsumedÚclearr�rŽÚargs)r;ÚsepÚseplenÚlineÚerrrÚreadline s $zStreamReader.readliner£cÃsüt|ƒ}|dkrtdƒ‚|jdur(|j‚d}t|jƒ}|||kr||j ||¡}|dkrZq´|d|}||jkr|t d|¡‚|jr¢t |jƒ}|j  ¡t  |d¡‚|  d¡IdHq,||jkrÊt d|¡‚|jd||…}|jd||…=|  ¡t |ƒS) aVRead data from the stream until ``separator`` is found. On success, the data and separator will be removed from the internal buffer (consumed). Returned data will include the separator at the end. Configured stream limit is used to check result. Limit sets the maximal length of data that can be returned, not counting the separator. If an EOF occurs and the complete separator is still not found, an IncompleteReadError exception will be raised, and the internal buffer will be reset. The IncompleteReadError.partial attribute may contain the separator partially. If the data cannot be read because of over limit, a LimitOverrunError exception will be raised, and the data will be left in the internal buffer, so it can be read again. rz,Separator should be at least one-byte stringNéÿÿÿÿr z2Separator is not found, and chunk exceed the limitr¤z2Separator is found, but chunk is longer than limit)r™rŽr”r‘Úfindr�r r§r’Úbytesrªr¥r¢r�)r;Ú separatorr­ÚoffsetÚbuflenZisepÚchunkrrrr¤(s<     þ    ÿzStreamReader.readuntilr±cÃsœ|jdur|j‚|dkrdS|dkrVg}| |j¡IdH}|s@qL| |¡q(d |¡S|jsr|jsr| d¡IdHt|jd|…ƒ}|jd|…=|  ¡|S)aÚRead up to `n` bytes from the stream. If n is not provided, or set to -1, read until EOF and return all read bytes. If the EOF was received and the internal buffer is empty, return an empty bytes object. If n is zero, return empty bytes object immediately. If n is positive, this function try to read `n` bytes, and may return less or equal bytes than requested, but at least one byte. If EOF was received before any byte is read, this function returns empty byte object. Returned value is not limited with limit, configured at stream creation. If stream was paused, this function will automatically resume it if needed. Nrr Úread) r”r¸r�r|r~r‘r’r¢r³r�)r;ÚnZblocksÚblockrnrrrr¸ƒs"     zStreamReader.readcÃsÀ|dkrtdƒ‚|jdur |j‚|dkr,dSt|jƒ|krr|jr`t|jƒ}|j ¡t ||¡‚|  d¡IdHq,t|jƒ|kr–t|jƒ}|j ¡nt|jd|…ƒ}|jd|…=|  ¡|S)aÏRead exactly `n` bytes. Raise an IncompleteReadError if EOF is reached before `n` bytes can be read. The IncompleteReadError.partial attribute of the exception will contain the partial read bytes. if n is zero, return empty bytes object. Returned value is not limited with limit, configured at stream creation. If stream was paused, this function will automatically resume it if needed. rz*readexactly size can not be less than zeroNr Ú readexactly) rŽr”r™r‘r’r³rªr r¥r¢r�)r;r¹Z incompleternrrrr»µs&       zStreamReader.readexactlycCs|Srrr?rrrÚ __aiter__ÞszStreamReader.__aiter__cÃs| ¡IdH}|dkrt‚|S)Nr )r°ÚStopAsyncIteration)r;ÚvalrrrÚ __anext__ászStreamReader.__anext__)r£)r±)rRrSrTrZršr<r�rrrGr›rgr�rlržrmr¢r°r¤r¸r»r¼r¿rrrrr†s$  [ 2)r)NN)NN)N)N)Ú__all__Úsocketr–r"rWÚhasattrÚr r r rrÚlogrZtasksrršrrrr ÚProtocolr5rrrrrrrÚs>        ÿ !ÿ ' ÿ ÿ DkP