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