self._proc = None
self._pid = None
self._returncode = None
- self._exit_waiters = []
+ self._exit_waiters = set()
self._pending_calls = collections.deque()
self._pipes = {}
self._finished = False
except (SystemExit, KeyboardInterrupt):
raise
except BaseException as exc:
+ # Close any pipes that were already connected before the
+ # error/cancellation to avoid leaking file descriptors.
+ for proto in self._pipes.values():
+ if proto is not None:
+ proto.pipe.close()
+ for raw_pipe in (proc.stdin, proc.stdout, proc.stderr):
+ if raw_pipe is not None:
+ raw_pipe.close()
if waiter is not None and not waiter.cancelled():
waiter.set_exception(exc)
else:
return self._returncode
waiter = self._loop.create_future()
- self._exit_waiters.append(waiter)
- return await waiter
+ self._exit_waiters.add(waiter)
+ try:
+ return await waiter
+ finally:
+ if self._exit_waiters is not None:
+ self._exit_waiters.discard(waiter)
def _try_finish(self):
assert not self._finished
from asyncio import subprocess
from test.test_asyncio import utils as test_utils
from test import support
-from test.support import os_helper, warnings_helper, gc_collect
+from test.support import os_helper, gc_collect
if not support.has_subprocess_support:
raise unittest.SkipTest("test module requires subprocess")
self.loop.run_until_complete(main())
- @warnings_helper.ignore_warnings(category=ResourceWarning)
def test_subprocess_read_pipe_cancelled(self):
async def main():
loop = asyncio.get_running_loop()
asyncio.run(main())
gc_collect()
- @warnings_helper.ignore_warnings(category=ResourceWarning)
def test_subprocess_write_pipe_cancelled(self):
async def main():
loop = asyncio.get_running_loop()
asyncio.run(main())
gc_collect()
- @warnings_helper.ignore_warnings(category=ResourceWarning)
def test_subprocess_read_write_pipe_cancelled(self):
async def main():
loop = asyncio.get_running_loop()