/
/
1"""FFMpeg related helpers."""
2
3from __future__ import annotations
4
5import asyncio
6import logging
7import re
8import time
9from collections import deque
10from collections.abc import AsyncGenerator, Sequence
11from contextlib import suppress
12from copy import copy
13from dataclasses import dataclass
14from typing import TYPE_CHECKING, Final
15
16from music_assistant_models.enums import ContentType
17from music_assistant_models.errors import AudioError
18from music_assistant_models.helpers import get_global_cache_value, set_global_cache_values
19
20from music_assistant.constants import VERBOSE_LOG_LEVEL
21
22from .dsp import ComplexFilter, ComplexFilterInput
23from .process import AsyncProcess, check_output
24from .util import close_async_generator
25
26if TYPE_CHECKING:
27 from music_assistant_models.media_items import AudioFormat
28
29LOGGER = logging.getLogger("ffmpeg")
30MINIMAL_FFMPEG_VERSION = 7
31CACHE_ATTR_LIBSOXR_PRESENT: Final[str] = "libsoxr_present"
32CACHE_ATTR_FFMPEG_VERSION: Final[str] = "ffmpeg_version"
33DEFAULT_MP3_BIT_RATE: Final[int] = 320
34
35# FFmpeg's mono->stereo rematrix spreads a source at 1/sqrt(2) per channel; this factor
36# restores its original level. _get_channel_conform_filter avoids the same loss on the
37# main decode path by duplicating the channel instead.
38_MONO_WIDEN_COMPENSATION: Final[float] = 2**0.5
39
40# FFmpeg applies these to the single input they precede, not to the command as a whole,
41# so every input we open has to bring its own copy.
42_INPUT_READ_ARGS: Final[list[str]] = [
43 "-protocol_whitelist",
44 "file,hls,http,https,tcp,tls,crypto,pipe,data,fd,rtp,udp,concat",
45 "-probesize",
46 "8096",
47 "-analyzeduration",
48 "500000", # 0.5 seconds should be enough to detect the format
49]
50
51# Regex patterns to extract audio format details from ffmpeg's stderr output.
52# Examples of the lines we parse:
53# Stream #0:0: Audio: mp3, 44100 Hz, stereo, fltp, 320 kb/s
54# Stream #0:0(eng): Audio: aac (LC) (mp4a / 0x6134706D), 44100 Hz, stereo, fltp, 254 kb/s
55# Stream #0:0: Audio: flac, 96000 Hz, stereo, s32 (24 bit)
56# Duration: 00:03:25.78, start: 0.000000, bitrate: 320 kb/s
57_FFMPEG_SAMPLE_RATE_RE: Final = re.compile(r"(\d+) Hz")
58_FFMPEG_BIT_RATE_RE: Final = re.compile(r"(\d+) kb/s")
59_FFMPEG_EXPLICIT_BIT_DEPTH_RE: Final = re.compile(r"\((\d+) bit\)")
60_FFMPEG_SAMPLE_FMT_RE: Final = re.compile(r"\b(u8p?|s16p?|s24p?|s32p?|fltp?|dblp?)\b")
61_FFMPEG_DURATION_RE: Final = re.compile(r"Duration: (\d+):(\d+):(\d+(?:\.\d+)?)")
62
63# Mapping from ffmpeg sample format token to bit depth.
64# Note: planar variants (suffix 'p') describe memory layout only.
65# Floating point formats (flt/fltp/dbl/dblp) are typically the decoder's internal
66# representation for lossy codecs and do not reflect source bit depth, so the
67# caller decides whether to apply them based on the codec.
68_SAMPLE_FMT_BIT_DEPTH: Final[dict[str, int]] = {
69 "u8": 8,
70 "u8p": 8,
71 "s16": 16,
72 "s16p": 16,
73 "s24": 24,
74 "s24p": 24,
75 "s32": 32,
76 "s32p": 32,
77 "flt": 32,
78 "fltp": 32,
79 "dbl": 64,
80 "dblp": 64,
81}
82
83
84@dataclass
85class FFMpegStreamInfo:
86 """Audio format details parsed from an ffmpeg 'Stream #' log line."""
87
88 codec: ContentType
89 sample_rate: int | None = None
90 bit_depth: int | None = None
91 bit_rate: int | None = None
92
93
94class FFMpeg(AsyncProcess):
95 """FFMpeg wrapped as AsyncProcess."""
96
97 def __init__(
98 self,
99 audio_input: AsyncGenerator[bytes] | str | int,
100 input_format: AudioFormat,
101 output_format: AudioFormat,
102 filter_params: Sequence[str | ComplexFilter] | None = None,
103 extra_input_args: list[str] | None = None,
104 extra_output_args: list[str] | None = None,
105 audio_output: str | int = "-",
106 collect_log_history: bool = False,
107 loglevel: str = "info",
108 ) -> None:
109 """Initialize AsyncProcess."""
110 ffmpeg_args = get_ffmpeg_args(
111 input_format=input_format,
112 output_format=output_format,
113 filter_params=filter_params or [],
114 input_path=audio_input if isinstance(audio_input, str) else "-",
115 output_path=audio_output if isinstance(audio_output, str) else "-",
116 extra_input_args=extra_input_args or [],
117 extra_output_args=extra_output_args or [],
118 loglevel=loglevel,
119 )
120 self.audio_input = audio_input
121 self.input_format = input_format
122 self.collect_log_history = collect_log_history
123 self.log_history: deque[str] = deque(maxlen=100)
124 self.concat_error = False # switch to True if concat demuxer fails on MultiPartFiles
125 # Audio format details for the input and output stream as detected from ffmpeg's
126 # own stderr probe output. input_stream_info is also mirrored onto self.input_format
127 # so callers that share the AudioFormat (e.g. streamdetails) pick up the corrected
128 # values; output_stream_info is informational (useful for logging / future UI use).
129 self.input_stream_info: FFMpegStreamInfo | None = None
130 self.output_stream_info: FFMpegStreamInfo | None = None
131 # Source duration in (whole) seconds as detected from the ffmpeg input log line,
132 # or None if not yet parsed / not reported (e.g. live radio streams).
133 self.parsed_duration: int | None = None
134 self._stdin_feeder_task: asyncio.Task[None] | None = None
135 self._stdin_feeder_exception: Exception | None = None
136 self._stderr_reader_task: asyncio.Task[None] | None = None
137 # holds the detached abort-on-corrupt-stream task from _log_reader_task so it
138 # isn't garbage collected mid-flight; not otherwise awaited
139 self._abort_task: asyncio.Task[None] | None = None
140 # ffmpeg emits 'Input #N, ...' and 'Output #N, ...' headers before each block of
141 # 'Stream #' lines; we track which block the next stream line belongs to.
142 # Defaults to "input" so a stray Stream # line before any header still routes there.
143 self._current_log_section: str = "input"
144 stdin: bool | int
145 if audio_input == "-" or isinstance(audio_input, AsyncGenerator):
146 stdin = True
147 else:
148 stdin = audio_input if isinstance(audio_input, int) else False
149 stdout = audio_output if isinstance(audio_output, int) else bool(audio_output == "-")
150 super().__init__(
151 ffmpeg_args,
152 stdin=stdin,
153 stdout=stdout,
154 stderr=True,
155 )
156 self.logger = LOGGER
157
158 @property
159 def stdin_feeder_exception(self) -> Exception | None:
160 """Return the exception raised by the stdin feeder task, if any."""
161 return self._stdin_feeder_exception
162
163 async def start(self) -> None:
164 """Perform Async init of process."""
165 await super().start()
166 if self.proc:
167 self.logger = LOGGER.getChild(str(self.proc.pid))
168 clean_args = []
169 for arg in self._args[1:]:
170 if arg.startswith("http"):
171 clean_args.append("<URL>")
172 elif "/" in arg and "." in arg:
173 clean_args.append("<FILE>")
174 elif arg.startswith("data:application/"):
175 clean_args.append("<DATA>")
176 else:
177 clean_args.append(arg)
178 args_str = " ".join(clean_args)
179 self.logger.log(VERBOSE_LOG_LEVEL, "started with args: %s", args_str)
180 self._stderr_reader_task = asyncio.create_task(self._log_reader_task())
181 if isinstance(self.audio_input, AsyncGenerator):
182 self._stdin_feeder_task = asyncio.create_task(self._feed_stdin())
183
184 async def communicate(
185 self,
186 input: bytes | None = None, # noqa: A002
187 timeout: float | None = None,
188 ) -> tuple[bytes, bytes]:
189 """Override communicate to avoid blocking."""
190 if self._stdin_feeder_task:
191 if not self._stdin_feeder_task.done():
192 self._stdin_feeder_task.cancel()
193 # Always await the task to consume any exception and prevent
194 # "Task exception was never retrieved" errors.
195 try:
196 await self._stdin_feeder_task
197 except asyncio.CancelledError:
198 pass # Expected when we cancel the task
199 except Exception as err:
200 # Log unexpected exceptions from the stdin feeder before suppressing
201 # The audio source may have failed, and we need visibility into this
202 self.logger.warning(
203 "FFMpeg stdin feeder task ended with error: %s",
204 err,
205 )
206 if self._stderr_reader_task:
207 if not self._stderr_reader_task.done():
208 self._stderr_reader_task.cancel()
209 with suppress(asyncio.CancelledError, Exception):
210 await self._stderr_reader_task
211 return await super().communicate(input, timeout)
212
213 async def _log_reader_task(self) -> None:
214 """Read ffmpeg log from stderr."""
215 decode_errors = 0
216 decode_errors_reported = False
217 async for line in self.iter_stderr():
218 if self.collect_log_history:
219 self.log_history.append(line)
220 # ffmpeg logging can be quite verbose, so we only log critical errors
221 # unless verbose logging is enabled
222 if "critical" in line:
223 self.logger.error(line)
224 elif self.logger.isEnabledFor(VERBOSE_LOG_LEVEL):
225 self.logger.log(VERBOSE_LOG_LEVEL, line)
226
227 if "Invalid data found when processing input" in line:
228 decode_errors += 1
229 if decode_errors >= 50 and not decode_errors_reported:
230 # stream is too corrupted to bother decoding further: report once (instead
231 # of promoting every remaining line to ERROR) and abort. close() awaits
232 # this very stderr reader task, and a task awaiting itself raises
233 # RuntimeError, so the abort must run as a detached task rather than be
234 # awaited here.
235 decode_errors_reported = True
236 self.logger.error(
237 "Excessive decode errors (%d+) for this stream; aborting", decode_errors
238 )
239 self._abort_task = asyncio.create_task(self.close())
240
241 # Log reconnection events for radio streams
242 if "Opening" in line or "Reconnect" in line or "reconnect" in line:
243 self.logger.debug("FFmpeg: %s", line)
244
245 if "Error during demuxing" in line:
246 # this can occur if using the concat demuxer for multipart files
247 # and should raise an exception to prevent false progress logging
248 self.concat_error = True
249
250 # Track which ffmpeg block we're currently parsing so the next 'Stream #'
251 # audio line is routed to the correct slot (input vs output).
252 if line.startswith("Input #"):
253 self._current_log_section = "input"
254 elif line.startswith("Output #"):
255 self._current_log_section = "output"
256
257 # Capture the first audio stream line per section. Provider-supplied input
258 # details are often incomplete (e.g. defaults to 44.1/16) or missing for
259 # lossy codecs, so we mirror the parsed input values onto input_format too.
260 if self._current_log_section == "input" and self.input_stream_info is None:
261 if stream_info := parse_ffmpeg_stream_info(line):
262 self.input_stream_info = stream_info
263 self._log_stream_info("input", stream_info)
264 self._apply_input_stream_info(stream_info)
265 elif self._current_log_section == "output" and self.output_stream_info is None:
266 if stream_info := parse_ffmpeg_stream_info(line):
267 self.output_stream_info = stream_info
268 self._log_stream_info("output", stream_info)
269
270 # Source duration is reported separately from the stream info. Useful when
271 # the provider didn't supply one (some podcast feeds report total_time=0).
272 if self.parsed_duration is None:
273 duration = parse_ffmpeg_duration(line)
274 if duration is not None:
275 self.parsed_duration = duration
276 self.logger.debug("Detected input duration: %s seconds", duration)
277 del line
278
279 async def _feed_stdin(self) -> None:
280 """Feed stdin with audio chunks from an AsyncGenerator."""
281 assert not isinstance(self.audio_input, str | int)
282 generator_exhausted = False
283 cancelled = False
284 status = "running"
285 chunk_count = 0
286 self.logger.log(VERBOSE_LOG_LEVEL, "Start reading audio data from source...")
287 try:
288 start = time.time()
289 while True:
290 try:
291 chunk = await anext(self.audio_input)
292 except StopAsyncIteration:
293 generator_exhausted = True
294 break
295 except Exception as err:
296 self._stdin_feeder_exception = err
297 raise
298 chunk_count += 1
299 if self.closed:
300 return
301 await self.write(chunk)
302 except asyncio.CancelledError:
303 status = "cancelled"
304 raise
305 except Exception:
306 status = "aborted with error"
307 raise
308 finally:
309 LOGGER.log(
310 VERBOSE_LOG_LEVEL,
311 "fill_buffer_task: %s (%s chunks received) in in %.2fs",
312 status,
313 chunk_count,
314 time.time() - start,
315 )
316 if not cancelled:
317 await self.write_eof()
318 # we need to ensure that we close the async generator
319 # if we get cancelled otherwise it keeps lingering forever
320 if not generator_exhausted:
321 await close_async_generator(self.audio_input)
322
323 def _apply_input_stream_info(self, info: FFMpegStreamInfo) -> None:
324 """Mirror values from a parsed ffmpeg input stream line onto self.input_format."""
325 # content_type is the container format; only fill it in if the provider didn't
326 # specify one. codec_type is the audio codec ffmpeg detected; only override
327 # if we actually parsed a known codec (don't clobber a provider value with UNKNOWN).
328 if info.codec != ContentType.UNKNOWN:
329 if self.input_format.content_type == ContentType.UNKNOWN:
330 self.input_format.content_type = info.codec
331 self.input_format.codec_type = info.codec
332 if info.sample_rate:
333 self.input_format.sample_rate = info.sample_rate
334 if info.bit_depth:
335 self.input_format.bit_depth = info.bit_depth
336 if info.bit_rate:
337 self.input_format.bit_rate = info.bit_rate
338
339 def _log_stream_info(self, label: str, info: FFMpegStreamInfo) -> None:
340 """Log a parsed FFMpegStreamInfo object at debug level."""
341 self.logger.debug(
342 "Detected %s stream info: codec=%s sample_rate=%s bit_depth=%s bit_rate=%s kb/s",
343 label,
344 info.codec,
345 info.sample_rate,
346 info.bit_depth,
347 info.bit_rate,
348 )
349
350
351def parse_ffmpeg_stream_info(line: str) -> FFMpegStreamInfo | None:
352 """
353 Extract audio format details from an ffmpeg 'Stream #X: Audio: ...' log line.
354
355 :param line: A single ffmpeg stderr log line.
356 :returns: FFMpegStreamInfo when the line describes an audio stream,
357 otherwise None.
358 """
359 if not (line.startswith("Stream #") and ": Audio: " in line):
360 return None
361
362 # the codec name is the first token right after "Audio: ", stripping
363 # any trailing profile annotation like "(LC)" or container suffix
364 codec_part = line.split(": Audio: ", 1)[1].split(" ", 1)[0].split(",", maxsplit=1)[0]
365 codec = ContentType.try_parse(codec_part)
366
367 info = FFMpegStreamInfo(codec=codec)
368 if match := _FFMPEG_SAMPLE_RATE_RE.search(line):
369 info.sample_rate = int(match.group(1))
370 if match := _FFMPEG_BIT_RATE_RE.search(line):
371 info.bit_rate = int(match.group(1))
372 # Bit depth: an explicit "(N bit)" annotation wins (this is how ffmpeg reports
373 # 24-bit FLAC stored in an s32 sample format), otherwise infer from the sample
374 # format token. Lossy codecs report the decoder's internal precision (typically
375 # fltp), so we ignore the sample format token for those.
376 if match := _FFMPEG_EXPLICIT_BIT_DEPTH_RE.search(line):
377 info.bit_depth = int(match.group(1))
378 elif codec.is_lossless() and (match := _FFMPEG_SAMPLE_FMT_RE.search(line)):
379 info.bit_depth = _SAMPLE_FMT_BIT_DEPTH.get(match.group(1))
380
381 return info
382
383
384def parse_ffmpeg_duration(line: str) -> int | None:
385 """
386 Extract the source duration in seconds from an ffmpeg 'Duration: ...' log line.
387
388 :param line: A single ffmpeg stderr log line.
389 :returns: Duration in whole seconds, or None if the line does not contain
390 a parseable duration (e.g. 'Duration: N/A' on live streams).
391 """
392 match = _FFMPEG_DURATION_RE.search(line)
393 if not match:
394 return None
395 hours, minutes, seconds = match.groups()
396 return int(hours) * 3600 + int(minutes) * 60 + int(float(seconds))
397
398
399async def get_ffmpeg_stream(
400 audio_input: AsyncGenerator[bytes] | str,
401 input_format: AudioFormat,
402 output_format: AudioFormat,
403 filter_params: Sequence[str | ComplexFilter] | None = None,
404 chunk_size: int | None = None,
405 extra_input_args: list[str] | None = None,
406 extra_output_args: list[str] | None = None,
407) -> AsyncGenerator[bytes]:
408 """
409 Get the ffmpeg audio stream as async generator.
410
411 Takes care of resampling and/or recoding if needed,
412 according to player preferences.
413 """
414 async with FFMpeg(
415 audio_input=audio_input,
416 input_format=input_format,
417 output_format=output_format,
418 filter_params=filter_params,
419 extra_input_args=extra_input_args,
420 extra_output_args=extra_output_args,
421 collect_log_history=True,
422 ) as ffmpeg_proc:
423 # read final chunks from stdout
424 iterator = ffmpeg_proc.iter_chunked(chunk_size) if chunk_size else ffmpeg_proc.iter_any()
425 async for chunk in iterator:
426 yield chunk
427 # reap the process before trusting returncode: a stream aborted mid-decode (e.g.
428 # excessive decode errors) closes stdout early, which ends the loop above before
429 # the OS process has actually exited, leaving returncode as None if checked directly
430 with suppress(TimeoutError):
431 await ffmpeg_proc.wait_with_timeout(5)
432 if ffmpeg_proc.returncode not in (None, 0) or ffmpeg_proc.concat_error:
433 # unclean exit of ffmpeg - raise error with log tail
434 log_lines = -20 if ffmpeg_proc.concat_error else -5
435 log_tail = "\n" + "\n".join(list(ffmpeg_proc.log_history)[log_lines:])
436 raise AudioError(log_tail)
437 if feeder_exception := ffmpeg_proc.stdin_feeder_exception:
438 raise AudioError("Error while feeding audio to FFmpeg") from feeder_exception
439
440
441async def get_ffmpeg_overlay_stream(
442 audio_input: AsyncGenerator[bytes],
443 overlay_input: str,
444 pcm_format: AudioFormat,
445 overlay_volume: int = 100,
446 chunk_size: int | None = None,
447) -> AsyncGenerator[bytes]:
448 """
449 Mix a looping audio overlay into a PCM audio stream.
450
451 The overlay is looped for the full duration of the main stream and the mixed
452 output has the exact same PCM format and duration as the main input. For a stereo
453 output, a mono overlay mixes in at the same level as an equivalent stereo one. If
454 the overlay input fails mid-stream, the main audio continues unaffected.
455
456 :param audio_input: The main audio stream (raw PCM in ``pcm_format``).
457 :param overlay_input: File path or URL of the overlay audio.
458 :param overlay_volume: Overlay loudness relative to the main audio in
459 percent (100 = equally loud, max 200).
460 :param pcm_format: PCM format of both the main input and the mixed output.
461 :param chunk_size: Optional exact chunk size for the yielded audio.
462 """
463 async with FFMpeg(
464 audio_input=audio_input,
465 # ffmpeg mirrors the metadata it probes from the input onto input_format,
466 # so hand it a copy to keep that mutation off the caller's format.
467 input_format=copy(pcm_format),
468 output_format=pcm_format,
469 filter_params=[_build_overlay_mixer(overlay_input, pcm_format, overlay_volume)],
470 collect_log_history=True,
471 ) as ffmpeg_proc:
472 iterator = ffmpeg_proc.iter_chunked(chunk_size) if chunk_size else ffmpeg_proc.iter_any()
473 async for chunk in iterator:
474 yield chunk
475 # reap the process before trusting returncode: a stream aborted mid-decode (e.g.
476 # excessive decode errors) closes stdout early, which ends the loop above before
477 # the OS process has actually exited, leaving returncode as None if checked directly
478 with suppress(TimeoutError):
479 await ffmpeg_proc.wait_with_timeout(5)
480 if ffmpeg_proc.returncode not in (None, 0):
481 # unclean exit of ffmpeg - raise error with log tail
482 log_tail = "\n" + "\n".join(list(ffmpeg_proc.log_history)[-5:])
483 raise AudioError(log_tail)
484 if feeder_exception := ffmpeg_proc.stdin_feeder_exception:
485 raise AudioError("Error while feeding audio to FFmpeg") from feeder_exception
486
487
488def get_ffmpeg_resample_filter(
489 input_format: AudioFormat,
490 output_format: AudioFormat,
491 filter_params: Sequence[str | ComplexFilter],
492) -> str | None:
493 """
494 Return the resampling and dithering filter required for a format conversion.
495
496 :param input_format: Format entering FFmpeg.
497 :param output_format: Requested FFmpeg output format.
498 :param filter_params: Filters that run before resampling.
499 """
500 if input_format.sample_rate == output_format.sample_rate and not (
501 input_format.bit_depth > 16 and output_format.bit_depth == 16
502 ):
503 return None
504 libsoxr_support = get_global_cache_value(CACHE_ATTR_LIBSOXR_PRESENT)
505 # loudnorm and libsoxr cannot be combined due to https://trac.ffmpeg.org/ticket/11323
506 if libsoxr_support and not any(
507 "loudnorm" in value for value in filter_params if isinstance(value, str)
508 ):
509 resample_filter = "aresample=resampler=soxr:precision=30"
510 else:
511 resample_filter = "aresample=resampler=swr"
512 if input_format.sample_rate != output_format.sample_rate:
513 resample_filter += f":osr={output_format.sample_rate}"
514 if output_format.bit_depth == 16 and input_format.bit_depth > 16:
515 resample_filter += ":osf=s16:dither_method=triangular_hp"
516 return resample_filter
517
518
519def get_ffmpeg_args(
520 input_format: AudioFormat,
521 output_format: AudioFormat,
522 filter_params: Sequence[str | ComplexFilter],
523 input_path: str = "-",
524 output_path: str = "-",
525 extra_input_args: list[str] | None = None,
526 extra_output_args: list[str] | None = None,
527 loglevel: str = "error",
528) -> list[str]:
529 """Collect all args to send to the ffmpeg process."""
530 filter_params = list(filter_params)
531 if extra_input_args is None:
532 extra_input_args = []
533 if extra_output_args is None:
534 extra_output_args = []
535 # the binary plus the options that apply to the command as a whole
536 global_args = [
537 "ffmpeg",
538 "-hide_banner",
539 "-loglevel",
540 loglevel,
541 "-nostats",
542 "-ignore_unknown",
543 ]
544 # collect args for the main input, mirroring how _build_filtergraph_args opens the
545 # extra inputs: the read args lead the group so the caller can still override them
546 input_args = [*_INPUT_READ_ARGS, *extra_input_args]
547 if "-f" not in extra_input_args:
548 # without an input format of their own, the caller leaves the input spec to us
549 if input_path.startswith("http"):
550 # append reconnect options for direct stream from http
551 input_args += [
552 # Reconnect automatically when disconnected before EOF is hit.
553 "-reconnect",
554 "1",
555 # Set the maximum delay in seconds after which to give up reconnecting.
556 "-reconnect_delay_max",
557 "10",
558 # If set then even streamed/non seekable streams will be reconnected on errors.
559 "-reconnect_streamed",
560 "1",
561 # Reconnect automatically in case of TCP/TLS errors during connect.
562 "-reconnect_on_network_error",
563 "0",
564 # A comma separated list of HTTP status codes to reconnect on.
565 # The list can include specific status codes (e.g. 503) or the strings 4xx / 5xx.
566 "-reconnect_on_http_error",
567 "5xx,429",
568 ]
569 if "-post_data" in extra_input_args:
570 # ffmpeg does not include Range headers on POST reconnects, so byte-range
571 # seeking via reconnect is not available. Mark the stream non-seekable so
572 # demuxers do not attempt end-of-file probes (e.g. OGG duration detection)
573 # that would trigger Range-less restarts from byte 0. MA-initiated seeks
574 # still work via -ss decode-and-discard.
575 input_args += ["-seekable", "0"]
576 if input_format.content_type.is_pcm():
577 input_args += [
578 *get_ffmpeg_channel_args(input_format),
579 "-ar",
580 str(input_format.sample_rate),
581 "-acodec",
582 input_format.content_type.name.lower(),
583 "-f",
584 input_format.content_type.value,
585 ]
586 if input_format.codec_type != ContentType.UNKNOWN:
587 input_args += ["-acodec", input_format.codec_type.name.lower()]
588
589 # add input path at the end
590 input_args += ["-i", input_path]
591
592 # collect output args
593 output_args = get_ffmpeg_channel_args(output_format)
594 if output_path.upper() == "NULL":
595 # devnull stream: nothing is encoded here, so there is no channel count to declare
596 output_path = "-"
597 output_args = ["-f", "null"]
598 elif output_format.content_type.is_pcm():
599 # use explicit format identifier for pcm formats
600 output_args += [
601 "-ar",
602 str(output_format.sample_rate),
603 "-acodec",
604 output_format.content_type.name.lower(),
605 "-f",
606 output_format.content_type.value,
607 ]
608 elif output_format.content_type == ContentType.NUT:
609 # passthrough-mode (for creating the cache) using NUT container.
610 # -acodec copy leaves the source untouched, so there is no channel count to declare
611 output_args = [
612 "-vn",
613 "-dn",
614 "-sn",
615 "-acodec",
616 "copy",
617 "-f",
618 "nut",
619 ]
620 elif output_format.content_type == ContentType.AAC:
621 output_args += ["-f", "adts", "-c:a", "aac", "-b:a", "256k"]
622 elif output_format.content_type == ContentType.MP3:
623 output_args += ["-f", "mp3", "-b:a", f"{DEFAULT_MP3_BIT_RATE}k"]
624 elif output_format.content_type == ContentType.WAV:
625 pcm_format = ContentType.from_bit_depth(output_format.bit_depth)
626 output_args += [
627 "-ar",
628 str(output_format.sample_rate),
629 "-acodec",
630 pcm_format.name.lower(),
631 "-f",
632 "wav",
633 ]
634 elif output_format.content_type == ContentType.FLAC:
635 # use level 0 compression for fastest encoding
636 sample_fmt = "s32" if output_format.bit_depth > 16 else "s16"
637 output_args += [
638 "-sample_fmt",
639 sample_fmt,
640 "-ar",
641 str(output_format.sample_rate),
642 "-f",
643 "flac",
644 "-compression_level",
645 "0",
646 ]
647 else:
648 raise RuntimeError("Invalid/unsupported output format specified")
649
650 output_args += extra_output_args # append the extra output args
651 # append (final) output path at the end of the args
652 output_args.append(output_path)
653
654 # runs ahead of the caller's own filters, so channel-aware ones such as the
655 # per-channel preamp see the conformed layout instead of the source layout
656 if channel_filter := _get_channel_conform_filter(input_format.channels, output_format.channels):
657 filter_params = [channel_filter, *filter_params]
658
659 if resample_filter := get_ffmpeg_resample_filter(
660 input_format,
661 output_format,
662 filter_params,
663 ):
664 filter_params.append(resample_filter)
665
666 # a complex fragment brings its own inputs, which must follow the main input
667 filter_input_args, filter_args = (
668 _build_filtergraph_args(filter_params) if filter_params else ([], [])
669 )
670
671 return global_args + input_args + filter_input_args + filter_args + output_args
672
673
674def get_ffmpeg_channel_args(audio_format: AudioFormat) -> list[str]:
675 """
676 Return the FFmpeg channel count/layout arguments for the given audio format.
677
678 The layout is only named for channel counts that map to exactly one layout.
679
680 :param audio_format: Format to describe.
681 """
682 args = ["-ac", str(audio_format.channels)]
683 if layout := _get_channel_layout_name(audio_format.channels):
684 args += ["-channel_layout", layout]
685 return args
686
687
688async def check_ffmpeg_version() -> None:
689 """Check if ffmpeg is present (with libsoxr support)."""
690 # check for FFmpeg presence
691 try:
692 returncode, output = await check_output("ffmpeg", "-version")
693 except FileNotFoundError:
694 raise AudioError(
695 "FFmpeg binary is missing from system. "
696 "Please install ffmpeg on your OS to enable playback."
697 )
698 if returncode != 0:
699 err_msg = "Error determining FFmpeg version on your system."
700 if returncode < 0:
701 # error below 0 is often illegal instruction
702 err_msg += " - Your CPU may be too old to run this version of FFmpeg."
703 err_msg += f" - Additional info: {returncode} {output.decode().strip()}"
704 raise AudioError(err_msg)
705 # parse version number from output
706 try:
707 version = output.decode().split("ffmpeg version ")[1].split(" ")[0].split("-")[0]
708 except IndexError:
709 raise AudioError(
710 "Error determining FFmpeg version on your system."
711 f"Additional info: {returncode} {output.decode().strip()}"
712 )
713 libsoxr_support = "enable-libsoxr" in output.decode()
714 # use globals as in-memory cache
715 await set_global_cache_values(
716 {CACHE_ATTR_LIBSOXR_PRESENT: libsoxr_support, CACHE_ATTR_FFMPEG_VERSION: version}
717 )
718
719 major_version = int("".join(char for char in version.split(".")[0] if not char.isalpha()))
720 if major_version < MINIMAL_FFMPEG_VERSION:
721 raise AudioError(
722 f"FFmpeg version {version} is not supported. "
723 f"Minimal version required is {MINIMAL_FFMPEG_VERSION}."
724 )
725
726 LOGGER.info(
727 "Detected ffmpeg version %s %s",
728 version,
729 "with libsoxr support" if libsoxr_support else "",
730 )
731
732
733def _get_channel_layout_name(channels: int) -> str | None:
734 """
735 Return FFmpeg's layout name for a channel count, or None when it has no unambiguous one.
736
737 :param channels: Number of channels to name.
738 """
739 if channels == 1:
740 return "mono"
741 if channels == 2:
742 return "stereo"
743 # a wider count maps to several possible layouts (5.1 vs 5.1(side), 7.1 vs 7.1(wide), ...)
744 # and a named layout wins over -ac, so naming the wrong one would make FFmpeg misread the
745 # stream as that layout. Left unnamed, it derives the default for the count itself.
746 return None
747
748
749def _get_channel_conform_filter(input_channels: int, output_channels: int) -> str | None:
750 """
751 Return the filter that maps the source onto the output channel count, if one is needed.
752
753 :param input_channels: Channel count entering FFmpeg.
754 :param output_channels: Channel count the output is encoded at.
755 :return: The filter to run before any caller supplied ones, or None when the
756 source already carries the requested channel count.
757 """
758 if input_channels > 2 and output_channels <= 2:
759 # a single channel output needs this fold too, otherwise a mono/left/right pan
760 # would only see the front channels and silently drop the center and surround.
761 # aformat leaves the rematrix to ffmpeg, which picks the correct coefficients
762 # for whatever layout the input turns out to have (and, for an integer output,
763 # scales them to stay clip-safe). A fixed pan expression, naming channels that
764 # a given layout may not even have, can do neither.
765 return "aformat=channel_layouts=stereo"
766 if input_channels == 1 and output_channels > 1:
767 # duplicate rather than leaving the widening to ffmpeg, whose rematrix
768 # spreads the source at 1/sqrt(2) per channel and so costs 3 dB
769 return "pan=stereo|c0=c0|c1=c0"
770 return None
771
772
773def _get_overlay_volume_filter(overlay_volume: int, output_channels: int) -> str:
774 """
775 Return the filter that scales an overlay source to the requested loudness.
776
777 :param overlay_volume: Requested overlay loudness in percent.
778 :param output_channels: Channel count of the mixed output.
779 """
780 gain = overlay_volume / 100
781 if output_channels != 2:
782 # a mono source widened to more than two channels is routed to the centre at full
783 # level, so only a stereo output loses any. No overlay call site is non-stereo today.
784 return f"volume={gain}"
785 # nb_channels is evaluated where this filter sits, ahead of any layout conversion, so it
786 # still reports the source's own count: only a mono source is scaled up, leaving a stereo
787 # one and its image untouched. Comma-free, as a comma would end this filter in the graph.
788 return f"volume={gain}*{_MONO_WIDEN_COMPENSATION}^not(nb_channels-1)"
789
790
791def _build_overlay_mixer(
792 overlay_input: str, pcm_format: AudioFormat, overlay_volume: int
793) -> ComplexFilter:
794 """
795 Build the filter that mixes a looping audio overlay into the main audio.
796
797 :param overlay_input: File path or URL of the overlay audio.
798 :param pcm_format: PCM format of the main input and the mixed output.
799 :param overlay_volume: Overlay loudness relative to the main audio in percent.
800 """
801 input_args = []
802 if overlay_input.startswith("http"):
803 input_args += [
804 "-reconnect",
805 "1",
806 "-reconnect_delay_max",
807 "10",
808 "-reconnect_streamed",
809 "1",
810 ]
811 input_args += ["-stream_loop", "-1"]
812 # conform the overlay to the main stream's layout so amix sees two matching inputs;
813 # an unnameable count is left to FFmpeg's own negotiation
814 layout = _get_channel_layout_name(pcm_format.channels)
815 conform_filter = f",aformat=channel_layouts={layout}" if layout else ""
816 return ComplexFilter(
817 # the main audio is amix's first input, so duration=first follows its length;
818 # normalize=0 keeps the original levels (no averaging)
819 body="amix=inputs=2:duration=first:normalize=0",
820 inputs=[
821 ComplexFilterInput(
822 path=overlay_input,
823 # silenceremove strips a near-silent intro from the overlay source (e.g. a
824 # soft fade-in) so it becomes audible right away; it is a no-op for sources
825 # that already start at full level. It runs before volume so detection is
826 # based on the source's own levels rather than the scaled output. volume
827 # in turn has to stay ahead of the resample and conform steps, which
828 # replace the source's own channel count with the output's.
829 filters=(
830 f"silenceremove=start_periods=1:start_threshold=-40dB,"
831 f"{_get_overlay_volume_filter(overlay_volume, pcm_format.channels)},"
832 f"aresample={pcm_format.sample_rate}"
833 f"{conform_filter}"
834 ),
835 input_args=input_args,
836 )
837 ],
838 )
839
840
841def _build_filtergraph_args(
842 filter_params: list[str | ComplexFilter],
843) -> tuple[list[str], list[str]]:
844 """
845 Render a DSP filter chain to FFmpeg command-line arguments.
846
847 :param filter_params: Ordered chain of plain filter strings and/or complex
848 fragments that need extra audio inputs.
849 :return: Extra input arguments to append after the main input, and the
850 filter arguments themselves.
851 """
852 if not any(isinstance(item, ComplexFilter) for item in filter_params):
853 simple = [item for item in filter_params if isinstance(item, str) and item]
854 return [], (["-af", ",".join(simple)] if simple else [])
855
856 input_args: list[str] = []
857 parts: list[str] = []
858 pending: list[str] = []
859 current = "0:a"
860 counter = 0
861 # the main input is 0, so extra inputs are numbered from 1 in the order added
862 next_input = 1
863
864 def next_label() -> str:
865 nonlocal counter
866 counter += 1
867 return f"dsp{counter}"
868
869 def flush_pending() -> None:
870 nonlocal current
871 if not pending:
872 return
873 label = next_label()
874 parts.append(f"[{current}]{','.join(pending)}[{label}]")
875 current = label
876 pending.clear()
877
878 for item in filter_params:
879 if isinstance(item, str):
880 if item:
881 pending.append(item)
882 continue
883 # a complex fragment closes the current simple run, adds its own inputs to
884 # the command, then consumes the main pad plus those inputs
885 flush_pending()
886 source_labels: list[str] = []
887 for extra_input in item.inputs:
888 input_args += [*_INPUT_READ_ARGS, *extra_input.input_args, "-i", extra_input.path]
889 source = f"{next_input}:a"
890 next_input += 1
891 if extra_input.filters:
892 label = next_label()
893 parts.append(f"[{source}]{extra_input.filters}[{label}]")
894 source = label
895 source_labels.append(source)
896 label = next_label()
897 inputs = f"[{current}]" + "".join(f"[{sl}]" for sl in source_labels)
898 parts.append(f"{inputs}{item.body}[{label}]")
899 current = label
900 flush_pending()
901
902 return input_args, ["-filter_complex", ";".join(parts), "-map", f"[{current}]"]
903