diff options
| author | Nikolay Kim <fafhrd91@gmail.com> | 2013-03-08 22:30:01 -0800 |
|---|---|---|
| committer | Nikolay Kim <fafhrd91@gmail.com> | 2013-03-08 22:30:01 -0800 |
| commit | ccdd90bbad0f951c06c5f80be93b0dbbbc847122 (patch) | |
| tree | cf592b28eec9ba8345fda4c91bfca08aad8420a7 | |
| parent | bca244f1642f1364c70b6e927e68fe4b43ab65d1 (diff) | |
| download | trollius-git-ccdd90bbad0f951c06c5f80be93b0dbbbc847122.tar.gz | |
proactor socket transport optimization
| -rw-r--r-- | tests/events_test.py | 2 | ||||
| -rw-r--r-- | tulip/proactor_events.py | 26 |
2 files changed, 16 insertions, 12 deletions
diff --git a/tests/events_test.py b/tests/events_test.py index 4c3056a..807aa50 100644 --- a/tests/events_test.py +++ b/tests/events_test.py @@ -513,7 +513,6 @@ class EventLoopTestsMixin: sock = self.event_loop.run_until_complete(f) host, port = sock.getsockname() self.assertEqual(host, '0.0.0.0') - self.event_loop.run_once(0.01) # for windows proactor selector client = socket.socket() client.connect(('127.0.0.1', port)) client.send(b'xxx') @@ -524,7 +523,6 @@ class EventLoopTestsMixin: self.assertEqual('CONNECTED', proto.state) self.assertEqual(0, proto.nbytes) self.event_loop.run_once() - self.event_loop.run_once(0.1) # for windows proactor selector self.assertEqual(3, proto.nbytes) # extra info is available diff --git a/tulip/proactor_events.py b/tulip/proactor_events.py index 0391eb4..45c075e 100644 --- a/tulip/proactor_events.py +++ b/tulip/proactor_events.py @@ -6,10 +6,8 @@ proactor is only implemented on Windows with IOCP. import logging - from . import base_events from . import transports -from . import winsocketpair class _ProactorSocketTransport(transports.Transport): @@ -29,16 +27,18 @@ class _ProactorSocketTransport(transports.Transport): if waiter is not None: self._event_loop.call_soon(waiter.set_result, None) - def _loop_reading(self, f=None): + def _loop_reading(self, fut=None): + data = None + try: - assert f is self._read_fut - if f: - data = f.result() + if fut is not None: + assert fut is self._read_fut + + data = fut.result() # deliver data later in "finally" clause if not data: - self._event_loop.call_soon(self._protocol.eof_received) self._read_fut = None return - self._event_loop.call_soon(self._protocol.data_received, data) + self._read_fut = self._event_loop._proactor.recv(self._sock, 4096) except ConnectionAbortedError as exc: if not self._closing: @@ -47,6 +47,11 @@ class _ProactorSocketTransport(transports.Transport): self._fatal_error(exc) else: self._read_fut.add_done_callback(self._loop_reading) + finally: + if data: + self._protocol.data_received(data) + elif data is not None: + self._protocol.eof_received() def write(self, data): assert isinstance(data, bytes) @@ -149,6 +154,7 @@ class BaseProactorEventLoop(base_events.BaseEventLoop): self._ssock.setblocking(False) self._csock.setblocking(False) self._internal_fds += 1 + def loop(f=None): try: if f: @@ -170,10 +176,10 @@ class BaseProactorEventLoop(base_events.BaseEventLoop): if f: conn, addr = f.result() protocol = protocol_factory() - transport = self._make_socket_transport( + self._make_socket_transport( conn, protocol, extra={'addr': addr}) f = self._proactor.accept(sock) - except OSError as exc: + except OSError: sock.close() logging.exception('Accept failed') else: |
