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