summaryrefslogtreecommitdiff
diff options
context:
space:
mode:
authorRichard Oudkerk <shibturn@gmail.com>2013-10-03 14:52:24 +0100
committerRichard Oudkerk <shibturn@gmail.com>2013-10-03 14:52:24 +0100
commit421672a2a52c68c206f5e328b77ba5138f14a080 (patch)
treef7dec6dffe274977adb3ec1629a929d5489e41e3
parente24683348114aad4ba9be61a819cf9e79a5442bf (diff)
downloadtrollius-git-421672a2a52c68c206f5e328b77ba5138f14a080.tar.gz
Fix pause() and resume() for proactor transports:
they should control reading not writing.
-rw-r--r--tests/proactor_events_test.py42
-rw-r--r--tulip/proactor_events.py41
2 files changed, 40 insertions, 43 deletions
diff --git a/tests/proactor_events_test.py b/tests/proactor_events_test.py
index 6a7391d..ae38c03 100644
--- a/tests/proactor_events_test.py
+++ b/tests/proactor_events_test.py
@@ -263,7 +263,7 @@ class ProactorSocketTransportTests(unittest.TestCase):
def test_write_eof_buffer(self):
tr = _ProactorSocketTransport(self.loop, self.sock, self.protocol)
f = tulip.Future(loop=self.loop)
- tr._loop._proactor.send.side_effect = f
+ tr._loop._proactor.send.return_value = f
tr.write(b'data')
tr.write_eof()
self.assertTrue(tr._eof_written)
@@ -287,7 +287,7 @@ class ProactorSocketTransportTests(unittest.TestCase):
def test_write_eof_buffer_write_pipe(self):
tr = _ProactorWritePipeTransport(self.loop, self.sock, self.protocol)
f = tulip.Future(loop=self.loop)
- tr._loop._proactor.send.side_effect = f
+ tr._loop._proactor.send.return_value = f
tr.write(b'data')
tr.write_eof()
self.assertTrue(tr._closing)
@@ -302,29 +302,29 @@ class ProactorSocketTransportTests(unittest.TestCase):
def test_pause_resume(self):
tr = _ProactorSocketTransport(
self.loop, self.sock, self.protocol)
- f = tulip.Future(loop=self.loop)
- tr._loop._proactor.send.side_effect = f
+ futures = []
+ for msg in [b'data1', b'data2', b'data3', b'data4', b'']:
+ f = tulip.Future(loop=self.loop)
+ f.set_result(msg)
+ futures.append(f)
+ self.loop._proactor.recv.side_effect = futures
+ self.loop._run_once()
self.assertFalse(tr._paused)
- tr.write(b'data1')
- tr._loop._proactor.send.assert_called_with(self.sock, b'data1')
- self.assertEqual(tr._buffer, [])
- tr.write(b'data2')
- self.assertEqual(tr._buffer, [b'data2'])
- tr.pause()
- tr.write(b'data3')
- self.assertEqual(tr._buffer, [b'data2', b'data3'])
- f.set_result(5)
self.loop._run_once()
- self.assertEqual(tr._buffer, [b'data2data3'])
+ self.protocol.data_received.assert_called_with(b'data1')
self.loop._run_once()
- self.assertEqual(tr._buffer, [b'data2data3'])
- f = tulip.Future(loop=self.loop)
- tr._loop._proactor.send.side_effect = f
+ self.protocol.data_received.assert_called_with(b'data2')
+ tr.pause()
+ self.assertTrue(tr._paused)
+ for i in range(10):
+ self.loop._run_once()
+ self.protocol.data_received.assert_called_with(b'data2')
tr.resume()
- tr._loop._proactor.send.assert_called_with(self.sock, b'data2data3')
- self.assertEqual(tr._buffer, [])
- tr.write(b'data4')
- self.assertEqual(tr._buffer, [b'data4'])
+ self.assertFalse(tr._paused)
+ self.loop._run_once()
+ self.protocol.data_received.assert_called_with(b'data3')
+ self.loop._run_once()
+ self.protocol.data_received.assert_called_with(b'data4')
tr.close()
diff --git a/tulip/proactor_events.py b/tulip/proactor_events.py
index f897f8c..5b631f6 100644
--- a/tulip/proactor_events.py
+++ b/tulip/proactor_events.py
@@ -28,7 +28,6 @@ class _ProactorBasePipeTransport(transports.BaseTransport):
self._conn_lost = 0
self._closing = False # Set when close() called.
self._eof_written = False
- self._paused = False
self._loop.call_soon(self._protocol.connection_made, self)
if waiter is not None:
self._loop.call_soon(waiter.set_result, None)
@@ -82,9 +81,25 @@ class _ProactorReadPipeTransport(_ProactorBasePipeTransport,
def __init__(self, loop, sock, protocol, waiter=None, extra=None):
super().__init__(loop, sock, protocol, waiter, extra)
+ self._read_fut = None
+ self._paused = False
self._loop.call_soon(self._loop_reading)
+ def pause(self):
+ assert not self._closing, 'Cannot pause() when closing'
+ assert not self._paused, 'Already paused'
+ self._paused = True
+
+ def resume(self):
+ assert self._paused, 'Not paused'
+ self._paused = False
+ if self._closing:
+ return
+ self._loop.call_soon(self._loop_reading, self._read_fut)
+
def _loop_reading(self, fut=None):
+ if self._paused:
+ return
data = None
try:
@@ -130,21 +145,6 @@ class _ProactorWritePipeTransport(_ProactorBasePipeTransport,
transports.WriteTransport):
"""Transport for write pipes."""
- def pause(self):
- assert not self._closing, 'Cannot pause() when closing'
- assert not self._paused, 'Already paused'
- # We don't try to cancel an existing overlapped write. Instead
- # we prevent new overlapped writes until resume() is called.
- self._paused = True
-
- def resume(self):
- assert self._paused, 'Not paused'
- self._paused = False
- if self._closing:
- return
- if self._buffer and self._write_fut is None:
- self._loop_writing()
-
def write(self, data):
assert isinstance(data, bytes), repr(data)
if self._closing or self._eof_written:
@@ -159,7 +159,7 @@ class _ProactorWritePipeTransport(_ProactorBasePipeTransport,
self._conn_lost += 1
return
self._buffer.append(data)
- if self._write_fut is None and not self._paused:
+ if self._write_fut is None:
self._loop_writing()
def _loop_writing(self, f=None):
@@ -176,11 +176,8 @@ class _ProactorWritePipeTransport(_ProactorBasePipeTransport,
if self._eof_written:
self._sock.shutdown(socket.SHUT_WR)
return
- if not self._paused:
- self._write_fut = self._loop._proactor.send(self._sock, data)
- self._write_fut.add_done_callback(self._loop_writing)
- else:
- self._buffer.append(data)
+ self._write_fut = self._loop._proactor.send(self._sock, data)
+ self._write_fut.add_done_callback(self._loop_writing)
except OSError as exc:
self._fatal_error(exc)