Skip to content
This repository was archived by the owner on Nov 23, 2017. It is now read-only.
Open
Changes from 1 commit
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
Prev Previous commit
Next Next commit
make asyncio backward compatible with older Popen
  • Loading branch information
Martiusweb committed Nov 9, 2016
commit ef41fc838a00ce182f58bd36b487e2fa7039d4ce
183 changes: 98 additions & 85 deletions asyncio/unix_events.py
Original file line number Diff line number Diff line change
Expand Up @@ -26,10 +26,11 @@

# XXX temporary: a monkey-patched subprocess.Popen
if compat.PY34:
from . import tmp_subprocess
from .tmp_subprocess import _Popen
else:
# Python 3.3 has a different version of Popen
from . import tmp_subprocess33 as tmp_subprocess
# shows that we can fallback to an older version of subprocess.Popen
# safely: it will block, but asyncio will still work.
_Popen = subprocess.Popen


__all__ = ['SelectorEventLoop',
Expand Down Expand Up @@ -665,97 +666,108 @@ def _set_inheritable(fd, inheritable):
fcntl.fcntl(fd, fcntl.F_SETFD, old & ~cloexec_flag)


class _NonBlockingPopen(tmp_subprocess._Popen):
"""A modified Popen which performs IO operations using an event loop."""
# TODO can we include the stdin trick in popen?
def __init__(self, loop, exec_waiter, watcher, *args, **kwargs):
self._loop = loop
self._watcher = watcher
self._exec_waiter = exec_waiter
super().__init__(*args, **kwargs)

def _cleanup_on_exec_failure(self):
super()._cleanup_on_exec_failure()
self._exec_waiter = None
self._loop = None
self._watcher = None

def _get_exec_err_pipe(self):
errpipe_read, errpipe_write = self._loop._socketpair()
errpipe_read.setblocking(False)
_set_inheritable(errpipe_write.fileno(), False)
return errpipe_read.detach(), errpipe_write.detach()
if hasattr(_Popen, "_wait_exec_done"):
class _NonBlockingPopen(_Popen):
"""A modified Popen which performs IO operations using an event loop."""
def __init__(self, loop, exec_waiter, watcher, *args, **kwargs):
self._loop = loop
self._watcher = watcher
self._exec_waiter = exec_waiter
super().__init__(*args, **kwargs)

def _wait_exec_done(self, orig_executable, cwd, errpipe_read):
errpipe_data = bytearray()
self._loop.add_reader(errpipe_read, self._read_errpipe,
orig_executable, cwd, errpipe_read, errpipe_data)

def _read_errpipe(self, orig_executable, cwd, errpipe_read, errpipe_data):
try:
part = os.read(errpipe_read, 50000)
except BlockingIOError:
return
except Exception as exc:
self._loop.remove_reader(errpipe_read)
os.close(errpipe_read)
self._exec_waiter.set_exception(exc)
self._cleanup_on_exec_failure()
else:
if part and len(errpipe_data) <= 50000:
errpipe_data.extend(part)
def _cleanup_on_exec_failure(self):
super()._cleanup_on_exec_failure()
self._exec_waiter = None
self._loop = None
self._watcher = None

def _get_exec_err_pipe(self):
errpipe_read, errpipe_write = self._loop._socketpair()
errpipe_read.setblocking(False)
_set_inheritable(errpipe_write.fileno(), False)
return errpipe_read.detach(), errpipe_write.detach()

def _wait_exec_done(self, orig_executable, cwd, errpipe_read):
errpipe_data = bytearray()
self._loop.add_reader(errpipe_read, self._read_errpipe,
orig_executable, cwd, errpipe_read,
errpipe_data)

def _read_errpipe(self, orig_executable, cwd, errpipe_read,
errpipe_data):
try:
part = os.read(errpipe_read, 50000)
except BlockingIOError:
return
except Exception as exc:
self._loop.remove_reader(errpipe_read)
os.close(errpipe_read)
self._exec_waiter.set_exception(exc)
self._cleanup_on_exec_failure()
else:
if part and len(errpipe_data) <= 50000:
errpipe_data.extend(part)
return

self._loop.remove_reader(errpipe_read)
os.close(errpipe_read)
self._loop.remove_reader(errpipe_read)
os.close(errpipe_read)

if errpipe_data:
# asynchronously wait until the process terminated
self._watcher.add_child_handler(
self.pid, self._check_exec_result, orig_executable,
cwd, errpipe_data)
else:
if errpipe_data:
# asynchronously wait until the process terminated
self._watcher.add_child_handler(
self.pid, self._check_exec_result, orig_executable,
cwd, errpipe_data)
else:
if not self._exec_waiter.cancelled():
self._exec_waiter.set_result(None)
self._exec_waiter = None
self._loop = None
self._watcher = None

def _check_exec_result(self, pid, returncode, orig_executable, cwd,
errpipe_data):
try:
super()._check_exec_result(orig_executable, cwd, errpipe_data)
except Exception as exc:
if not self._exec_waiter.cancelled():
self._exec_waiter.set_result(None)
self._exec_waiter = None
self._loop = None
self._watcher = None

def _check_exec_result(self, pid, returncode, orig_executable, cwd,
errpipe_data):
try:
super()._check_exec_result(orig_executable, cwd, errpipe_data)
except Exception as exc:
if not self._exec_waiter.cancelled():
self._exec_waiter.set_exception(exc)
self._cleanup_on_exec_failure()
self._exec_waiter.set_exception(exc)
self._cleanup_on_exec_failure()
else:
_NonBlockingPopen = None


class _UnixSubprocessTransport(base_subprocess.BaseSubprocessTransport):
@coroutine
def _start(self, args, shell, stdin, stdout, stderr, bufsize, **kwargs):
stdin_w = None
if stdin == subprocess.PIPE:
# Use a socket pair for stdin, since not all platforms
# support selecting read events on the write end of a
# socket (which we use in order to detect closing of the
# other end). Notably this is needed on AIX, and works
# just fine on other platforms.
stdin, stdin_w = self._loop._socketpair()

# Mark the write end of the stdin pipe as non-inheritable,
# needed by close_fds=False on Python 3.3 and older
# (Python 3.4 implements the PEP 446, socketpair returns
# non-inheritable sockets)
_set_inheritable(stdin_w.fileno(), False)

with events.get_child_watcher() as watcher:
stdin_w = None
if stdin == subprocess.PIPE:
# Use a socket pair for stdin, since not all platforms
# support selecting read events on the write end of a
# socket (which we use in order to detect closing of the
# other end). Notably this is needed on AIX, and works
# just fine on other platforms.
stdin, stdin_w = self._loop._socketpair()

# Mark the write end of the stdin pipe as non-inheritable,
# needed by close_fds=False on Python 3.3 and older
# (Python 3.4 implements the PEP 446, socketpair returns
# non-inheritable sockets)
_set_inheritable(stdin_w.fileno(), False)
exec_waiter = self._loop.create_future()
try:
self._proc = _NonBlockingPopen(
self._loop, exec_waiter, watcher, args, shell=shell,
stdin=stdin, stdout=stdout, stderr=stderr,
universal_newlines=False, bufsize=bufsize, **kwargs)
yield from exec_waiter
if _NonBlockingPopen:
exec_waiter = self._loop.create_future()
self._proc = _NonBlockingPopen(
self._loop, exec_waiter, watcher, args, shell=shell,
stdin=stdin, stdout=stdout, stderr=stderr,
universal_newlines=False, bufsize=bufsize, **kwargs)
yield from exec_waiter
else:
self._proc = subprocess.Popen(
args, shell=shell, stdin=stdin, stdout=stdout,
stderr=stderr, universal_newlines=False,
bufsize=bufsize, **kwargs)
except:
self._failed_before_start = True
# TODO stdin is probably closed by proc, but what about stdin_w
Expand All @@ -766,9 +778,10 @@ def _start(self, args, shell, stdin, stdout, stderr, bufsize, **kwargs):
else:
watcher.add_child_handler(self._proc.pid,
self._child_watcher_callback)
if stdin_w is not None:
stdin.close()
self._proc.stdin = open(stdin_w.detach(), 'wb', buffering=bufsize)

if stdin_w is not None:
stdin.close()
self._proc.stdin = open(stdin_w.detach(), 'wb', buffering=bufsize)

def _child_watcher_callback(self, pid, returncode):
self._loop.call_soon_threadsafe(self._process_exited, returncode)
Expand Down