From 99aa5f066b8f1b5d6c7822cf970a7984d4ebb528 Mon Sep 17 00:00:00 2001 From: Andrey Egorov Date: Mon, 13 Nov 2017 16:03:59 +0300 Subject: [PATCH 1/5] bpo-32015: Asyncio cycling during simultaneously socket read/write and reconnection --- Lib/asyncio/selector_events.py | 37 +++++++++++++++++----------------- 1 file changed, 18 insertions(+), 19 deletions(-) diff --git a/Lib/asyncio/selector_events.py b/Lib/asyncio/selector_events.py index 7143ca2660afaa..54e8f2b9cbcffb 100644 --- a/Lib/asyncio/selector_events.py +++ b/Lib/asyncio/selector_events.py @@ -362,25 +362,25 @@ def sock_recv(self, sock, n): if self._debug and sock.gettimeout() != 0: raise ValueError("the socket must be non-blocking") fut = self.create_future() - self._sock_recv(fut, False, sock, n) + self._sock_recv(fut, None, sock, n) return fut - def _sock_recv(self, fut, registered, sock, n): + def _sock_recv(self, fut, registered_fd, sock, n): # _sock_recv() can add itself as an I/O callback if the operation can't # be done immediately. Don't use it directly, call sock_recv(). - fd = sock.fileno() - if registered: + if registered_fd is not None: # Remove the callback early. It should be rare that the # selector says the fd is ready but the call still returns # EAGAIN, and I am willing to take a hit in that case in # order to simplify the common case. - self.remove_reader(fd) + self.remove_reader(registered_fd) if fut.cancelled(): return try: data = sock.recv(n) except (BlockingIOError, InterruptedError): - self.add_reader(fd, self._sock_recv, fut, True, sock, n) + fd = sock.fileno() + self.add_reader(fd, self._sock_recv, fut, fd, sock, n) except Exception as exc: fut.set_exception(exc) else: @@ -397,25 +397,25 @@ def sock_recv_into(self, sock, buf): if self._debug and sock.gettimeout() != 0: raise ValueError("the socket must be non-blocking") fut = self.create_future() - self._sock_recv_into(fut, False, sock, buf) + self._sock_recv_into(fut, None, sock, buf) return fut - def _sock_recv_into(self, fut, registered, sock, buf): + def _sock_recv_into(self, fut, registered_fd, sock, buf): # _sock_recv_into() can add itself as an I/O callback if the operation # can't be done immediately. Don't use it directly, call sock_recv_into(). - fd = sock.fileno() - if registered: + if registered_fd is not None: # Remove the callback early. It should be rare that the # selector says the fd is ready but the call still returns # EAGAIN, and I am willing to take a hit in that case in # order to simplify the common case. - self.remove_reader(fd) + self.remove_reader(registered_fd) if fut.cancelled(): return try: nbytes = sock.recv_into(buf) except (BlockingIOError, InterruptedError): - self.add_reader(fd, self._sock_recv_into, fut, True, sock, buf) + fd = sock.fileno() + self.add_reader(fd, self._sock_recv_into, fut, fd, sock, buf) except Exception as exc: fut.set_exception(exc) else: @@ -436,16 +436,14 @@ def sock_sendall(self, sock, data): raise ValueError("the socket must be non-blocking") fut = self.create_future() if data: - self._sock_sendall(fut, False, sock, data) + self._sock_sendall(fut, None, sock, data) else: fut.set_result(None) return fut - def _sock_sendall(self, fut, registered, sock, data): - fd = sock.fileno() - - if registered: - self.remove_writer(fd) + def _sock_sendall(self, fut, registered_fd, sock, data): + if registered_fd is not None: + self.remove_writer(registered_fd) if fut.cancelled(): return @@ -462,7 +460,8 @@ def _sock_sendall(self, fut, registered, sock, data): else: if n: data = data[n:] - self.add_writer(fd, self._sock_sendall, fut, True, sock, data) + fd = sock.fileno() + self.add_writer(fd, self._sock_sendall, fut, fd, sock, data) @coroutine def sock_connect(self, sock, address): From 02d9c370485a481acb77eb719b81f655dcfa1725 Mon Sep 17 00:00:00 2001 From: Andrey Egorov Date: Mon, 13 Nov 2017 16:59:48 +0300 Subject: [PATCH 2/5] Tests fix --- Lib/test/test_asyncio/test_selector_events.py | 42 +++++++++---------- 1 file changed, 21 insertions(+), 21 deletions(-) diff --git a/Lib/test/test_asyncio/test_selector_events.py b/Lib/test/test_asyncio/test_selector_events.py index c50b3e49565c92..bde686693db720 100644 --- a/Lib/test/test_asyncio/test_selector_events.py +++ b/Lib/test/test_asyncio/test_selector_events.py @@ -182,7 +182,7 @@ def test_sock_recv(self): f = self.loop.sock_recv(sock, 1024) self.assertIsInstance(f, asyncio.Future) - self.loop._sock_recv.assert_called_with(f, False, sock, 1024) + self.loop._sock_recv.assert_called_with(f, None, sock, 1024) def test__sock_recv_canceled_fut(self): sock = mock.Mock() @@ -190,7 +190,7 @@ def test__sock_recv_canceled_fut(self): f = asyncio.Future(loop=self.loop) f.cancel() - self.loop._sock_recv(f, False, sock, 1024) + self.loop._sock_recv(f, None, sock, 1024) self.assertFalse(sock.recv.called) def test__sock_recv_unregister(self): @@ -201,8 +201,8 @@ def test__sock_recv_unregister(self): f.cancel() self.loop.remove_reader = mock.Mock() - self.loop._sock_recv(f, True, sock, 1024) - self.assertEqual((10,), self.loop.remove_reader.call_args[0]) + self.loop._sock_recv(f, 10, sock, 1024) + self.assertEqual((True,), self.loop.remove_reader.call_args[0]) def test__sock_recv_tryagain(self): f = asyncio.Future(loop=self.loop) @@ -211,8 +211,8 @@ def test__sock_recv_tryagain(self): sock.recv.side_effect = BlockingIOError self.loop.add_reader = mock.Mock() - self.loop._sock_recv(f, False, sock, 1024) - self.assertEqual((10, self.loop._sock_recv, f, True, sock, 1024), + self.loop._sock_recv(f, None, sock, 1024) + self.assertEqual((10, self.loop._sock_recv, f, 10, sock, 1024), self.loop.add_reader.call_args[0]) def test__sock_recv_exception(self): @@ -221,7 +221,7 @@ def test__sock_recv_exception(self): sock.fileno.return_value = 10 err = sock.recv.side_effect = OSError() - self.loop._sock_recv(f, False, sock, 1024) + self.loop._sock_recv(f, None, sock, 1024) self.assertIs(err, f.exception()) def test_sock_sendall(self): @@ -231,7 +231,7 @@ def test_sock_sendall(self): f = self.loop.sock_sendall(sock, b'data') self.assertIsInstance(f, asyncio.Future) self.assertEqual( - (f, False, sock, b'data'), + (f, None, sock, b'data'), self.loop._sock_sendall.call_args[0]) def test_sock_sendall_nodata(self): @@ -250,7 +250,7 @@ def test__sock_sendall_canceled_fut(self): f = asyncio.Future(loop=self.loop) f.cancel() - self.loop._sock_sendall(f, False, sock, b'data') + self.loop._sock_sendall(f, None, sock, b'data') self.assertFalse(sock.send.called) def test__sock_sendall_unregister(self): @@ -261,8 +261,8 @@ def test__sock_sendall_unregister(self): f.cancel() self.loop.remove_writer = mock.Mock() - self.loop._sock_sendall(f, True, sock, b'data') - self.assertEqual((10,), self.loop.remove_writer.call_args[0]) + self.loop._sock_sendall(f, 10, sock, b'data') + self.assertEqual((True,), self.loop.remove_writer.call_args[0]) def test__sock_sendall_tryagain(self): f = asyncio.Future(loop=self.loop) @@ -271,9 +271,9 @@ def test__sock_sendall_tryagain(self): sock.send.side_effect = BlockingIOError self.loop.add_writer = mock.Mock() - self.loop._sock_sendall(f, False, sock, b'data') + self.loop._sock_sendall(f, None, sock, b'data') self.assertEqual( - (10, self.loop._sock_sendall, f, True, sock, b'data'), + (10, self.loop._sock_sendall, f, 10, sock, b'data'), self.loop.add_writer.call_args[0]) def test__sock_sendall_interrupted(self): @@ -283,9 +283,9 @@ def test__sock_sendall_interrupted(self): sock.send.side_effect = InterruptedError self.loop.add_writer = mock.Mock() - self.loop._sock_sendall(f, False, sock, b'data') + self.loop._sock_sendall(f, None, sock, b'data') self.assertEqual( - (10, self.loop._sock_sendall, f, True, sock, b'data'), + (10, self.loop._sock_sendall, f, 10, sock, b'data'), self.loop.add_writer.call_args[0]) def test__sock_sendall_exception(self): @@ -294,7 +294,7 @@ def test__sock_sendall_exception(self): sock.fileno.return_value = 10 err = sock.send.side_effect = OSError() - self.loop._sock_sendall(f, False, sock, b'data') + self.loop._sock_sendall(f, None, sock, b'data') self.assertIs(f.exception(), err) def test__sock_sendall(self): @@ -304,7 +304,7 @@ def test__sock_sendall(self): sock.fileno.return_value = 10 sock.send.return_value = 4 - self.loop._sock_sendall(f, False, sock, b'data') + self.loop._sock_sendall(f, None, sock, b'data') self.assertTrue(f.done()) self.assertIsNone(f.result()) @@ -316,10 +316,10 @@ def test__sock_sendall_partial(self): sock.send.return_value = 2 self.loop.add_writer = mock.Mock() - self.loop._sock_sendall(f, False, sock, b'data') + self.loop._sock_sendall(f, None, sock, b'data') self.assertFalse(f.done()) self.assertEqual( - (10, self.loop._sock_sendall, f, True, sock, b'ta'), + (10, self.loop._sock_sendall, f, 10, sock, b'ta'), self.loop.add_writer.call_args[0]) def test__sock_sendall_none(self): @@ -330,10 +330,10 @@ def test__sock_sendall_none(self): sock.send.return_value = 0 self.loop.add_writer = mock.Mock() - self.loop._sock_sendall(f, False, sock, b'data') + self.loop._sock_sendall(f, None, sock, b'data') self.assertFalse(f.done()) self.assertEqual( - (10, self.loop._sock_sendall, f, True, sock, b'data'), + (10, self.loop._sock_sendall, f, 10, sock, b'data'), self.loop.add_writer.call_args[0]) def test_sock_connect_timeout(self): From 8ac2de80dd879642b3e985e23ac0cd53ce8ed931 Mon Sep 17 00:00:00 2001 From: Andrey Egorov Date: Mon, 13 Nov 2017 17:33:21 +0300 Subject: [PATCH 3/5] Tests fix --- Lib/test/test_asyncio/test_selector_events.py | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/Lib/test/test_asyncio/test_selector_events.py b/Lib/test/test_asyncio/test_selector_events.py index bde686693db720..0f65f95d64f957 100644 --- a/Lib/test/test_asyncio/test_selector_events.py +++ b/Lib/test/test_asyncio/test_selector_events.py @@ -202,7 +202,7 @@ def test__sock_recv_unregister(self): self.loop.remove_reader = mock.Mock() self.loop._sock_recv(f, 10, sock, 1024) - self.assertEqual((True,), self.loop.remove_reader.call_args[0]) + self.assertEqual((10,), self.loop.remove_reader.call_args[0]) def test__sock_recv_tryagain(self): f = asyncio.Future(loop=self.loop) @@ -262,7 +262,7 @@ def test__sock_sendall_unregister(self): self.loop.remove_writer = mock.Mock() self.loop._sock_sendall(f, 10, sock, b'data') - self.assertEqual((True,), self.loop.remove_writer.call_args[0]) + self.assertEqual((10,), self.loop.remove_writer.call_args[0]) def test__sock_sendall_tryagain(self): f = asyncio.Future(loop=self.loop) From 82368d2836e08104ebe752af95a10c62cf51e401 Mon Sep 17 00:00:00 2001 From: Andrey Egorov Date: Mon, 13 Nov 2017 17:49:02 +0300 Subject: [PATCH 4/5] News add --- .../next/Library/2017-11-13-17-48-33.bpo-32015.4nqRTD.rst | 2 ++ 1 file changed, 2 insertions(+) create mode 100644 Misc/NEWS.d/next/Library/2017-11-13-17-48-33.bpo-32015.4nqRTD.rst diff --git a/Misc/NEWS.d/next/Library/2017-11-13-17-48-33.bpo-32015.4nqRTD.rst b/Misc/NEWS.d/next/Library/2017-11-13-17-48-33.bpo-32015.4nqRTD.rst new file mode 100644 index 00000000000000..6117e5625d7edb --- /dev/null +++ b/Misc/NEWS.d/next/Library/2017-11-13-17-48-33.bpo-32015.4nqRTD.rst @@ -0,0 +1,2 @@ +Fixed the looping of asyncio in the case of reconnection the socket during +waiting async read/write from/to the socket. From fa541f1000fbea13629fbe8db868f8c6f672539e Mon Sep 17 00:00:00 2001 From: Andrey Egorov Date: Tue, 14 Nov 2017 03:09:39 +0300 Subject: [PATCH 5/5] Add new unit tests --- Lib/test/test_asyncio/test_selector_events.py | 40 +++++++++++++++++++ 1 file changed, 40 insertions(+) diff --git a/Lib/test/test_asyncio/test_selector_events.py b/Lib/test/test_asyncio/test_selector_events.py index 0f65f95d64f957..a3d118e18812f2 100644 --- a/Lib/test/test_asyncio/test_selector_events.py +++ b/Lib/test/test_asyncio/test_selector_events.py @@ -184,6 +184,26 @@ def test_sock_recv(self): self.assertIsInstance(f, asyncio.Future) self.loop._sock_recv.assert_called_with(f, None, sock, 1024) + def test_sock_recv_reconnection(self): + sock = mock.Mock() + sock.fileno.return_value = 10 + sock.recv.side_effect = BlockingIOError + + self.loop.add_reader = mock.Mock() + self.loop.remove_reader = mock.Mock() + fut = self.loop.sock_recv(sock, 1024) + callback = self.loop.add_reader.call_args[0][1] + params = self.loop.add_reader.call_args[0][2:] + + # emulate the old socket has closed, but the new one has + # the same fileno, so callback is called with old (closed) socket + sock.fileno.return_value = -1 + sock.recv.side_effect = OSError(9) + callback(*params) + + self.assertIsInstance(fut.exception(), OSError) + self.assertEqual((10,), self.loop.remove_reader.call_args[0]) + def test__sock_recv_canceled_fut(self): sock = mock.Mock() @@ -244,6 +264,26 @@ def test_sock_sendall_nodata(self): self.assertIsNone(f.result()) self.assertFalse(self.loop._sock_sendall.called) + def test_sock_sendall_reconnection(self): + sock = mock.Mock() + sock.fileno.return_value = 10 + sock.send.side_effect = BlockingIOError + + self.loop.add_writer = mock.Mock() + self.loop.remove_writer = mock.Mock() + fut = self.loop.sock_sendall(sock, b'data') + callback = self.loop.add_writer.call_args[0][1] + params = self.loop.add_writer.call_args[0][2:] + + # emulate the old socket has closed, but the new one has + # the same fileno, so callback is called with old (closed) socket + sock.fileno.return_value = -1 + sock.send.side_effect = OSError(9) + callback(*params) + + self.assertIsInstance(fut.exception(), OSError) + self.assertEqual((10,), self.loop.remove_writer.call_args[0]) + def test__sock_sendall_canceled_fut(self): sock = mock.Mock()