streams.py 26 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752
  1. __all__ = (
  2. 'StreamReader', 'StreamWriter', 'StreamReaderProtocol',
  3. 'open_connection', 'start_server')
  4. import collections
  5. import socket
  6. import sys
  7. import weakref
  8. if hasattr(socket, 'AF_UNIX'):
  9. __all__ += ('open_unix_connection', 'start_unix_server')
  10. from . import coroutines
  11. from . import events
  12. from . import exceptions
  13. from . import format_helpers
  14. from . import protocols
  15. from .log import logger
  16. from .tasks import sleep
  17. _DEFAULT_LIMIT = 2 ** 16 # 64 KiB
  18. async def open_connection(host=None, port=None, *,
  19. limit=_DEFAULT_LIMIT, **kwds):
  20. """A wrapper for create_connection() returning a (reader, writer) pair.
  21. The reader returned is a StreamReader instance; the writer is a
  22. StreamWriter instance.
  23. The arguments are all the usual arguments to create_connection()
  24. except protocol_factory; most common are positional host and port,
  25. with various optional keyword arguments following.
  26. Additional optional keyword arguments are loop (to set the event loop
  27. instance to use) and limit (to set the buffer limit passed to the
  28. StreamReader).
  29. (If you want to customize the StreamReader and/or
  30. StreamReaderProtocol classes, just copy the code -- there's
  31. really nothing special here except some convenience.)
  32. """
  33. loop = events.get_running_loop()
  34. reader = StreamReader(limit=limit, loop=loop)
  35. protocol = StreamReaderProtocol(reader, loop=loop)
  36. transport, _ = await loop.create_connection(
  37. lambda: protocol, host, port, **kwds)
  38. writer = StreamWriter(transport, protocol, reader, loop)
  39. return reader, writer
  40. async def start_server(client_connected_cb, host=None, port=None, *,
  41. limit=_DEFAULT_LIMIT, **kwds):
  42. """Start a socket server, call back for each client connected.
  43. The first parameter, `client_connected_cb`, takes two parameters:
  44. client_reader, client_writer. client_reader is a StreamReader
  45. object, while client_writer is a StreamWriter object. This
  46. parameter can either be a plain callback function or a coroutine;
  47. if it is a coroutine, it will be automatically converted into a
  48. Task.
  49. The rest of the arguments are all the usual arguments to
  50. loop.create_server() except protocol_factory; most common are
  51. positional host and port, with various optional keyword arguments
  52. following. The return value is the same as loop.create_server().
  53. Additional optional keyword arguments are loop (to set the event loop
  54. instance to use) and limit (to set the buffer limit passed to the
  55. StreamReader).
  56. The return value is the same as loop.create_server(), i.e. a
  57. Server object which can be used to stop the service.
  58. """
  59. loop = events.get_running_loop()
  60. def factory():
  61. reader = StreamReader(limit=limit, loop=loop)
  62. protocol = StreamReaderProtocol(reader, client_connected_cb,
  63. loop=loop)
  64. return protocol
  65. return await loop.create_server(factory, host, port, **kwds)
  66. if hasattr(socket, 'AF_UNIX'):
  67. # UNIX Domain Sockets are supported on this platform
  68. async def open_unix_connection(path=None, *,
  69. limit=_DEFAULT_LIMIT, **kwds):
  70. """Similar to `open_connection` but works with UNIX Domain Sockets."""
  71. loop = events.get_running_loop()
  72. reader = StreamReader(limit=limit, loop=loop)
  73. protocol = StreamReaderProtocol(reader, loop=loop)
  74. transport, _ = await loop.create_unix_connection(
  75. lambda: protocol, path, **kwds)
  76. writer = StreamWriter(transport, protocol, reader, loop)
  77. return reader, writer
  78. async def start_unix_server(client_connected_cb, path=None, *,
  79. limit=_DEFAULT_LIMIT, **kwds):
  80. """Similar to `start_server` but works with UNIX Domain Sockets."""
  81. loop = events.get_running_loop()
  82. def factory():
  83. reader = StreamReader(limit=limit, loop=loop)
  84. protocol = StreamReaderProtocol(reader, client_connected_cb,
  85. loop=loop)
  86. return protocol
  87. return await loop.create_unix_server(factory, path, **kwds)
  88. class FlowControlMixin(protocols.Protocol):
  89. """Reusable flow control logic for StreamWriter.drain().
  90. This implements the protocol methods pause_writing(),
  91. resume_writing() and connection_lost(). If the subclass overrides
  92. these it must call the super methods.
  93. StreamWriter.drain() must wait for _drain_helper() coroutine.
  94. """
  95. def __init__(self, loop=None):
  96. if loop is None:
  97. self._loop = events.get_event_loop()
  98. else:
  99. self._loop = loop
  100. self._paused = False
  101. self._drain_waiters = collections.deque()
  102. self._connection_lost = False
  103. def pause_writing(self):
  104. assert not self._paused
  105. self._paused = True
  106. if self._loop.get_debug():
  107. logger.debug("%r pauses writing", self)
  108. def resume_writing(self):
  109. assert self._paused
  110. self._paused = False
  111. if self._loop.get_debug():
  112. logger.debug("%r resumes writing", self)
  113. for waiter in self._drain_waiters:
  114. if not waiter.done():
  115. waiter.set_result(None)
  116. def connection_lost(self, exc):
  117. self._connection_lost = True
  118. # Wake up the writer(s) if currently paused.
  119. if not self._paused:
  120. return
  121. for waiter in self._drain_waiters:
  122. if not waiter.done():
  123. if exc is None:
  124. waiter.set_result(None)
  125. else:
  126. waiter.set_exception(exc)
  127. async def _drain_helper(self):
  128. if self._connection_lost:
  129. raise ConnectionResetError('Connection lost')
  130. if not self._paused:
  131. return
  132. waiter = self._loop.create_future()
  133. self._drain_waiters.append(waiter)
  134. try:
  135. await waiter
  136. finally:
  137. self._drain_waiters.remove(waiter)
  138. def _get_close_waiter(self, stream):
  139. raise NotImplementedError
  140. class StreamReaderProtocol(FlowControlMixin, protocols.Protocol):
  141. """Helper class to adapt between Protocol and StreamReader.
  142. (This is a helper class instead of making StreamReader itself a
  143. Protocol subclass, because the StreamReader has other potential
  144. uses, and to prevent the user of the StreamReader to accidentally
  145. call inappropriate methods of the protocol.)
  146. """
  147. _source_traceback = None
  148. def __init__(self, stream_reader, client_connected_cb=None, loop=None):
  149. super().__init__(loop=loop)
  150. if stream_reader is not None:
  151. self._stream_reader_wr = weakref.ref(stream_reader)
  152. self._source_traceback = stream_reader._source_traceback
  153. else:
  154. self._stream_reader_wr = None
  155. if client_connected_cb is not None:
  156. # This is a stream created by the `create_server()` function.
  157. # Keep a strong reference to the reader until a connection
  158. # is established.
  159. self._strong_reader = stream_reader
  160. self._reject_connection = False
  161. self._stream_writer = None
  162. self._task = None
  163. self._transport = None
  164. self._client_connected_cb = client_connected_cb
  165. self._over_ssl = False
  166. self._closed = self._loop.create_future()
  167. @property
  168. def _stream_reader(self):
  169. if self._stream_reader_wr is None:
  170. return None
  171. return self._stream_reader_wr()
  172. def _replace_writer(self, writer):
  173. loop = self._loop
  174. transport = writer.transport
  175. self._stream_writer = writer
  176. self._transport = transport
  177. self._over_ssl = transport.get_extra_info('sslcontext') is not None
  178. def connection_made(self, transport):
  179. if self._reject_connection:
  180. context = {
  181. 'message': ('An open stream was garbage collected prior to '
  182. 'establishing network connection; '
  183. 'call "stream.close()" explicitly.')
  184. }
  185. if self._source_traceback:
  186. context['source_traceback'] = self._source_traceback
  187. self._loop.call_exception_handler(context)
  188. transport.abort()
  189. return
  190. self._transport = transport
  191. reader = self._stream_reader
  192. if reader is not None:
  193. reader.set_transport(transport)
  194. self._over_ssl = transport.get_extra_info('sslcontext') is not None
  195. if self._client_connected_cb is not None:
  196. self._stream_writer = StreamWriter(transport, self,
  197. reader,
  198. self._loop)
  199. res = self._client_connected_cb(reader,
  200. self._stream_writer)
  201. if coroutines.iscoroutine(res):
  202. self._task = self._loop.create_task(res)
  203. self._strong_reader = None
  204. def connection_lost(self, exc):
  205. reader = self._stream_reader
  206. if reader is not None:
  207. if exc is None:
  208. reader.feed_eof()
  209. else:
  210. reader.set_exception(exc)
  211. if not self._closed.done():
  212. if exc is None:
  213. self._closed.set_result(None)
  214. else:
  215. self._closed.set_exception(exc)
  216. super().connection_lost(exc)
  217. self._stream_reader_wr = None
  218. self._stream_writer = None
  219. self._task = None
  220. self._transport = None
  221. def data_received(self, data):
  222. reader = self._stream_reader
  223. if reader is not None:
  224. reader.feed_data(data)
  225. def eof_received(self):
  226. reader = self._stream_reader
  227. if reader is not None:
  228. reader.feed_eof()
  229. if self._over_ssl:
  230. # Prevent a warning in SSLProtocol.eof_received:
  231. # "returning true from eof_received()
  232. # has no effect when using ssl"
  233. return False
  234. return True
  235. def _get_close_waiter(self, stream):
  236. return self._closed
  237. def __del__(self):
  238. # Prevent reports about unhandled exceptions.
  239. # Better than self._closed._log_traceback = False hack
  240. try:
  241. closed = self._closed
  242. except AttributeError:
  243. pass # failed constructor
  244. else:
  245. if closed.done() and not closed.cancelled():
  246. closed.exception()
  247. class StreamWriter:
  248. """Wraps a Transport.
  249. This exposes write(), writelines(), [can_]write_eof(),
  250. get_extra_info() and close(). It adds drain() which returns an
  251. optional Future on which you can wait for flow control. It also
  252. adds a transport property which references the Transport
  253. directly.
  254. """
  255. def __init__(self, transport, protocol, reader, loop):
  256. self._transport = transport
  257. self._protocol = protocol
  258. # drain() expects that the reader has an exception() method
  259. assert reader is None or isinstance(reader, StreamReader)
  260. self._reader = reader
  261. self._loop = loop
  262. self._complete_fut = self._loop.create_future()
  263. self._complete_fut.set_result(None)
  264. def __repr__(self):
  265. info = [self.__class__.__name__, f'transport={self._transport!r}']
  266. if self._reader is not None:
  267. info.append(f'reader={self._reader!r}')
  268. return '<{}>'.format(' '.join(info))
  269. @property
  270. def transport(self):
  271. return self._transport
  272. def write(self, data):
  273. self._transport.write(data)
  274. def writelines(self, data):
  275. self._transport.writelines(data)
  276. def write_eof(self):
  277. return self._transport.write_eof()
  278. def can_write_eof(self):
  279. return self._transport.can_write_eof()
  280. def close(self):
  281. return self._transport.close()
  282. def is_closing(self):
  283. return self._transport.is_closing()
  284. async def wait_closed(self):
  285. await self._protocol._get_close_waiter(self)
  286. def get_extra_info(self, name, default=None):
  287. return self._transport.get_extra_info(name, default)
  288. async def drain(self):
  289. """Flush the write buffer.
  290. The intended use is to write
  291. w.write(data)
  292. await w.drain()
  293. """
  294. if self._reader is not None:
  295. exc = self._reader.exception()
  296. if exc is not None:
  297. raise exc
  298. if self._transport.is_closing():
  299. # Wait for protocol.connection_lost() call
  300. # Raise connection closing error if any,
  301. # ConnectionResetError otherwise
  302. # Yield to the event loop so connection_lost() may be
  303. # called. Without this, _drain_helper() would return
  304. # immediately, and code that calls
  305. # write(...); await drain()
  306. # in a loop would never call connection_lost(), so it
  307. # would not see an error when the socket is closed.
  308. await sleep(0)
  309. await self._protocol._drain_helper()
  310. async def start_tls(self, sslcontext, *,
  311. server_hostname=None,
  312. ssl_handshake_timeout=None,
  313. ssl_shutdown_timeout=None):
  314. """Upgrade an existing stream-based connection to TLS."""
  315. server_side = self._protocol._client_connected_cb is not None
  316. protocol = self._protocol
  317. await self.drain()
  318. new_transport = await self._loop.start_tls( # type: ignore
  319. self._transport, protocol, sslcontext,
  320. server_side=server_side, server_hostname=server_hostname,
  321. ssl_handshake_timeout=ssl_handshake_timeout,
  322. ssl_shutdown_timeout=ssl_shutdown_timeout)
  323. self._transport = new_transport
  324. protocol._replace_writer(self)
  325. def __del__(self):
  326. if not self._transport.is_closing():
  327. self.close()
  328. class StreamReader:
  329. _source_traceback = None
  330. def __init__(self, limit=_DEFAULT_LIMIT, loop=None):
  331. # The line length limit is a security feature;
  332. # it also doubles as half the buffer limit.
  333. if limit <= 0:
  334. raise ValueError('Limit cannot be <= 0')
  335. self._limit = limit
  336. if loop is None:
  337. self._loop = events.get_event_loop()
  338. else:
  339. self._loop = loop
  340. self._buffer = bytearray()
  341. self._eof = False # Whether we're done.
  342. self._waiter = None # A future used by _wait_for_data()
  343. self._exception = None
  344. self._transport = None
  345. self._paused = False
  346. if self._loop.get_debug():
  347. self._source_traceback = format_helpers.extract_stack(
  348. sys._getframe(1))
  349. def __repr__(self):
  350. info = ['StreamReader']
  351. if self._buffer:
  352. info.append(f'{len(self._buffer)} bytes')
  353. if self._eof:
  354. info.append('eof')
  355. if self._limit != _DEFAULT_LIMIT:
  356. info.append(f'limit={self._limit}')
  357. if self._waiter:
  358. info.append(f'waiter={self._waiter!r}')
  359. if self._exception:
  360. info.append(f'exception={self._exception!r}')
  361. if self._transport:
  362. info.append(f'transport={self._transport!r}')
  363. if self._paused:
  364. info.append('paused')
  365. return '<{}>'.format(' '.join(info))
  366. def exception(self):
  367. return self._exception
  368. def set_exception(self, exc):
  369. self._exception = exc
  370. waiter = self._waiter
  371. if waiter is not None:
  372. self._waiter = None
  373. if not waiter.cancelled():
  374. waiter.set_exception(exc)
  375. def _wakeup_waiter(self):
  376. """Wakeup read*() functions waiting for data or EOF."""
  377. waiter = self._waiter
  378. if waiter is not None:
  379. self._waiter = None
  380. if not waiter.cancelled():
  381. waiter.set_result(None)
  382. def set_transport(self, transport):
  383. assert self._transport is None, 'Transport already set'
  384. self._transport = transport
  385. def _maybe_resume_transport(self):
  386. if self._paused and len(self._buffer) <= self._limit:
  387. self._paused = False
  388. self._transport.resume_reading()
  389. def feed_eof(self):
  390. self._eof = True
  391. self._wakeup_waiter()
  392. def at_eof(self):
  393. """Return True if the buffer is empty and 'feed_eof' was called."""
  394. return self._eof and not self._buffer
  395. def feed_data(self, data):
  396. assert not self._eof, 'feed_data after feed_eof'
  397. if not data:
  398. return
  399. self._buffer.extend(data)
  400. self._wakeup_waiter()
  401. if (self._transport is not None and
  402. not self._paused and
  403. len(self._buffer) > 2 * self._limit):
  404. try:
  405. self._transport.pause_reading()
  406. except NotImplementedError:
  407. # The transport can't be paused.
  408. # We'll just have to buffer all data.
  409. # Forget the transport so we don't keep trying.
  410. self._transport = None
  411. else:
  412. self._paused = True
  413. async def _wait_for_data(self, func_name):
  414. """Wait until feed_data() or feed_eof() is called.
  415. If stream was paused, automatically resume it.
  416. """
  417. # StreamReader uses a future to link the protocol feed_data() method
  418. # to a read coroutine. Running two read coroutines at the same time
  419. # would have an unexpected behaviour. It would not possible to know
  420. # which coroutine would get the next data.
  421. if self._waiter is not None:
  422. raise RuntimeError(
  423. f'{func_name}() called while another coroutine is '
  424. f'already waiting for incoming data')
  425. assert not self._eof, '_wait_for_data after EOF'
  426. # Waiting for data while paused will make deadlock, so prevent it.
  427. # This is essential for readexactly(n) for case when n > self._limit.
  428. if self._paused:
  429. self._paused = False
  430. self._transport.resume_reading()
  431. self._waiter = self._loop.create_future()
  432. try:
  433. await self._waiter
  434. finally:
  435. self._waiter = None
  436. async def readline(self):
  437. """Read chunk of data from the stream until newline (b'\n') is found.
  438. On success, return chunk that ends with newline. If only partial
  439. line can be read due to EOF, return incomplete line without
  440. terminating newline. When EOF was reached while no bytes read, empty
  441. bytes object is returned.
  442. If limit is reached, ValueError will be raised. In that case, if
  443. newline was found, complete line including newline will be removed
  444. from internal buffer. Else, internal buffer will be cleared. Limit is
  445. compared against part of the line without newline.
  446. If stream was paused, this function will automatically resume it if
  447. needed.
  448. """
  449. sep = b'\n'
  450. seplen = len(sep)
  451. try:
  452. line = await self.readuntil(sep)
  453. except exceptions.IncompleteReadError as e:
  454. return e.partial
  455. except exceptions.LimitOverrunError as e:
  456. if self._buffer.startswith(sep, e.consumed):
  457. del self._buffer[:e.consumed + seplen]
  458. else:
  459. self._buffer.clear()
  460. self._maybe_resume_transport()
  461. raise ValueError(e.args[0])
  462. return line
  463. async def readuntil(self, separator=b'\n'):
  464. """Read data from the stream until ``separator`` is found.
  465. On success, the data and separator will be removed from the
  466. internal buffer (consumed). Returned data will include the
  467. separator at the end.
  468. Configured stream limit is used to check result. Limit sets the
  469. maximal length of data that can be returned, not counting the
  470. separator.
  471. If an EOF occurs and the complete separator is still not found,
  472. an IncompleteReadError exception will be raised, and the internal
  473. buffer will be reset. The IncompleteReadError.partial attribute
  474. may contain the separator partially.
  475. If the data cannot be read because of over limit, a
  476. LimitOverrunError exception will be raised, and the data
  477. will be left in the internal buffer, so it can be read again.
  478. """
  479. seplen = len(separator)
  480. if seplen == 0:
  481. raise ValueError('Separator should be at least one-byte string')
  482. if self._exception is not None:
  483. raise self._exception
  484. # Consume whole buffer except last bytes, which length is
  485. # one less than seplen. Let's check corner cases with
  486. # separator='SEPARATOR':
  487. # * we have received almost complete separator (without last
  488. # byte). i.e buffer='some textSEPARATO'. In this case we
  489. # can safely consume len(separator) - 1 bytes.
  490. # * last byte of buffer is first byte of separator, i.e.
  491. # buffer='abcdefghijklmnopqrS'. We may safely consume
  492. # everything except that last byte, but this require to
  493. # analyze bytes of buffer that match partial separator.
  494. # This is slow and/or require FSM. For this case our
  495. # implementation is not optimal, since require rescanning
  496. # of data that is known to not belong to separator. In
  497. # real world, separator will not be so long to notice
  498. # performance problems. Even when reading MIME-encoded
  499. # messages :)
  500. # `offset` is the number of bytes from the beginning of the buffer
  501. # where there is no occurrence of `separator`.
  502. offset = 0
  503. # Loop until we find `separator` in the buffer, exceed the buffer size,
  504. # or an EOF has happened.
  505. while True:
  506. buflen = len(self._buffer)
  507. # Check if we now have enough data in the buffer for `separator` to
  508. # fit.
  509. if buflen - offset >= seplen:
  510. isep = self._buffer.find(separator, offset)
  511. if isep != -1:
  512. # `separator` is in the buffer. `isep` will be used later
  513. # to retrieve the data.
  514. break
  515. # see upper comment for explanation.
  516. offset = buflen + 1 - seplen
  517. if offset > self._limit:
  518. raise exceptions.LimitOverrunError(
  519. 'Separator is not found, and chunk exceed the limit',
  520. offset)
  521. # Complete message (with full separator) may be present in buffer
  522. # even when EOF flag is set. This may happen when the last chunk
  523. # adds data which makes separator be found. That's why we check for
  524. # EOF *ater* inspecting the buffer.
  525. if self._eof:
  526. chunk = bytes(self._buffer)
  527. self._buffer.clear()
  528. raise exceptions.IncompleteReadError(chunk, None)
  529. # _wait_for_data() will resume reading if stream was paused.
  530. await self._wait_for_data('readuntil')
  531. if isep > self._limit:
  532. raise exceptions.LimitOverrunError(
  533. 'Separator is found, but chunk is longer than limit', isep)
  534. chunk = self._buffer[:isep + seplen]
  535. del self._buffer[:isep + seplen]
  536. self._maybe_resume_transport()
  537. return bytes(chunk)
  538. async def read(self, n=-1):
  539. """Read up to `n` bytes from the stream.
  540. If `n` is not provided or set to -1,
  541. read until EOF, then return all read bytes.
  542. If EOF was received and the internal buffer is empty,
  543. return an empty bytes object.
  544. If `n` is 0, return an empty bytes object immediately.
  545. If `n` is positive, return at most `n` available bytes
  546. as soon as at least 1 byte is available in the internal buffer.
  547. If EOF is received before any byte is read, return an empty
  548. bytes object.
  549. Returned value is not limited with limit, configured at stream
  550. creation.
  551. If stream was paused, this function will automatically resume it if
  552. needed.
  553. """
  554. if self._exception is not None:
  555. raise self._exception
  556. if n == 0:
  557. return b''
  558. if n < 0:
  559. # This used to just loop creating a new waiter hoping to
  560. # collect everything in self._buffer, but that would
  561. # deadlock if the subprocess sends more than self.limit
  562. # bytes. So just call self.read(self._limit) until EOF.
  563. blocks = []
  564. while True:
  565. block = await self.read(self._limit)
  566. if not block:
  567. break
  568. blocks.append(block)
  569. return b''.join(blocks)
  570. if not self._buffer and not self._eof:
  571. await self._wait_for_data('read')
  572. # This will work right even if buffer is less than n bytes
  573. data = bytes(memoryview(self._buffer)[:n])
  574. del self._buffer[:n]
  575. self._maybe_resume_transport()
  576. return data
  577. async def readexactly(self, n):
  578. """Read exactly `n` bytes.
  579. Raise an IncompleteReadError if EOF is reached before `n` bytes can be
  580. read. The IncompleteReadError.partial attribute of the exception will
  581. contain the partial read bytes.
  582. if n is zero, return empty bytes object.
  583. Returned value is not limited with limit, configured at stream
  584. creation.
  585. If stream was paused, this function will automatically resume it if
  586. needed.
  587. """
  588. if n < 0:
  589. raise ValueError('readexactly size can not be less than zero')
  590. if self._exception is not None:
  591. raise self._exception
  592. if n == 0:
  593. return b''
  594. while len(self._buffer) < n:
  595. if self._eof:
  596. incomplete = bytes(self._buffer)
  597. self._buffer.clear()
  598. raise exceptions.IncompleteReadError(incomplete, n)
  599. await self._wait_for_data('readexactly')
  600. if len(self._buffer) == n:
  601. data = bytes(self._buffer)
  602. self._buffer.clear()
  603. else:
  604. data = bytes(memoryview(self._buffer)[:n])
  605. del self._buffer[:n]
  606. self._maybe_resume_transport()
  607. return data
  608. def __aiter__(self):
  609. return self
  610. async def __anext__(self):
  611. val = await self.readline()
  612. if val == b'':
  613. raise StopAsyncIteration
  614. return val