Build and publish container / build (pull_request) Successful in 6m57s
Copies of already-MP3 sources were written straight to their destination while encodes went via a temporary file and a rename. A copy interrupted by a full disk, a killed container or an I/O error therefore left a truncated MP3 in the mirror -- and because shutil.copy2 reproduces the source's mtime along with its bytes, staleness detection would read that fragment as up to date and never replace it. The damage is silent and permanent until someone plays the track. Give copy the same temporary-file-and-rename path encode already uses, so the destination either has the whole file or has nothing.
590 lines
20 KiB
Python
590 lines
20 KiB
Python
"""Maintain a lossy MP3 mirror of a lossless music library.
|
|
|
|
Walks a source library and reproduces it, path for path, as MP3 in a separate
|
|
directory tree: FLAC in, MP3 out, same relative layout, tags and cover art
|
|
carried across. Sources that are already MP3 are copied rather than re-encoded.
|
|
|
|
The mirror is derived state. It is only ever written to, never read as a
|
|
source of truth, so it can be deleted and rebuilt at any time. Nothing here
|
|
writes to the source library.
|
|
|
|
Staleness is tracked by modification time: an encoded file is given its
|
|
source's mtime, so a file is out of date exactly when the two differ. That
|
|
makes runs idempotent without a database to keep in step.
|
|
"""
|
|
|
|
import argparse
|
|
import concurrent.futures
|
|
import fcntl
|
|
import functools
|
|
import logging
|
|
import os
|
|
import re
|
|
import shutil
|
|
import signal
|
|
import subprocess
|
|
import sys
|
|
import tempfile
|
|
import time
|
|
from dataclasses import dataclass
|
|
from pathlib import Path
|
|
|
|
logger = logging.getLogger("music-mirror")
|
|
|
|
# Sources that are transcoded. Anything ffmpeg can decode works; this list
|
|
# decides what the walker picks up in the first place.
|
|
SOURCE_EXTENSIONS = {
|
|
".flac",
|
|
".wav",
|
|
".aif",
|
|
".aiff",
|
|
".ape",
|
|
".wv",
|
|
".m4a",
|
|
".alac",
|
|
".ogg",
|
|
".opus",
|
|
".wma",
|
|
}
|
|
|
|
# Already MP3: copied through. Re-encoding lossy audio to lossy audio costs
|
|
# quality for nothing.
|
|
COPY_EXTENSIONS = {".mp3"}
|
|
|
|
# Best first. Used only to settle which source wins when two of them want the
|
|
# same mirror path; see plan().
|
|
SOURCE_PRIORITY = [
|
|
".flac",
|
|
".wav",
|
|
".aif",
|
|
".aiff",
|
|
".ape",
|
|
".wv",
|
|
".alac",
|
|
".m4a",
|
|
".ogg",
|
|
".opus",
|
|
".wma",
|
|
".mp3",
|
|
]
|
|
|
|
# Looked for in the source directory when a file has no embedded picture.
|
|
COVER_NAMES = ("cover.jpg", "folder.jpg", "front.jpg", "cover.png", "folder.png")
|
|
|
|
# Files the mirror is allowed to contain, and therefore allowed to delete.
|
|
MIRROR_SUFFIX = ".mp3"
|
|
|
|
# Filesystems disagree about mtime precision; SMB in particular rounds.
|
|
MTIME_TOLERANCE_SECONDS = 2
|
|
|
|
# The mirror exists to be read back by something else -- an SMB share, another
|
|
# account on the box -- so everything written into it has to be group-readable.
|
|
# Neither writer manages that unaided: tempfile.mkstemp forces 0600 whatever the
|
|
# umask, and shutil.copy2 carries the source file's mode across from a library
|
|
# that may be tighter still. Directories need the execute bit too, or the group
|
|
# cannot enter them to reach the readable files inside.
|
|
GROUP_READ = 0o040
|
|
GROUP_ENTER = 0o050
|
|
|
|
|
|
@dataclass
|
|
class Result:
|
|
"""Outcome of processing one file."""
|
|
|
|
action: str # encoded | copied | skipped | failed
|
|
path: Path
|
|
error: str = ""
|
|
|
|
|
|
def parse_quality(quality):
|
|
"""Return the ffmpeg arguments for a quality setting.
|
|
|
|
Accepts LAME VBR levels (``V0``..``V9``) or a constant bitrate in kbps
|
|
(``256``). VBR is the better trade at a given average bitrate; CBR is for
|
|
when a fixed size matters more.
|
|
"""
|
|
text = str(quality).strip().lower()
|
|
if re.fullmatch(r"v[0-9]", text):
|
|
return ["-q:a", text[1:]]
|
|
if re.fullmatch(r"[0-9]{2,3}", text):
|
|
return ["-b:a", f"{text}k"]
|
|
raise ValueError(f"unrecognised quality {quality!r}: expected V0-V9 or a bitrate like 256")
|
|
|
|
|
|
def parse_interval(interval):
|
|
"""Return seconds for an interval such as ``30m``, ``6h`` or ``90``."""
|
|
text = str(interval).strip().lower()
|
|
match = re.fullmatch(r"([0-9]+)([smhd]?)", text)
|
|
if not match:
|
|
raise ValueError(f"unrecognised interval {interval!r}: expected e.g. 45m, 6h, 1d")
|
|
value = int(match.group(1))
|
|
return value * {"": 1, "s": 1, "m": 60, "h": 3600, "d": 86400}[match.group(2)]
|
|
|
|
|
|
def mirror_path_for(source, source_root, mirror_root):
|
|
"""Return the mirror path corresponding to a source file."""
|
|
return (mirror_root / source.relative_to(source_root)).with_suffix(MIRROR_SUFFIX)
|
|
|
|
|
|
def is_current(source, mirror):
|
|
"""Return whether the mirror file is up to date with its source."""
|
|
if not mirror.exists():
|
|
return False
|
|
return abs(source.stat().st_mtime - mirror.stat().st_mtime) <= MTIME_TOLERANCE_SECONDS
|
|
|
|
|
|
def make_group_readable(path):
|
|
"""Add the group-read bit to a mirror file, leaving the rest of the mode alone."""
|
|
mode = path.stat().st_mode
|
|
if not mode & GROUP_READ:
|
|
path.chmod(mode | GROUP_READ)
|
|
|
|
|
|
@functools.lru_cache(maxsize=4096)
|
|
def find_cover(directory):
|
|
"""Return an external cover image for a directory, if one is present.
|
|
|
|
Cached because an album's tracks all ask the same question, and the answer
|
|
costs one stat per candidate name.
|
|
"""
|
|
for name in COVER_NAMES:
|
|
candidate = directory / name
|
|
if candidate.is_file():
|
|
return candidate
|
|
return None
|
|
|
|
|
|
def has_embedded_picture(source):
|
|
"""Return whether the source carries its own cover art."""
|
|
try:
|
|
probe = subprocess.run(
|
|
[
|
|
"ffprobe",
|
|
"-v",
|
|
"error",
|
|
"-select_streams",
|
|
"v",
|
|
"-show_entries",
|
|
"stream=index",
|
|
"-of",
|
|
"csv=p=0",
|
|
str(source),
|
|
],
|
|
capture_output=True,
|
|
text=True,
|
|
check=True,
|
|
)
|
|
except (subprocess.CalledProcessError, FileNotFoundError):
|
|
return False
|
|
return bool(probe.stdout.strip())
|
|
|
|
|
|
def build_command(source, destination, quality_args, cover):
|
|
"""Return the ffmpeg command that encodes one file."""
|
|
command = ["ffmpeg", "-nostdin", "-hide_banner", "-loglevel", "error", "-y", "-i", str(source)]
|
|
|
|
if cover is not None:
|
|
command += ["-i", str(cover), "-map", "1:v:0"]
|
|
else:
|
|
# Optional: the source may have no picture stream at all.
|
|
command += ["-map", "0:v:0?"]
|
|
|
|
command += [
|
|
"-map",
|
|
"0:a:0",
|
|
"-map_metadata",
|
|
"0",
|
|
"-c:a",
|
|
"libmp3lame",
|
|
*quality_args,
|
|
"-c:v",
|
|
"copy",
|
|
"-disposition:v",
|
|
"attached_pic",
|
|
# ID3v2.3 is the widest-compatibility tag version, and what the iPod
|
|
# firmware is happiest with; the v1 tag costs 128 bytes.
|
|
"-id3v2_version",
|
|
"3",
|
|
"-write_id3v1",
|
|
"1",
|
|
# Stated rather than inferred: the destination is a temporary file
|
|
# whose suffix ffmpeg would not recognise.
|
|
"-f",
|
|
"mp3",
|
|
str(destination),
|
|
]
|
|
return command
|
|
|
|
|
|
def encode(source, mirror, quality_args, dry_run):
|
|
"""Encode one source file into the mirror, atomically."""
|
|
if dry_run:
|
|
logger.info("would encode %s", source)
|
|
return Result("encoded", mirror)
|
|
|
|
mirror.parent.mkdir(parents=True, exist_ok=True)
|
|
# Probing costs an ffprobe process per file, so only ask when the answer
|
|
# can change the command. With no cover file beside the track, `-map
|
|
# 0:v:0?` carries embedded art if there is any and shrugs if there is not.
|
|
cover = find_cover(source.parent)
|
|
if cover is not None and has_embedded_picture(source):
|
|
cover = None
|
|
|
|
# Read the source's mtime before encoding, not after. If the file is still
|
|
# being written -- a Lidarr import landing mid-pass -- stamping the mirror
|
|
# with the later mtime would mark truncated output as current. Stamping the
|
|
# earlier one leaves the two mismatched, so the next pass re-encodes it.
|
|
stat = source.stat()
|
|
|
|
# Encode to a temporary file in the destination directory and rename it
|
|
# into place, so an interrupted run cannot leave a truncated MP3 that the
|
|
# next run would treat as complete.
|
|
handle, temporary = tempfile.mkstemp(dir=mirror.parent, suffix=".mp3.part")
|
|
os.close(handle)
|
|
temporary = Path(temporary)
|
|
|
|
try:
|
|
command = build_command(source, temporary, quality_args, cover)
|
|
completed = subprocess.run(command, capture_output=True, text=True)
|
|
if completed.returncode != 0:
|
|
lines = completed.stderr.strip().splitlines()
|
|
return Result("failed", source, lines[-1] if lines else "ffmpeg failed")
|
|
os.utime(temporary, (stat.st_atime, stat.st_mtime))
|
|
# Before the rename, so the file is never visible in the mirror without
|
|
# the bit.
|
|
make_group_readable(temporary)
|
|
os.replace(temporary, mirror)
|
|
except Exception as error: # noqa: BLE001 - reported per file, run continues
|
|
return Result("failed", source, str(error))
|
|
finally:
|
|
temporary.unlink(missing_ok=True)
|
|
|
|
logger.info("encoded %s", source)
|
|
return Result("encoded", mirror)
|
|
|
|
|
|
def copy(source, mirror, dry_run):
|
|
"""Copy an already-MP3 source into the mirror, atomically."""
|
|
if dry_run:
|
|
logger.info("would copy %s", source)
|
|
return Result("copied", mirror)
|
|
|
|
mirror.parent.mkdir(parents=True, exist_ok=True)
|
|
|
|
# Through a temporary file and a rename, for the same reason encodes go
|
|
# that way, and a sharper one: copy2 reproduces the source's mtime as well
|
|
# as its bytes, so a copy cut short by a full disk or a killed container
|
|
# would leave a truncated MP3 that every later pass reads as current.
|
|
handle, temporary = tempfile.mkstemp(dir=mirror.parent, suffix=".mp3.part")
|
|
os.close(handle)
|
|
temporary = Path(temporary)
|
|
|
|
try:
|
|
shutil.copy2(source, temporary)
|
|
# copy2 brings the source's mode with it, and the source library is not
|
|
# ours to have permissions opinions about.
|
|
make_group_readable(temporary)
|
|
os.replace(temporary, mirror)
|
|
except OSError as error:
|
|
return Result("failed", source, str(error))
|
|
finally:
|
|
temporary.unlink(missing_ok=True)
|
|
|
|
logger.info("copied %s", source)
|
|
return Result("copied", mirror)
|
|
|
|
|
|
def process(source, mirror, quality_args, dry_run):
|
|
"""Bring one source file's mirror entry up to date."""
|
|
if is_current(source, mirror):
|
|
# A mirror written before this bit was set has a correct mtime, so
|
|
# nothing else in the pass would ever revisit it. Top it up here
|
|
# instead: one stat per file, and no chmod at all once it is right.
|
|
if not dry_run:
|
|
try:
|
|
make_group_readable(mirror)
|
|
except OSError as error:
|
|
return Result("failed", mirror, str(error))
|
|
return Result("skipped", mirror)
|
|
if source.suffix.lower() in COPY_EXTENSIONS:
|
|
return copy(source, mirror, dry_run)
|
|
return encode(source, mirror, quality_args, dry_run)
|
|
|
|
|
|
def find_sources(root):
|
|
"""Yield every audio file under a root, in a stable order."""
|
|
extensions = SOURCE_EXTENSIONS | COPY_EXTENSIONS
|
|
for path in sorted(root.rglob("*")):
|
|
if path.is_file() and path.suffix.lower() in extensions:
|
|
yield path
|
|
|
|
|
|
def plan(scan_root, source_root, mirror_root):
|
|
"""Map each mirror path to the one source that should produce it.
|
|
|
|
Two sources can want the same mirror path -- `01 Song.flac` alongside a
|
|
leftover `01 Song.mp3`, which is what an interrupted Lidarr upgrade leaves
|
|
behind. Without a decision here both would encode to the same destination,
|
|
each pass would find the loser stale, and the mirror would be rewritten
|
|
forever. Preferring the highest-quality source, ties broken by path, makes
|
|
the outcome stable and predictable instead.
|
|
"""
|
|
chosen = {}
|
|
for source in find_sources(scan_root):
|
|
mirror = mirror_path_for(source, source_root, mirror_root)
|
|
rival = chosen.get(mirror)
|
|
if rival is None:
|
|
chosen[mirror] = source
|
|
continue
|
|
winner, loser = sorted((source, rival), key=source_rank)
|
|
logger.warning("%s and %s both map to %s; using %s", rival, source, mirror, winner)
|
|
chosen[mirror] = winner
|
|
return chosen
|
|
|
|
|
|
def source_rank(source):
|
|
"""Sort key preferring better source formats, then a stable path order."""
|
|
suffix = source.suffix.lower()
|
|
position = SOURCE_PRIORITY.index(suffix) if suffix in SOURCE_PRIORITY else len(SOURCE_PRIORITY)
|
|
return (position, str(source))
|
|
|
|
|
|
def prune(mirror_root, expected, dry_run):
|
|
"""Delete mirror files this pass did not account for, and empty dirs.
|
|
|
|
Driven by the set of paths the pass expects to exist rather than by
|
|
probing the source tree for names, which would disagree with it over
|
|
letter case and over any extension the walker does not collect.
|
|
"""
|
|
removed = 0
|
|
|
|
for mirror in sorted(mirror_root.rglob(f"*{MIRROR_SUFFIX}")):
|
|
if mirror in expected:
|
|
continue
|
|
removed += 1
|
|
if dry_run:
|
|
logger.info("would remove orphan %s", mirror)
|
|
continue
|
|
logger.info("removing orphan %s", mirror)
|
|
mirror.unlink(missing_ok=True)
|
|
|
|
if not dry_run:
|
|
# Deepest first, so a directory emptied by the loop above is caught.
|
|
for directory in sorted(mirror_root.rglob("*"), reverse=True):
|
|
if directory.is_dir() and not any(directory.iterdir()):
|
|
directory.rmdir()
|
|
|
|
return removed
|
|
|
|
|
|
def run_once(scan_root, source_root, mirror_root, quality_args, jobs, dry_run, do_prune):
|
|
"""Run a single pass. Returns the number of failures.
|
|
|
|
`scan_root` is what gets walked and `source_root` is what mirror paths are
|
|
computed against; they differ only for a partial pass over one directory.
|
|
"""
|
|
started = time.monotonic()
|
|
logger.info("pass starting with %d concurrent encoders", jobs)
|
|
counts = {"encoded": 0, "copied": 0, "skipped": 0, "failed": 0}
|
|
failures = []
|
|
|
|
work = plan(scan_root, source_root, mirror_root)
|
|
|
|
with concurrent.futures.ThreadPoolExecutor(max_workers=jobs) as pool:
|
|
futures = [
|
|
pool.submit(process, source, mirror, quality_args, dry_run)
|
|
for mirror, source in work.items()
|
|
]
|
|
for future in concurrent.futures.as_completed(futures):
|
|
result = future.result()
|
|
counts[result.action] += 1
|
|
if result.action == "failed":
|
|
failures.append(result)
|
|
|
|
removed = prune(mirror_root, set(work), dry_run) if do_prune else 0
|
|
|
|
for failure in failures:
|
|
logger.error("failed: %s: %s", failure.path, failure.error)
|
|
|
|
logger.info(
|
|
"pass complete in %.1fs: %d encoded, %d copied, %d up to date, %d removed, %d failed",
|
|
time.monotonic() - started,
|
|
counts["encoded"],
|
|
counts["copied"],
|
|
counts["skipped"],
|
|
removed,
|
|
counts["failed"],
|
|
)
|
|
return counts["failed"]
|
|
|
|
|
|
def acquire_lock(mirror_root):
|
|
"""Take an exclusive lock so two passes cannot run over one mirror."""
|
|
mirror_root.mkdir(parents=True, exist_ok=True)
|
|
handle = open(mirror_root / ".music-mirror.lock", "w") # noqa: SIM115 - held for the process
|
|
try:
|
|
fcntl.flock(handle, fcntl.LOCK_EX | fcntl.LOCK_NB)
|
|
except OSError:
|
|
handle.close()
|
|
return None
|
|
return handle
|
|
|
|
|
|
def default_jobs():
|
|
"""Return the number of CPUs this process may actually use.
|
|
|
|
os.cpu_count() reports the host's total, which in a container with a `cpus:`
|
|
limit means starting several times more encoders than there is CPU to run
|
|
them. libmp3lame is single-threaded, so one process per available CPU is the
|
|
whole of the concurrency story.
|
|
"""
|
|
try:
|
|
quota, period = Path("/sys/fs/cgroup/cpu.max").read_text().split()
|
|
if quota != "max":
|
|
return max(1, round(int(quota) / int(period)))
|
|
except (OSError, ValueError):
|
|
pass
|
|
if hasattr(os, "process_cpu_count"): # 3.13+, respects CPU affinity
|
|
return os.process_cpu_count() or 4
|
|
return os.cpu_count() or 4
|
|
|
|
|
|
def build_parser():
|
|
"""Return the argument parser. Every option also reads an env var, so the
|
|
container can be configured without a command line."""
|
|
parser = argparse.ArgumentParser(
|
|
prog="music-mirror",
|
|
description="Maintain a lossy MP3 mirror of a lossless music library.",
|
|
)
|
|
parser.add_argument(
|
|
"--source",
|
|
default=os.getenv("MUSIC_MIRROR_SOURCE"),
|
|
help="root of the lossless library; never written to (env MUSIC_MIRROR_SOURCE)",
|
|
)
|
|
parser.add_argument(
|
|
"--mirror",
|
|
default=os.getenv("MUSIC_MIRROR_MIRROR"),
|
|
help="root of the MP3 mirror (env MUSIC_MIRROR_MIRROR)",
|
|
)
|
|
parser.add_argument(
|
|
"--quality",
|
|
default=os.getenv("MUSIC_MIRROR_QUALITY", "V0"),
|
|
help="LAME VBR level (V0-V9) or a CBR bitrate in kbps (env MUSIC_MIRROR_QUALITY)",
|
|
)
|
|
parser.add_argument(
|
|
"--jobs",
|
|
type=int,
|
|
default=int(os.getenv("MUSIC_MIRROR_JOBS", "0")) or default_jobs(),
|
|
help="concurrent encodes (env MUSIC_MIRROR_JOBS; default: available CPUs)",
|
|
)
|
|
parser.add_argument(
|
|
"--interval",
|
|
default=os.getenv("MUSIC_MIRROR_INTERVAL"),
|
|
help="repeat forever, waiting this long between passes, e.g. 6h (env MUSIC_MIRROR_INTERVAL)",
|
|
)
|
|
parser.add_argument(
|
|
"--subdir",
|
|
default=None,
|
|
help="limit the pass to one directory below --source; skips pruning",
|
|
)
|
|
parser.add_argument(
|
|
"--no-prune",
|
|
action="store_true",
|
|
help="keep mirror files whose source has been deleted",
|
|
)
|
|
parser.add_argument(
|
|
"--dry-run",
|
|
action="store_true",
|
|
help="report what would change without touching the mirror",
|
|
)
|
|
return parser
|
|
|
|
|
|
def main(argv=None):
|
|
"""Entry point. Returns a process exit code."""
|
|
logging.basicConfig(format="%(asctime)s %(levelname)s %(message)s", level=logging.INFO)
|
|
args = build_parser().parse_args(argv)
|
|
|
|
# Directories are created with 0o777 masked by the umask, so clear the group
|
|
# bits from it once here rather than chmod'ing every directory the walk
|
|
# creates. Files cannot be handled this way -- mkstemp and copy2 both set a
|
|
# mode outright -- so they get an explicit chmod instead.
|
|
inherited = os.umask(0o077)
|
|
os.umask(inherited & ~GROUP_ENTER)
|
|
|
|
if not args.source or not args.mirror:
|
|
logger.error("both --source and --mirror are required")
|
|
return 2
|
|
|
|
source_root = Path(args.source).resolve()
|
|
mirror_root = Path(args.mirror).resolve()
|
|
|
|
if not source_root.is_dir():
|
|
logger.error("source %s is not a directory", source_root)
|
|
return 2
|
|
if mirror_root == source_root or mirror_root.is_relative_to(source_root):
|
|
logger.error("mirror %s must not sit inside the source library", mirror_root)
|
|
return 2
|
|
|
|
try:
|
|
quality_args = parse_quality(args.quality)
|
|
interval = parse_interval(args.interval) if args.interval else None
|
|
except ValueError as error:
|
|
logger.error("%s", error)
|
|
return 2
|
|
|
|
scan_root = source_root
|
|
do_prune = not args.no_prune
|
|
if args.subdir:
|
|
scan_root = (source_root / args.subdir).resolve()
|
|
if not scan_root.is_relative_to(source_root) or not scan_root.is_dir():
|
|
logger.error("--subdir %s is not a directory below the source", args.subdir)
|
|
return 2
|
|
# A partial pass cannot tell an orphan from a file outside its scope.
|
|
do_prune = False
|
|
|
|
lock = acquire_lock(mirror_root)
|
|
if lock is None:
|
|
logger.error("another pass is already running over %s", mirror_root)
|
|
return 3
|
|
|
|
stopping = False
|
|
|
|
def stop(signum, _frame):
|
|
nonlocal stopping
|
|
stopping = True
|
|
logger.info("signal %d received; finishing the current pass", signum)
|
|
|
|
signal.signal(signal.SIGTERM, stop)
|
|
signal.signal(signal.SIGINT, stop)
|
|
|
|
try:
|
|
while True:
|
|
failures = run_once(
|
|
scan_root,
|
|
source_root,
|
|
mirror_root,
|
|
quality_args,
|
|
args.jobs,
|
|
args.dry_run,
|
|
do_prune,
|
|
)
|
|
if interval is None or stopping:
|
|
return 1 if failures else 0
|
|
logger.info("sleeping %ds", interval)
|
|
for _ in range(interval):
|
|
if stopping:
|
|
return 1 if failures else 0
|
|
time.sleep(1)
|
|
finally:
|
|
lock.close()
|
|
|
|
|
|
def run():
|
|
"""Console-script entry point."""
|
|
sys.exit(main())
|
|
|
|
|
|
if __name__ == "__main__":
|
|
run()
|