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