3 \˜_ã@sLdZdddddddgZdd lZeed ƒr6ejd d gƒd dlmZd dlmZd dlmZd dlm Z d dlm Z d dl m Z d"Z Gdd„deƒZGdd„deƒZe d#d e dœdd„ƒZe d$d e dœdd„ƒZeed ƒ�re d%d e dœdd „ƒZe d&d e dœdd „ƒZGdd„de jƒZGdd„dee jƒZGd d„dƒZGd!d„dƒZd S)'zStream-related things.Ú StreamReaderÚ StreamWriterÚStreamReaderProtocolÚopen_connectionÚ start_serverÚIncompleteReadErrorÚLimitOverrunErroréNZAF_UNIXÚopen_unix_connectionÚstart_unix_serveré)Ú coroutines)Úcompat)Úevents)Ú protocols)Ú coroutine)Úloggeréécs(eZdZdZ‡fdd„Zdd„Z‡ZS)rz· Incomplete read error. Attributes: - partial: read bytes string before the end of stream was reached - expected: total number of expected bytes (or None if unknown) cs(tƒjdt|ƒ|fƒ||_||_dS)Nz-%d bytes read on a total of %r expected bytes)ÚsuperÚ__init__ÚlenÚpartialÚexpected)Úselfrr)Ú __class__©ú'/usr/lib64/python3.6/asyncio/streams.pyr szIncompleteReadError.__init__cCst|ƒ|j|jffS)N)Útyperr)rrrrÚ __reduce__&szIncompleteReadError.__reduce__)Ú__name__Ú __module__Ú __qualname__Ú__doc__rrÚ __classcell__rr)rrrs cs(eZdZdZ‡fdd„Zdd„Z‡ZS)rzƒReached the buffer limit while looking for a separator. Attributes: - consumed: total number of to be consumed bytes. cstƒj|ƒ||_dS)N)rrÚconsumed)rÚmessager$)rrrr0s zLimitOverrunError.__init__cCst|ƒ|jd|jffS)Nr)rÚargsr$)rrrrr4szLimitOverrunError.__reduce__)rr r!r"rrr#rr)rrr*s )ÚloopÚlimitc +sb|dkrtjƒ}t||d�}t||d�‰|j‡fdd„||f|ŽEdH\}}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)r(r')r'csˆS)Nrr)ÚprotocolrrÚQsz!open_connection..)rÚget_event_looprrZcreate_connectionr) ÚhostÚportr'r(ÚkwdsÚreaderÚ transportÚ_Úwriterr)r)rr8s   c+s8ˆdkrtjƒ‰‡‡‡fdd„}ˆj|||f|ŽEdHS)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. Ncstˆˆd�}t|ˆˆd�}|S)N)r(r')r')rr)r/r))Úclient_connected_cbr(r'rrÚfactoryqs zstart_server..factory)rr+Z create_server)r3r,r-r'r(r.r4r)r3r(r'rrVsc+s`|dkrtjƒ}t||d�}t||d�‰|j‡fdd„|f|ŽEdH\}}t|ˆ||ƒ}||fS)z@Similar to `open_connection` but works with UNIX Domain Sockets.N)r(r')r'csˆS)Nrr)r)rrr*†sz&open_unix_connection..)rr+rrZcreate_unix_connectionr)Úpathr'r(r.r/r0r1r2r)r)rr }s  c+s6ˆdkrtjƒ‰‡‡‡fdd„}ˆj||f|ŽEdHS)z=Similar to `start_server` but works with UNIX Domain Sockets.Ncstˆˆd�}t|ˆˆd�}|S)N)r(r')r')rr)r/r))r3r(r'rrr4‘s z"start_unix_server..factory)rr+Zcreate_unix_server)r3r5r'r(r.r4r)r3r(r'rr Šsc@s>eZdZdZd dd„Zdd„Zdd„Zd d „Zed d „ƒZ dS)ÚFlowControlMixina)Reusable flow control logic for StreamWriter.drain(). This implements the protocol methods pause_writing(), resume_reading() and connection_lost(). If the subclass overrides these it must call the super methods. StreamWriter.drain() must wait for _drain_helper() coroutine. NcCs0|dkrtjƒ|_n||_d|_d|_d|_dS)NF)rr+Ú_loopÚ_pausedÚ _drain_waiterÚ_connection_lost)rr'rrrr¤s  zFlowControlMixin.__init__cCs,|j s t‚d|_|jjƒr(tjd|ƒdS)NTz%r pauses writing)r8ÚAssertionErrorr7Ú get_debugrÚdebug)rrrrÚ pause_writing­s  zFlowControlMixin.pause_writingcCsP|js t‚d|_|jjƒr&tjd|ƒ|j}|dk rLd|_|jƒsL|jdƒdS)NFz%r resumes writing) r8r;r7r<rr=r9ÚdoneÚ set_result)rÚwaiterrrrÚresume_writing³s   zFlowControlMixin.resume_writingcCsVd|_|jsdS|j}|dkr"dSd|_|jƒr4dS|dkrH|jdƒn |j|ƒdS)NT)r:r8r9r?r@Ú set_exception)rÚexcrArrrÚconnection_lost¿s z FlowControlMixin.connection_lostccsP|jrtdƒ‚|jsdS|j}|dks2|jƒs2t‚|jjƒ}||_|EdHdS)NzConnection lost)r:ÚConnectionResetErrorr8r9Ú cancelledr;r7Ú create_future)rrArrrÚ _drain_helperÏs zFlowControlMixin._drain_helper)N) rr r!r"rr>rBrErrIrrrrr6šs   r6csFeZdZdZd ‡fdd„ Zdd„Z‡fdd„Zd d „Zd d „Z‡Z S)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.) Ncs*tƒj|d�||_d|_||_d|_dS)N)r'F)rrÚ_stream_readerÚ_stream_writerÚ_client_connected_cbÚ _over_ssl)rZ stream_readerr3r')rrrrås zStreamReaderProtocol.__init__cCsd|jj|ƒ|jdƒdk |_|jdk r`t|||j|jƒ|_|j|j|jƒ}tj |ƒr`|jj |ƒdS)NZ sslcontext) rJÚ set_transportÚget_extra_inforMrLrr7rKr Z iscoroutineZ create_task)rr0ÚresrrrÚconnection_madeìs    z$StreamReaderProtocol.connection_madecsF|jdk r*|dkr|jjƒn |jj|ƒtƒj|ƒd|_d|_dS)N)rJÚfeed_eofrCrrErK)rrD)rrrrEøs    z$StreamReaderProtocol.connection_lostcCs|jj|ƒdS)N)rJÚ feed_data)rÚdatarrrÚ data_receivedsz"StreamReaderProtocol.data_receivedcCs|jjƒ|jrdSdS)NFT)rJrRrM)rrrrÚ eof_receiveds z!StreamReaderProtocol.eof_received)NN) rr r!r"rrQrErUrVr#rr)rrrÜs  c@sjeZdZdZdd„Zdd„Zedd„ƒZdd „Zd d „Z d d „Z dd„Z dd„Z ddd„Z edd„ƒ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. cCs2||_||_|dks"t|tƒs"t‚||_||_dS)N)Ú _transportÚ _protocolÚ isinstancerr;Ú_readerr7)rr0r)r/r'rrrrs zStreamWriter.__init__cCs:|jjd|jg}|jdk r,|jd|jƒddj|ƒS)Nz transport=%rz reader=%rz<%s>ú )rrrWrZÚappendÚjoin)rÚinforrrÚ__repr__!s zStreamWriter.__repr__cCs|jS)N)rW)rrrrr0'szStreamWriter.transportcCs|jj|ƒdS)N)rWÚwrite)rrTrrrr`+szStreamWriter.writecCs|jj|ƒdS)N)rWÚ writelines)rrTrrrra.szStreamWriter.writelinescCs |jjƒS)N)rWÚ write_eof)rrrrrb1szStreamWriter.write_eofcCs |jjƒS)N)rWÚ can_write_eof)rrrrrc4szStreamWriter.can_write_eofcCs |jjƒS)N)rWÚclose)rrrrrd7szStreamWriter.closeNcCs|jj||ƒS)N)rWrO)rÚnameÚdefaultrrrrO:szStreamWriter.get_extra_infoccsN|jdk r |jjƒ}|dk r |‚|jdk r:|jjƒr:dV|jjƒEdHdS)z~Flush the write buffer. The intended use is to write w.write(data) yield from w.drain() N)rZÚ exceptionrWZ is_closingrXrI)rrDrrrÚdrain=s    zStreamWriter.drain)N)rr r!r"rr_Úpropertyr0r`rarbrcrdrOrrhrrrrrs  c@sÎeZdZedfdd„Zdd„Zdd„Zdd „Zd d „Zd d „Z dd„Z dd„Z dd„Z dd„Z edd„ƒZedd„ƒZed'dd„ƒZed)dd„ƒZed d!„ƒZejr¼ed"d#„ƒZed$d%„ƒZejrÊd&d#„ZdS)*rNcCsZ|dkrtdƒ‚||_|dkr*tjƒ|_n||_tƒ|_d|_d|_d|_ d|_ d|_ dS)NrzLimit cannot be <= 0F) Ú ValueErrorÚ_limitrr+r7Ú bytearrayÚ_bufferÚ_eofÚ_waiterÚ _exceptionrWr8)rr(r'rrrrXs zStreamReader.__init__cCsªdg}|jr |jdt|jƒƒ|jr0|jdƒ|jtkrJ|jd|jƒ|jr`|jd|jƒ|jrv|jd|jƒ|jrŒ|jd|jƒ|j rœ|jdƒd d j |ƒS) Nrz%d bytesÚeofzl=%dzw=%rze=%rzt=%rZpausedz<%s>r[) rmr\rrnrkÚ_DEFAULT_LIMITrorprWr8r])rr^rrrr_ks    zStreamReader.__repr__cCs|jS)N)rp)rrrrrg}szStreamReader.exceptioncCs0||_|j}|dk r,d|_|jƒs,|j|ƒdS)N)rprorGrC)rrDrArrrrC€s zStreamReader.set_exceptioncCs*|j}|dk r&d|_|jƒs&|jdƒdS)z1Wakeup read*() functions waiting for data or EOF.N)rorGr@)rrArrrÚ_wakeup_waiter‰s zStreamReader._wakeup_waitercCs|jdkstdƒ‚||_dS)NzTransport already set)rWr;)rr0rrrrN‘szStreamReader.set_transportcCs*|jr&t|jƒ|jkr&d|_|jjƒdS)NF)r8rrmrkrWÚresume_reading)rrrrÚ_maybe_resume_transport•sz$StreamReader._maybe_resume_transportcCsd|_|jƒdS)NT)rnrs)rrrrrRšszStreamReader.feed_eofcCs|jo |j S)z=Return True if the buffer is empty and 'feed_eof' was called.)rnrm)rrrrÚat_eofžszStreamReader.at_eofc Cs†|j stdƒ‚|sdS|jj|ƒ|jƒ|jdk r‚|j r‚t|jƒd|jkr‚y|jj ƒWnt k rzd|_YnXd|_dS)Nzfeed_data after feed_eofrT) rnr;rmÚextendrsrWr8rrkZ pause_readingÚNotImplementedError)rrTrrrrS¢s   zStreamReader.feed_datac csf|jdk rtd|ƒ‚|j s&tdƒ‚|jrsB       "  B3G