/
/
1"""
2AsyncProcess.
3
4Wrapper around asyncio subprocess to help with using pipe streams and
5taking care of properly closing the process in case of exit (on both success and failures),
6without deadlocking.
7"""
8
9from __future__ import annotations
10
11import asyncio
12import logging
13import os
14
15# if TYPE_CHECKING:
16from collections.abc import AsyncGenerator, AsyncIterator
17from contextlib import asynccontextmanager, suppress
18from signal import SIGINT
19from types import TracebackType
20from typing import Self
21
22from music_assistant.constants import MASS_LOGGER_NAME, VERBOSE_LOG_LEVEL
23
24LOGGER = logging.getLogger(f"{MASS_LOGGER_NAME}.helpers.process")
25
26DEFAULT_CHUNKSIZE = 64000
27
28
29def get_subprocess_env(env: dict[str, str] | None = None) -> dict[str, str]:
30 """Get environment for subprocess, stripping LD_PRELOAD to avoid jemalloc warnings."""
31 result = dict(os.environ)
32 result.pop("LD_PRELOAD", None)
33 if env:
34 result.update(env)
35 return result
36
37
38class AsyncProcess:
39 """
40 AsyncProcess.
41
42 Wrapper around asyncio subprocess to help with using pipe streams and
43 taking care of properly closing the process in case of exit (on both success and failures),
44 without deadlocking.
45 """
46
47 _stdin_feeder_task: asyncio.Task[None] | None = None # used for ffmpeg
48 _stderr_reader_task: asyncio.Task[None] | None = None # used for ffmpeg
49
50 def __init__(
51 self,
52 args: list[str],
53 stdin: bool | int | None = None,
54 stdout: bool | int | None = None,
55 stderr: bool | int | None = False,
56 name: str | None = None,
57 env: dict[str, str] | None = None,
58 ) -> None:
59 """
60 Initialize AsyncProcess.
61
62 :param args: Command and arguments to execute.
63 :param stdin: Stdin configuration (True for PIPE, False for None, or custom).
64 :param stdout: Stdout configuration (True for PIPE, False for None, or custom).
65 :param stderr: Stderr configuration (True for PIPE, False for DEVNULL, or custom).
66 :param name: Process name for logging.
67 :param env: Environment variables for the subprocess (None inherits parent env).
68 """
69 self.proc: asyncio.subprocess.Process | None = None
70 if name is None:
71 name = args[0].split(os.sep)[-1]
72 self.name = name
73 self.logger = LOGGER.getChild(name)
74 self._args = args
75 self._stdin = None if stdin is False else stdin
76 self._stdout = None if stdout is False else stdout
77 self._stderr = asyncio.subprocess.DEVNULL if stderr is False else stderr
78 self._env = get_subprocess_env(env)
79 self._stderr_lock = asyncio.Lock()
80 self._stdout_lock = asyncio.Lock()
81 self._stdin_lock = asyncio.Lock()
82 self._close_called = False
83 self._returncode: int | None = None
84
85 @property
86 def closed(self) -> bool:
87 """Return if the process was closed."""
88 return self._close_called or self.returncode is not None
89
90 @property
91 def returncode(self) -> int | None:
92 """Return the erturncode of the process."""
93 if self._returncode is not None:
94 return self._returncode
95 if self.proc is None:
96 return None
97 if (ret_code := self.proc.returncode) is not None:
98 self._returncode = ret_code
99 return ret_code
100
101 async def __aenter__(self) -> Self:
102 """Enter context manager."""
103 await self.start()
104 return self
105
106 async def __aexit__(
107 self,
108 exc_type: type[BaseException] | None,
109 exc_val: BaseException | None,
110 exc_tb: TracebackType | None,
111 ) -> bool | None:
112 """Exit context manager."""
113 # make sure we close and cleanup the process
114 await self.close()
115 self._returncode = self.returncode
116 return None
117
118 async def start(self) -> None:
119 """Perform Async init of process."""
120 self.proc = await asyncio.create_subprocess_exec(
121 *self._args,
122 stdin=asyncio.subprocess.PIPE if self._stdin is True else self._stdin,
123 stdout=asyncio.subprocess.PIPE if self._stdout is True else self._stdout,
124 stderr=asyncio.subprocess.PIPE if self._stderr is True else self._stderr,
125 env=self._env,
126 bufsize=0,
127 )
128 self.logger.log(
129 VERBOSE_LOG_LEVEL, "Process %s started with PID %s", self.name, self.proc.pid
130 )
131
132 async def iter_chunked(self, n: int = DEFAULT_CHUNKSIZE) -> AsyncGenerator[bytes]:
133 """Yield chunks of n size from the process stdout."""
134 while True:
135 chunk = await self.readexactly(n)
136 if len(chunk) == 0:
137 break
138 yield chunk
139
140 async def iter_any(self, n: int = DEFAULT_CHUNKSIZE) -> AsyncGenerator[bytes]:
141 """Yield chunks as they come in from process stdout."""
142 while True:
143 chunk = await self.read(n)
144 if len(chunk) == 0:
145 break
146 yield chunk
147
148 async def readexactly(self, n: int) -> bytes:
149 """Read exactly n bytes from the process stdout (or less if eof)."""
150 if self._close_called:
151 return b""
152 assert self.proc is not None # for type checking
153 assert self.proc.stdout is not None # for type checking
154 async with self._stdout_lock:
155 try:
156 return await self.proc.stdout.readexactly(n)
157 except asyncio.IncompleteReadError as err:
158 return err.partial
159
160 async def read(self, n: int) -> bytes:
161 """
162 Read up to n bytes from the stdout stream.
163
164 If n is positive, this function try to read n bytes,
165 and may return less or equal bytes than requested, but at least one byte.
166 If EOF was received before any byte is read, this function returns empty byte object.
167 """
168 if self._close_called:
169 return b""
170 assert self.proc is not None # for type checking
171 assert self.proc.stdout is not None # for type checking
172 async with self._stdout_lock:
173 return await self.proc.stdout.read(n)
174
175 async def write(self, data: bytes) -> None:
176 """Write data to process stdin."""
177 if self._close_called or self.proc is None:
178 return
179 if self.proc.stdin is None:
180 return
181 async with self._stdin_lock:
182 self.proc.stdin.write(data)
183 await self.proc.stdin.drain()
184
185 @asynccontextmanager
186 async def stdin_quiesced(self, timeout: float = 5.0) -> AsyncIterator[bool]:
187 """
188 Hold stdin quiet for a block, with what was already written seen through to the pipe.
189
190 :meth:`write` only waits while the transport is paused, which it is only
191 above the high-water mark, so it returns with up to that much still queued
192 locally (64 KiB by default). This first sees those bytes through to the
193 kernel pipe -- as far as it can guarantee; whether the process has read
194 them is its own business -- and then keeps the write lock for the body, so
195 no :meth:`write` or :meth:`write_eof` can interleave. For a caller telling
196 the process something about the bytes it has been handed -- out of band,
197 and in a sequence the process must not see a write inside -- that turns
198 "we happen to have stopped writing" into something the block enforces.
199
200 Yields True when stdin was emptied, False when it could not be: the
201 process is then still owed bytes, so a caller whose message depends on it
202 having received everything must give up rather than send it.
203
204 :param timeout: Seconds to wait for the buffer to empty.
205 """
206 if self._close_called or self.proc is None or self.proc.stdin is None:
207 yield True
208 return
209 async with self._stdin_lock:
210 yield await self._drain_stdin_locked(timeout)
211
212 async def write_eof(self) -> None:
213 """Write end of file to to process stdin."""
214 if self._close_called or self.proc is None:
215 return
216 if self.proc.stdin is None:
217 return
218 async with self._stdin_lock:
219 try:
220 if self.proc.stdin.can_write_eof():
221 self.proc.stdin.write_eof()
222 await self.proc.stdin.drain()
223 except (
224 AttributeError,
225 AssertionError,
226 BrokenPipeError,
227 RuntimeError,
228 ConnectionResetError,
229 ):
230 # already exited, race condition
231 pass
232
233 async def read_stderr(self) -> bytes:
234 """Read line from stderr."""
235 if self.returncode is not None:
236 return b""
237 assert self.proc is not None # for type checking
238 assert self.proc.stderr is not None # for type checking
239 async with self._stderr_lock:
240 try:
241 return await self.proc.stderr.readline()
242 except ValueError as err:
243 # we're waiting for a line (separator found), but the line was too big
244 # this may happen with ffmpeg during a long (radio) stream where progress
245 # gets outputted to the stderr but no newline
246 # https://stackoverflow.com/questions/55457370/how-to-avoid-valueerror-separator-is-not-found-and-chunk-exceed-the-limit
247 # NOTE: this consumes the line that was too big
248 if "chunk exceed the limit" in str(err):
249 return await self.proc.stderr.readline()
250 # raise for all other (value) errors
251 raise
252
253 async def iter_stderr(self) -> AsyncGenerator[str]:
254 """Iterate lines from the stderr stream as string."""
255 line: str | bytes
256 while True:
257 line = await self.read_stderr()
258 if line == b"":
259 break
260 line = line.decode("utf-8", errors="ignore").strip()
261 if not line:
262 continue
263 yield line
264
265 async def communicate(
266 self,
267 input: bytes | None = None, # noqa: A002
268 timeout: float | None = None,
269 ) -> tuple[bytes, bytes]:
270 """Communicate with the process and return stdout and stderr."""
271 if self.closed:
272 raise RuntimeError("communicate called while process already done")
273 # abort existing readers on stderr/stdout first before we send communicate
274 await self._stderr_lock.acquire()
275 await self._stdout_lock.acquire()
276 assert self.proc is not None # for type checking
277 stdout, stderr = await asyncio.wait_for(self.proc.communicate(input), timeout)
278 return (stdout, stderr)
279
280 async def close(self) -> None:
281 """Close/terminate the process and wait for exit."""
282 self._close_called = True
283 if not self.proc:
284 return
285
286 # cancel existing stdin feeder task if any
287 if self._stdin_feeder_task:
288 if not self._stdin_feeder_task.done():
289 self._stdin_feeder_task.cancel()
290 # Always await the task to consume any exception and prevent
291 # "Task exception was never retrieved" errors.
292 try:
293 await self._stdin_feeder_task
294 except asyncio.CancelledError:
295 pass # Expected when we cancel the task
296 except Exception as err:
297 # Log unexpected exceptions from the stdin feeder before suppressing
298 LOGGER.warning(
299 "Process stdin feeder task ended with error: %s",
300 err,
301 )
302
303 # close stdin to signal we're done sending data
304 with suppress(TimeoutError, asyncio.CancelledError):
305 await asyncio.wait_for(self._stdin_lock.acquire(), 5)
306 if self.proc.stdin and not self.proc.stdin.is_closing():
307 self.proc.stdin.close()
308 elif not self.proc.stdin and self.proc.returncode is None:
309 # the process may exit between the returncode check and the signal; guard the
310 # race the same way the SIGKILL delivery below does
311 with suppress(ProcessLookupError, OSError):
312 self.proc.send_signal(SIGINT)
313
314 # ensure we have no more readers active and stdout is drained
315 with suppress(TimeoutError, asyncio.CancelledError):
316 await asyncio.wait_for(self._stdout_lock.acquire(), 5)
317 if self.proc.stdout and not self.proc.stdout.at_eof():
318 with suppress(Exception):
319 await self.proc.stdout.read(-1)
320 # if we have a stderr task active, allow it to finish
321 if self._stderr_reader_task:
322 with suppress(TimeoutError, asyncio.CancelledError):
323 await asyncio.wait_for(self._stderr_reader_task, 5)
324 elif self.proc.stderr and not self.proc.stderr.at_eof():
325 with suppress(TimeoutError, asyncio.CancelledError):
326 await asyncio.wait_for(self._stderr_lock.acquire(), 5)
327 # drain stderr
328 with suppress(Exception):
329 await self.proc.stderr.read(-1)
330
331 # make sure the process is really cleaned up.
332 # especially with pipes this can cause deadlocks if not properly guarded
333 # we need to ensure stdout and stderr are flushed and stdin closed
334 pid = self.proc.pid
335 terminate_attempts = 0
336 while self.returncode is None:
337 try:
338 # use communicate to flush all pipe buffers
339 await asyncio.wait_for(self.proc.communicate(), 2)
340 except TimeoutError:
341 terminate_attempts += 1
342 self.logger.debug(
343 "Process %s with PID %s did not stop in time (attempt %d). Sending SIGKILL...",
344 self.name,
345 pid,
346 terminate_attempts,
347 )
348 # Use os.kill for more direct signal delivery
349 with suppress(ProcessLookupError, OSError):
350 os.kill(pid, 9) # SIGKILL = 9
351 # Give up after 5 attempts - process may be zombie
352 if terminate_attempts >= 5:
353 self.logger.warning(
354 "Process %s (PID %s) did not terminate after %d SIGKILL attempts",
355 self.name,
356 pid,
357 terminate_attempts,
358 )
359 break
360 self.logger.log(
361 VERBOSE_LOG_LEVEL,
362 "Process %s with PID %s stopped with returncode %s",
363 self.name,
364 self.proc.pid,
365 self.returncode,
366 )
367
368 async def kill(self) -> None:
369 """
370 Immediately kill the process with SIGKILL.
371
372 Use this for forceful termination when the process doesn't respond to
373 normal termination signals. Unlike close(), this doesn't attempt graceful
374 shutdown - it immediately sends SIGKILL.
375 """
376 self._close_called = True
377 if not self.proc or self.returncode is not None:
378 return
379
380 pid = self.proc.pid
381
382 # Cancel stdin feeder task if any
383 if self._stdin_feeder_task and not self._stdin_feeder_task.done():
384 self._stdin_feeder_task.cancel()
385 with suppress(asyncio.CancelledError, Exception):
386 await self._stdin_feeder_task
387
388 # Cancel stderr reader task if any
389 if self._stderr_reader_task and not self._stderr_reader_task.done():
390 self._stderr_reader_task.cancel()
391 with suppress(asyncio.CancelledError, Exception):
392 await self._stderr_reader_task
393
394 # Close stdin to signal we're done sending data
395 # Note: Don't manually call feed_eof() on stdout/stderr - this causes
396 # "feed_data after feed_eof" assertion errors when the subprocess transport
397 # still has buffered data to deliver. Let the process termination naturally
398 # close the streams.
399 if self.proc.stdin and not self.proc.stdin.is_closing():
400 self.proc.stdin.close()
401
402 # Send SIGKILL immediately using os.kill for more direct signal delivery
403 self.logger.debug("Killing process %s with PID %s", self.name, pid)
404 with suppress(ProcessLookupError, OSError):
405 os.kill(pid, 9) # SIGKILL = 9
406
407 # Wait for process to actually terminate
408 try:
409 await asyncio.wait_for(self.proc.wait(), 2)
410 except TimeoutError:
411 # Try one more time with os.kill
412 with suppress(ProcessLookupError, OSError):
413 os.kill(pid, 9)
414 try:
415 await asyncio.wait_for(self.proc.wait(), 2)
416 except TimeoutError:
417 self.logger.warning(
418 "Process %s with PID %s did not terminate after SIGKILL - may be zombie",
419 self.name,
420 pid,
421 )
422
423 self.logger.log(
424 VERBOSE_LOG_LEVEL,
425 "Process %s with PID %s killed with returncode %s",
426 self.name,
427 pid,
428 self.returncode,
429 )
430
431 async def wait(self) -> int:
432 """Wait for the process and return the returncode."""
433 if self._returncode is None:
434 assert self.proc is not None
435 self._returncode = await self.proc.wait()
436 return self._returncode
437
438 async def wait_with_timeout(self, timeout: int) -> int:
439 """Wait for the process and return the returncode with a timeout."""
440 return await asyncio.wait_for(self.wait(), timeout)
441
442 def attach_stderr_reader(self, task: asyncio.Task[None]) -> None:
443 """Attach a stderr reader task to this process."""
444 self._stderr_reader_task = task
445
446 async def _drain_stdin_locked(self, timeout: float) -> bool:
447 """
448 Empty the stdin write buffer, with the write lock already held.
449
450 :param timeout: Seconds to wait for the buffer to empty.
451 :return: True once the buffer is empty, False when the wait timed out.
452 """
453 assert self.proc is not None # for type checking
454 assert self.proc.stdin is not None # for type checking
455 transport = self.proc.stdin.transport
456 low, high = transport.get_write_buffer_limits()
457 try:
458 # Pausing the transport at a zero high-water mark is what makes
459 # drain() resolve only once the buffer is completely empty: it
460 # otherwise resolves as soon as the transport is not paused.
461 transport.set_write_buffer_limits(high=0)
462 await asyncio.wait_for(self.proc.stdin.drain(), timeout)
463 except TimeoutError:
464 return False
465 except BrokenPipeError, RuntimeError, ConnectionResetError:
466 # already exited, race condition: nothing is left to arrive
467 return True
468 finally:
469 # Restore what this process was configured with rather than the
470 # asyncio defaults a bare call would reinstate.
471 with suppress(RuntimeError):
472 transport.set_write_buffer_limits(high=high, low=low)
473 return True
474
475
476async def check_output(
477 *args: str, env: dict[str, str] | None = None, timeout: float | None = None
478) -> tuple[int, bytes]:
479 """
480 Run subprocess and return returncode and output.
481
482 :param env: Optional environment overrides for the subprocess.
483 :param timeout: Maximum seconds to wait for the process to exit. On expiry the
484 process is killed and TimeoutError is raised; None (default) waits forever.
485 """
486 proc = await asyncio.create_subprocess_exec(
487 *args,
488 stderr=asyncio.subprocess.STDOUT,
489 stdout=asyncio.subprocess.PIPE,
490 env=get_subprocess_env(env),
491 )
492 try:
493 async with asyncio.timeout(timeout):
494 stdout, _ = await proc.communicate()
495 except TimeoutError:
496 proc.kill()
497 with suppress(ProcessLookupError):
498 await proc.wait()
499 raise
500 assert proc.returncode is not None # for type checking
501 return (proc.returncode, stdout)
502
503
504async def communicate(
505 args: list[str],
506 input: bytes | None = None, # noqa: A002
507) -> tuple[int, bytes, bytes]:
508 """Communicate with subprocess and return returncode, stdout and stderr output."""
509 proc = await asyncio.create_subprocess_exec(
510 *args,
511 stderr=asyncio.subprocess.PIPE,
512 stdout=asyncio.subprocess.PIPE,
513 stdin=asyncio.subprocess.PIPE if input is not None else None,
514 env=get_subprocess_env(),
515 )
516 stdout, stderr = await proc.communicate(input)
517 assert proc.returncode is not None # for type checking
518 return (proc.returncode, stdout, stderr)
519