summaryrefslogtreecommitdiffstats
path: root/Lib/asyncio
diff options
context:
space:
mode:
authorGuido van Rossum <guido@python.org>2014-01-29 21:20:39 (GMT)
committerGuido van Rossum <guido@python.org>2014-01-29 21:20:39 (GMT)
commit47fb97e4e6b3d11435bd0051a1203ee7b3bea34f (patch)
tree19f27a09db65248a17c54d95b0b09ec2a52f4bf0 /Lib/asyncio
parent3ccead1f6a36770b38c3fb561e74eb2e1dcbe76c (diff)
downloadcpython-47fb97e4e6b3d11435bd0051a1203ee7b3bea34f.zip
cpython-47fb97e4e6b3d11435bd0051a1203ee7b3bea34f.tar.gz
cpython-47fb97e4e6b3d11435bd0051a1203ee7b3bea34f.tar.bz2
asyncio: Add write flow control to unix pipes.
Diffstat (limited to 'Lib/asyncio')
-rw-r--r--Lib/asyncio/unix_events.py14
1 files changed, 11 insertions, 3 deletions
diff --git a/Lib/asyncio/unix_events.py b/Lib/asyncio/unix_events.py
index a1aff3f..05aa272 100644
--- a/Lib/asyncio/unix_events.py
+++ b/Lib/asyncio/unix_events.py
@@ -246,7 +246,8 @@ class _UnixReadPipeTransport(transports.ReadTransport):
self._loop = None
-class _UnixWritePipeTransport(transports.WriteTransport):
+class _UnixWritePipeTransport(selector_events._FlowControlMixin,
+ transports.WriteTransport):
def __init__(self, loop, pipe, protocol, waiter=None, extra=None):
super().__init__(extra)
@@ -277,12 +278,17 @@ class _UnixWritePipeTransport(transports.WriteTransport):
if waiter is not None:
self._loop.call_soon(waiter.set_result, None)
+ def get_write_buffer_size(self):
+ return sum(len(data) for data in self._buffer)
+
def _read_ready(self):
# Pipe was closed by peer.
self._close()
def write(self, data):
- assert isinstance(data, bytes), repr(data)
+ assert isinstance(data, (bytes, bytearray, memoryview)), repr(data)
+ if isinstance(data, bytearray):
+ data = memoryview(data)
if not data:
return
@@ -310,6 +316,7 @@ class _UnixWritePipeTransport(transports.WriteTransport):
self._loop.add_writer(self._fileno, self._write_ready)
self._buffer.append(data)
+ self._maybe_pause_protocol()
def _write_ready(self):
data = b''.join(self._buffer)
@@ -329,7 +336,8 @@ class _UnixWritePipeTransport(transports.WriteTransport):
else:
if n == len(data):
self._loop.remove_writer(self._fileno)
- if self._closing:
+ self._maybe_resume_protocol() # May append to buffer.
+ if not self._buffer and self._closing:
self._loop.remove_reader(self._fileno)
self._call_connection_lost(None)
return