Files
mcma-backend/app/workers/tasks/transcode_task.py
Цвылев Александр Вадимович 9a78cf5261 refactor(transcode): thread fs ops + atomic cache publish
- run blocking Path.exists/read_text off the event loop via anyio.to_thread
- publish HLS renditions by building into a temp dir then atomic rename, so a
  reader never sees a playlist referencing half-written segments
- per-writer temp name for opus so two concurrent jobs can't interleave
- drop a track's cached renditions on delete so they don't dangle

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-07-29 10:53:35 +03:00

83 lines
3.1 KiB
Python

"""arq task: transcode a track to a cached Opus rendition and/or HLS (§6.6).
Triggered by ``POST /tracks/{id}/optimize`` and by a streaming cache miss. Reads
the master once (via ``storage.as_local_path``) and writes into
``transcode_cache_path``. Idempotent — existing outputs are skipped, so repeated
enqueues (e.g. several plays before the first finishes) are cheap. The DB session
is released before ffmpeg runs; the track entity is a detached value object.
"""
import shutil
import uuid
from pathlib import Path
from typing import Any
import anyio
from app.application.transcode_service import (
HLS_BITRATE,
bitrate_for_quality,
hls_dir,
hls_playlist_path,
opus_path,
)
from app.core.config import get_settings
from app.core.logging import get_logger
from app.infrastructure.db import session_scope
from app.infrastructure.db.repositories import SqlAlchemyTrackRepository
from app.infrastructure.storage.provider import get_file_storage
from app.infrastructure.transcode.ffmpeg import FfmpegTranscoder
log = get_logger("worker.transcode")
def _swap_dir(tmp_dir: Path, final_dir: Path) -> None:
"""Publish a freshly-built HLS rendition atomically: replace the final dir in
one rename so a reader never sees a playlist referencing half-written
segments. Any stale partial at the destination is cleared first."""
if final_dir.exists():
shutil.rmtree(final_dir)
tmp_dir.replace(final_dir)
async def transcode_track(
_ctx: dict[str, Any],
*,
track_id: str,
quality: str = "high",
hls: bool = True,
) -> dict[str, Any]:
settings = get_settings()
tid = uuid.UUID(track_id)
async with session_scope() as session:
track = await SqlAlchemyTrackRepository(session).get_by_id(tid)
if track is None or track.storage_uri is None:
log.warning("transcode_skip_no_file", track_id=track_id)
return {"track_id": track_id, "opus": False, "hls": False}
storage = get_file_storage()
transcoder = FfmpegTranscoder(settings.ffmpeg_path)
root = settings.transcode_cache_path
did_opus = False
did_hls = False
async with storage.as_local_path(track.storage_uri) as src:
bitrate = bitrate_for_quality(quality)
if bitrate is not None:
dest = opus_path(root, tid, bitrate)
if not dest.exists():
await transcoder.to_opus(src, dest, bitrate_kbps=bitrate)
did_opus = True
if hls and not hls_playlist_path(root, tid).exists():
# Build into a private temp dir, then swap it in atomically — the
# cached playlist only becomes visible once every segment is written.
final_dir = hls_dir(root, tid)
tmp_dir = final_dir.with_name(f"hls.{uuid.uuid4().hex}.tmp")
await transcoder.to_hls(src, tmp_dir, bitrate_kbps=HLS_BITRATE)
await anyio.to_thread.run_sync(_swap_dir, tmp_dir, final_dir)
did_hls = True
log.info("transcode_done", track_id=track_id, opus=did_opus, hls=did_hls)
return {"track_id": track_id, "opus": did_opus, "hls": did_hls}