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