Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
36 changes: 33 additions & 3 deletions Lib/asyncio/selector_events.py
Original file line number Diff line number Diff line change
Expand Up @@ -878,11 +878,39 @@ def close(self):
self._call_soon(self._call_connection_lost, None)

def __del__(self, _warn=warnings.warn):
# The transport can be resurrected after this runs, so leave it in the
# state _force_close() and _call_connection_lost() would: the fd is
# gone, and a stale self._sock_fd would later be used to unregister a
# descriptor the OS has handed out to somebody else.
if self._sock is not None:
_warn(f"unclosed transport {self!r}", ResourceWarning, source=self)
# Warn before cleaning up: the message embeds repr(self), which
# reports the fd number.
if self._protocol_connected:
self._protocol_connected = False
_warn(f"unclosed transport {self!r}", ResourceWarning,
source=self)

if self._buffer:
self._buffer.clear()
self._buffer_size = 0
self._loop._remove_writer(self._sock_fd)

if not self._closing:
self._closing = True
self._loop._remove_reader(self._sock_fd)

self._conn_lost += 1

self._sock_fd = -1
self._sock.close()
if self._server is not None:
self._server._detach(self)
self._sock = None
self._protocol = None
self._loop = None

server = self._server
if server is not None:
self._server = None
server._detach(self)

def _fatal_error(self, exc, message='Fatal error on transport'):
# Should be called from exception handler only.
Expand Down Expand Up @@ -914,8 +942,10 @@ def _force_close(self, exc):
def _call_connection_lost(self, exc):
try:
if self._protocol_connected:
self._protocol_connected = False
self._protocol.connection_lost(exc)
finally:
self._sock_fd = -1
self._sock.close()
self._sock = None
self._protocol = None
Expand Down
94 changes: 94 additions & 0 deletions Lib/test/test_asyncio/test_selector_events.py
Original file line number Diff line number Diff line change
Expand Up @@ -542,6 +542,75 @@ def test_force_close(self):
self.assertFalse(self.loop.readers)
self.assertEqual(1, self.loop.remove_reader_count[7])

def test_del(self):
tr = self.create_transport()
self.loop._add_reader(7, mock.sentinel)

with self.assertWarns(ResourceWarning):
tr.__del__()

# The socket is closed, so fd 7 may be handed out to an unrelated
# file at any moment: the loop must not be left polling it, and the
# cached fd must no longer look valid.
self.assertFalse(self.loop.readers)
self.assertEqual(1, self.loop.remove_reader_count[7])
self.assertEqual(-1, tr._sock_fd)
self.sock.close.assert_called_with()
self.assertIsNone(tr._sock)
self.assertIsNone(tr._protocol)
self.assertIsNone(tr._loop)
self.assertTrue(tr.is_closing())

def test_del_write_buffer(self):
tr = self.create_transport()
tr._buffer.extend(b'data')
tr._buffer_size = 4
self.loop._add_reader(7, mock.sentinel)
self.loop._add_writer(7, mock.sentinel)

with self.assertWarns(ResourceWarning):
tr.__del__()

self.assertFalse(self.loop.readers)
self.assertFalse(self.loop.writers)
self.assertEqual(1, self.loop.remove_writer_count[7])
self.assertEqual(tr._buffer, list_to_buffer())
self.assertEqual(0, tr.get_write_buffer_size())

def test_del_warning_names_the_fd(self):
# The warning has to be issued before the cleanup, otherwise its
# repr(self) reports a transport that is already closed and the fd
# number, the only actionable part of the message, is lost.
tr = self.create_transport()

with self.assertWarns(ResourceWarning) as cm:
tr.__del__()

self.assertIn('fd=7', str(cm.warning))

def test_del_then_close_leaves_reused_fd_alone(self):
tr = self.create_transport()
self.loop._add_reader(7, mock.sentinel)

with self.assertWarns(ResourceWarning):
tr.__del__()

# Something else in the process now owns fd 7 and waits on it.
self.loop._add_reader(7, mock.sentinel.other)
self.loop._add_writer(7, mock.sentinel.other)
self.loop.reset_counters()

# The transport got resurrected during garbage collection and is
# closed properly by its new owner. It no longer owns fd 7, so it
# must keep its hands off the loop.
tr.close()
tr.abort()

self.assertEqual(0, self.loop.remove_reader_count[7])
self.assertEqual(0, self.loop.remove_writer_count[7])
self.assertIs(mock.sentinel.other, self.loop.readers[7]._callback)
self.assertIs(mock.sentinel.other, self.loop.writers[7]._callback)

@mock.patch('asyncio.log.logger.error')
def test_fatal_error(self, m_exc):
exc = OSError()
Expand Down Expand Up @@ -599,6 +668,31 @@ def test__add_reader(self):
self.assertFalse(self.loop.readers)


class SelectorTransportDelTests(test_utils.TestCase):
"""__del__ against a real event loop and a real socket."""

def test_del_unregisters_fd_from_the_selector(self):
loop = asyncio.SelectorEventLoop()
self.set_event_loop(loop)
rsock, wsock = socket.socketpair()
self.addCleanup(wsock.close)
self.addCleanup(rsock.close)

protocol = test_utils.make_test_protocol(asyncio.Protocol)
tr = _SelectorSocketTransport(loop, rsock, protocol)
test_utils.run_briefly(loop) # let connection_made() and _add_reader()
fd = tr._sock_fd
self.assertIn(fd, loop._selector.get_map())

with self.assertWarns(ResourceWarning):
tr.__del__()

# Leaving the fd registered here would make the loop poll a
# descriptor owned by whatever opens a file next.
self.assertNotIn(fd, loop._selector.get_map())
self.assertEqual(-1, tr._sock_fd)


class SelectorSocketTransportTests(test_utils.TestCase):

def setUp(self):
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,5 @@
Fix :mod:`asyncio` selector transports unregistering a file descriptor they
no longer own. ``_SelectorTransport.__del__`` closed the socket but left the
descriptor registered with the event loop, so a transport resurrected during
garbage collection could later remove an unrelated descriptor from the
selector. ``repr()`` of a closed transport now reports ``fd=-1``.
Loading