diff options
| author | Richard Oudkerk <shibturn@gmail.com> | 2013-10-03 14:52:24 +0100 |
|---|---|---|
| committer | Richard Oudkerk <shibturn@gmail.com> | 2013-10-03 14:52:24 +0100 |
| commit | 421672a2a52c68c206f5e328b77ba5138f14a080 (patch) | |
| tree | f7dec6dffe274977adb3ec1629a929d5489e41e3 | |
| parent | e24683348114aad4ba9be61a819cf9e79a5442bf (diff) | |
| download | trollius-git-421672a2a52c68c206f5e328b77ba5138f14a080.tar.gz | |
Fix pause() and resume() for proactor transports:
they should control reading not writing.
| -rw-r--r-- | tests/proactor_events_test.py | 42 | ||||
| -rw-r--r-- | tulip/proactor_events.py | 41 |
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) |
