| 1 | n/a | __all__ = ['create_subprocess_exec', 'create_subprocess_shell'] |
|---|
| 2 | n/a | |
|---|
| 3 | n/a | import subprocess |
|---|
| 4 | n/a | |
|---|
| 5 | n/a | from . import events |
|---|
| 6 | n/a | from . import protocols |
|---|
| 7 | n/a | from . import streams |
|---|
| 8 | n/a | from . import tasks |
|---|
| 9 | n/a | from .coroutines import coroutine |
|---|
| 10 | n/a | from .log import logger |
|---|
| 11 | n/a | |
|---|
| 12 | n/a | |
|---|
| 13 | n/a | PIPE = subprocess.PIPE |
|---|
| 14 | n/a | STDOUT = subprocess.STDOUT |
|---|
| 15 | n/a | DEVNULL = subprocess.DEVNULL |
|---|
| 16 | n/a | |
|---|
| 17 | n/a | |
|---|
| 18 | n/a | class SubprocessStreamProtocol(streams.FlowControlMixin, |
|---|
| 19 | n/a | protocols.SubprocessProtocol): |
|---|
| 20 | n/a | """Like StreamReaderProtocol, but for a subprocess.""" |
|---|
| 21 | n/a | |
|---|
| 22 | n/a | def __init__(self, limit, loop): |
|---|
| 23 | n/a | super().__init__(loop=loop) |
|---|
| 24 | n/a | self._limit = limit |
|---|
| 25 | n/a | self.stdin = self.stdout = self.stderr = None |
|---|
| 26 | n/a | self._transport = None |
|---|
| 27 | n/a | |
|---|
| 28 | n/a | def __repr__(self): |
|---|
| 29 | n/a | info = [self.__class__.__name__] |
|---|
| 30 | n/a | if self.stdin is not None: |
|---|
| 31 | n/a | info.append('stdin=%r' % self.stdin) |
|---|
| 32 | n/a | if self.stdout is not None: |
|---|
| 33 | n/a | info.append('stdout=%r' % self.stdout) |
|---|
| 34 | n/a | if self.stderr is not None: |
|---|
| 35 | n/a | info.append('stderr=%r' % self.stderr) |
|---|
| 36 | n/a | return '<%s>' % ' '.join(info) |
|---|
| 37 | n/a | |
|---|
| 38 | n/a | def connection_made(self, transport): |
|---|
| 39 | n/a | self._transport = transport |
|---|
| 40 | n/a | |
|---|
| 41 | n/a | stdout_transport = transport.get_pipe_transport(1) |
|---|
| 42 | n/a | if stdout_transport is not None: |
|---|
| 43 | n/a | self.stdout = streams.StreamReader(limit=self._limit, |
|---|
| 44 | n/a | loop=self._loop) |
|---|
| 45 | n/a | self.stdout.set_transport(stdout_transport) |
|---|
| 46 | n/a | |
|---|
| 47 | n/a | stderr_transport = transport.get_pipe_transport(2) |
|---|
| 48 | n/a | if stderr_transport is not None: |
|---|
| 49 | n/a | self.stderr = streams.StreamReader(limit=self._limit, |
|---|
| 50 | n/a | loop=self._loop) |
|---|
| 51 | n/a | self.stderr.set_transport(stderr_transport) |
|---|
| 52 | n/a | |
|---|
| 53 | n/a | stdin_transport = transport.get_pipe_transport(0) |
|---|
| 54 | n/a | if stdin_transport is not None: |
|---|
| 55 | n/a | self.stdin = streams.StreamWriter(stdin_transport, |
|---|
| 56 | n/a | protocol=self, |
|---|
| 57 | n/a | reader=None, |
|---|
| 58 | n/a | loop=self._loop) |
|---|
| 59 | n/a | |
|---|
| 60 | n/a | def pipe_data_received(self, fd, data): |
|---|
| 61 | n/a | if fd == 1: |
|---|
| 62 | n/a | reader = self.stdout |
|---|
| 63 | n/a | elif fd == 2: |
|---|
| 64 | n/a | reader = self.stderr |
|---|
| 65 | n/a | else: |
|---|
| 66 | n/a | reader = None |
|---|
| 67 | n/a | if reader is not None: |
|---|
| 68 | n/a | reader.feed_data(data) |
|---|
| 69 | n/a | |
|---|
| 70 | n/a | def pipe_connection_lost(self, fd, exc): |
|---|
| 71 | n/a | if fd == 0: |
|---|
| 72 | n/a | pipe = self.stdin |
|---|
| 73 | n/a | if pipe is not None: |
|---|
| 74 | n/a | pipe.close() |
|---|
| 75 | n/a | self.connection_lost(exc) |
|---|
| 76 | n/a | return |
|---|
| 77 | n/a | if fd == 1: |
|---|
| 78 | n/a | reader = self.stdout |
|---|
| 79 | n/a | elif fd == 2: |
|---|
| 80 | n/a | reader = self.stderr |
|---|
| 81 | n/a | else: |
|---|
| 82 | n/a | reader = None |
|---|
| 83 | n/a | if reader != None: |
|---|
| 84 | n/a | if exc is None: |
|---|
| 85 | n/a | reader.feed_eof() |
|---|
| 86 | n/a | else: |
|---|
| 87 | n/a | reader.set_exception(exc) |
|---|
| 88 | n/a | |
|---|
| 89 | n/a | def process_exited(self): |
|---|
| 90 | n/a | self._transport.close() |
|---|
| 91 | n/a | self._transport = None |
|---|
| 92 | n/a | |
|---|
| 93 | n/a | |
|---|
| 94 | n/a | class Process: |
|---|
| 95 | n/a | def __init__(self, transport, protocol, loop): |
|---|
| 96 | n/a | self._transport = transport |
|---|
| 97 | n/a | self._protocol = protocol |
|---|
| 98 | n/a | self._loop = loop |
|---|
| 99 | n/a | self.stdin = protocol.stdin |
|---|
| 100 | n/a | self.stdout = protocol.stdout |
|---|
| 101 | n/a | self.stderr = protocol.stderr |
|---|
| 102 | n/a | self.pid = transport.get_pid() |
|---|
| 103 | n/a | |
|---|
| 104 | n/a | def __repr__(self): |
|---|
| 105 | n/a | return '<%s %s>' % (self.__class__.__name__, self.pid) |
|---|
| 106 | n/a | |
|---|
| 107 | n/a | @property |
|---|
| 108 | n/a | def returncode(self): |
|---|
| 109 | n/a | return self._transport.get_returncode() |
|---|
| 110 | n/a | |
|---|
| 111 | n/a | @coroutine |
|---|
| 112 | n/a | def wait(self): |
|---|
| 113 | n/a | """Wait until the process exit and return the process return code. |
|---|
| 114 | n/a | |
|---|
| 115 | n/a | This method is a coroutine.""" |
|---|
| 116 | n/a | return (yield from self._transport._wait()) |
|---|
| 117 | n/a | |
|---|
| 118 | n/a | def send_signal(self, signal): |
|---|
| 119 | n/a | self._transport.send_signal(signal) |
|---|
| 120 | n/a | |
|---|
| 121 | n/a | def terminate(self): |
|---|
| 122 | n/a | self._transport.terminate() |
|---|
| 123 | n/a | |
|---|
| 124 | n/a | def kill(self): |
|---|
| 125 | n/a | self._transport.kill() |
|---|
| 126 | n/a | |
|---|
| 127 | n/a | @coroutine |
|---|
| 128 | n/a | def _feed_stdin(self, input): |
|---|
| 129 | n/a | debug = self._loop.get_debug() |
|---|
| 130 | n/a | self.stdin.write(input) |
|---|
| 131 | n/a | if debug: |
|---|
| 132 | n/a | logger.debug('%r communicate: feed stdin (%s bytes)', |
|---|
| 133 | n/a | self, len(input)) |
|---|
| 134 | n/a | try: |
|---|
| 135 | n/a | yield from self.stdin.drain() |
|---|
| 136 | n/a | except (BrokenPipeError, ConnectionResetError) as exc: |
|---|
| 137 | n/a | # communicate() ignores BrokenPipeError and ConnectionResetError |
|---|
| 138 | n/a | if debug: |
|---|
| 139 | n/a | logger.debug('%r communicate: stdin got %r', self, exc) |
|---|
| 140 | n/a | |
|---|
| 141 | n/a | if debug: |
|---|
| 142 | n/a | logger.debug('%r communicate: close stdin', self) |
|---|
| 143 | n/a | self.stdin.close() |
|---|
| 144 | n/a | |
|---|
| 145 | n/a | @coroutine |
|---|
| 146 | n/a | def _noop(self): |
|---|
| 147 | n/a | return None |
|---|
| 148 | n/a | |
|---|
| 149 | n/a | @coroutine |
|---|
| 150 | n/a | def _read_stream(self, fd): |
|---|
| 151 | n/a | transport = self._transport.get_pipe_transport(fd) |
|---|
| 152 | n/a | if fd == 2: |
|---|
| 153 | n/a | stream = self.stderr |
|---|
| 154 | n/a | else: |
|---|
| 155 | n/a | assert fd == 1 |
|---|
| 156 | n/a | stream = self.stdout |
|---|
| 157 | n/a | if self._loop.get_debug(): |
|---|
| 158 | n/a | name = 'stdout' if fd == 1 else 'stderr' |
|---|
| 159 | n/a | logger.debug('%r communicate: read %s', self, name) |
|---|
| 160 | n/a | output = yield from stream.read() |
|---|
| 161 | n/a | if self._loop.get_debug(): |
|---|
| 162 | n/a | name = 'stdout' if fd == 1 else 'stderr' |
|---|
| 163 | n/a | logger.debug('%r communicate: close %s', self, name) |
|---|
| 164 | n/a | transport.close() |
|---|
| 165 | n/a | return output |
|---|
| 166 | n/a | |
|---|
| 167 | n/a | @coroutine |
|---|
| 168 | n/a | def communicate(self, input=None): |
|---|
| 169 | n/a | if input is not None: |
|---|
| 170 | n/a | stdin = self._feed_stdin(input) |
|---|
| 171 | n/a | else: |
|---|
| 172 | n/a | stdin = self._noop() |
|---|
| 173 | n/a | if self.stdout is not None: |
|---|
| 174 | n/a | stdout = self._read_stream(1) |
|---|
| 175 | n/a | else: |
|---|
| 176 | n/a | stdout = self._noop() |
|---|
| 177 | n/a | if self.stderr is not None: |
|---|
| 178 | n/a | stderr = self._read_stream(2) |
|---|
| 179 | n/a | else: |
|---|
| 180 | n/a | stderr = self._noop() |
|---|
| 181 | n/a | stdin, stdout, stderr = yield from tasks.gather(stdin, stdout, stderr, |
|---|
| 182 | n/a | loop=self._loop) |
|---|
| 183 | n/a | yield from self.wait() |
|---|
| 184 | n/a | return (stdout, stderr) |
|---|
| 185 | n/a | |
|---|
| 186 | n/a | |
|---|
| 187 | n/a | @coroutine |
|---|
| 188 | n/a | def create_subprocess_shell(cmd, stdin=None, stdout=None, stderr=None, |
|---|
| 189 | n/a | loop=None, limit=streams._DEFAULT_LIMIT, **kwds): |
|---|
| 190 | n/a | if loop is None: |
|---|
| 191 | n/a | loop = events.get_event_loop() |
|---|
| 192 | n/a | protocol_factory = lambda: SubprocessStreamProtocol(limit=limit, |
|---|
| 193 | n/a | loop=loop) |
|---|
| 194 | n/a | transport, protocol = yield from loop.subprocess_shell( |
|---|
| 195 | n/a | protocol_factory, |
|---|
| 196 | n/a | cmd, stdin=stdin, stdout=stdout, |
|---|
| 197 | n/a | stderr=stderr, **kwds) |
|---|
| 198 | n/a | return Process(transport, protocol, loop) |
|---|
| 199 | n/a | |
|---|
| 200 | n/a | @coroutine |
|---|
| 201 | n/a | def create_subprocess_exec(program, *args, stdin=None, stdout=None, |
|---|
| 202 | n/a | stderr=None, loop=None, |
|---|
| 203 | n/a | limit=streams._DEFAULT_LIMIT, **kwds): |
|---|
| 204 | n/a | if loop is None: |
|---|
| 205 | n/a | loop = events.get_event_loop() |
|---|
| 206 | n/a | protocol_factory = lambda: SubprocessStreamProtocol(limit=limit, |
|---|
| 207 | n/a | loop=loop) |
|---|
| 208 | n/a | transport, protocol = yield from loop.subprocess_exec( |
|---|
| 209 | n/a | protocol_factory, |
|---|
| 210 | n/a | program, *args, |
|---|
| 211 | n/a | stdin=stdin, stdout=stdout, |
|---|
| 212 | n/a | stderr=stderr, **kwds) |
|---|
| 213 | n/a | return Process(transport, protocol, loop) |
|---|