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*|jr J‚d|_|j ¡r&t d|¡dS)NTz%r pauses writing)r8r7Ú get_debugrÚdebug©r;rrrÚ pause_writingšs  zFlowControlMixin.pause_writingcCsP|js J‚d|_|j ¡r&t d|¡|j}|durLd|_| ¡sL| 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ÃsP|jrtdƒ‚|jsdS|j}|dus2| ¡s2J‚|j ¡}||_|IdHdS)NzConnection lost)r:ÚConnectionResetErrorr8r9Ú cancelledr7Ú 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@rErJrNrRrrrrr5‡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_sslr7rMÚ_closed)r;Z stream_readerr1r©Ú __class__rrr<Ös  zStreamReaderProtocol.__init__cCs|jdurdS| ¡Sr)rZr?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]r[r7Zcall_exception_handlerÚabortr_reÚ set_transportÚget_extra_inforar`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) reÚfeed_eofrGrbrArBrWrJrZr^r_)r;rIr)rcrrrJ s     z$StreamReaderProtocol.connection_lostcCs|j}|dur| |¡dSr)reÚ feed_data)r;Údatar)rrrÚ data_receivedsz"StreamReaderProtocol.data_receivedcCs$|j}|dur| ¡|jr dSdS)NFT)rermra)r;r)rrrÚ eof_received s z!StreamReaderProtocol.eof_receivedcCs|jSr)rbrPrrrrR+sz&StreamReaderProtocol._get_close_waitercCs"|j}| ¡r| ¡s| ¡dSr)rbrArLÚ exception)r;ÚclosedrrrÚ__del__.szStreamReaderProtocol.__del__)NN)rSrTrUrVr[r<ÚpropertyrerlrJrprqrRrtÚ __classcell__rrrcrrË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. cCsJ||_||_|dus"t|tƒs"J‚||_||_|j ¡|_|j d¡dSr) r_Ú _protocolÚ isinstancerÚ_readerr7rMZ _complete_futrB)r;r*rr)rrrrr<@s zStreamWriter.__init__cCs@|jjd|j›�g}|jdur0| d|j›�¡d d |¡¡S)Nú transport=zreader=ú<{}>ú )rdrSr_ryÚappendÚformatÚjoin©r;ÚinforrrÚ__repr__Js zStreamWriter.__repr__cCs|jSr©r_r?rrrr*PszStreamWriter.transportcCs|j |¡dSr)r_Úwrite©r;rorrrr„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)rwrRr?rrrÚ wait_closedfszStreamWriter.wait_closedNcCs|j ||¡Sr)r_ri)r;ÚnameÚdefaultrrrriiszStreamWriter.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)ryrrr_rŠrrwrN)r;rIrrrÚdrainls   zStreamWriter.drain)N)rSrTrUrVr<r‚rur*r„r†r‡rˆr‰rŠr‹rirŽ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Ú _getframer[)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=rzZpausedr{r|) r’r}Úlenr“r�Ú_DEFAULT_LIMITr”r•r_r8r~rr€rrrr‚ s    zStreamReader.__repr__cCs|jSr)r•r?rrrrr²szStreamReader.exceptioncCs0||_|j}|dur,d|_| ¡s,| |¡dSr)r•r”rLrGrHrrrrGµs zStreamReader.set_exceptioncCs*|j}|dur&d|_| ¡s&| d¡dS)z1Wakeup read*() functions waiting for data or EOF.N)r”rLrBrCrrrÚ_wakeup_waiter¾s zStreamReader._wakeup_waitercCs|jdusJdƒ‚||_dS)NzTransport already setrƒ)r;r*rrrrhÆ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?rrrrmÏ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_eofcCs€|jrJdƒ‚|sdS|j |¡| ¡|jdur||js|t|jƒd|jkr|z|j ¡Wnt ytd|_Yn0d|_dS)Nzfeed_data after feed_eofrT) r“r’Úextendrœr_r8ršr�Z pause_readingrOr…rrrrn×s  ÿþ  zStreamReader.feed_datacÃsl|jdurt|›d�ƒ‚|jr&Jdƒ‚|jrs>        ÿ !ÿ ' ÿ ÿ DkP