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>
This commit is contained in:
@@ -10,6 +10,7 @@ import re
|
|||||||
import uuid
|
import uuid
|
||||||
from typing import Annotated
|
from typing import Annotated
|
||||||
|
|
||||||
|
import anyio
|
||||||
from fastapi import APIRouter, Header, Query, Response
|
from fastapi import APIRouter, Header, Query, Response
|
||||||
from fastapi.responses import FileResponse, StreamingResponse
|
from fastapi.responses import FileResponse, StreamingResponse
|
||||||
|
|
||||||
@@ -78,7 +79,7 @@ async def stream_hls_playlist(
|
|||||||
await enqueue_transcode_quiet(track_id, hls=True)
|
await enqueue_transcode_quiet(track_id, hls=True)
|
||||||
raise NotFoundError("HLS rendition is being prepared; retry shortly.")
|
raise NotFoundError("HLS rendition is being prepared; retry shortly.")
|
||||||
|
|
||||||
body = path.read_text()
|
body = await anyio.to_thread.run_sync(path.read_text)
|
||||||
if token:
|
if token:
|
||||||
body = _SEGMENT_LINE_RE.sub(rf"\1?token={token}", body)
|
body = _SEGMENT_LINE_RE.sub(rf"\1?token={token}", body)
|
||||||
return Response(body, media_type=_HLS_PLAYLIST_TYPE)
|
return Response(body, media_type=_HLS_PLAYLIST_TYPE)
|
||||||
|
|||||||
@@ -3,6 +3,7 @@
|
|||||||
import uuid
|
import uuid
|
||||||
from typing import Annotated
|
from typing import Annotated
|
||||||
|
|
||||||
|
import anyio
|
||||||
from fastapi import APIRouter, Query, Response
|
from fastapi import APIRouter, Query, Response
|
||||||
from fastapi.responses import StreamingResponse
|
from fastapi.responses import StreamingResponse
|
||||||
|
|
||||||
@@ -33,7 +34,8 @@ from app.api.schemas.track import (
|
|||||||
TrackUpdate,
|
TrackUpdate,
|
||||||
)
|
)
|
||||||
from app.api.schemas.transcode import OptimizeEnqueuedOut
|
from app.api.schemas.transcode import OptimizeEnqueuedOut
|
||||||
from app.application.transcode_service import bitrate_for_quality
|
from app.application.transcode_service import bitrate_for_quality, remove_track_cache
|
||||||
|
from app.core.config import get_settings
|
||||||
from app.domain.entities.album import Album
|
from app.domain.entities.album import Album
|
||||||
from app.domain.entities.track import Artist, Track
|
from app.domain.entities.track import Artist, Track
|
||||||
from app.domain.errors import NotFoundError, ValidationError
|
from app.domain.errors import NotFoundError, ValidationError
|
||||||
@@ -219,6 +221,10 @@ async def delete_track(
|
|||||||
await track_repo.delete(track_id)
|
await track_repo.delete(track_id)
|
||||||
if track.storage_uri is not None:
|
if track.storage_uri is not None:
|
||||||
await storage.delete(track.storage_uri)
|
await storage.delete(track.storage_uri)
|
||||||
|
# Drop any cached transcode renditions (Opus + HLS) so they don't dangle.
|
||||||
|
await anyio.to_thread.run_sync(
|
||||||
|
remove_track_cache, get_settings().transcode_cache_path, track_id
|
||||||
|
)
|
||||||
return Response(status_code=204)
|
return Response(status_code=204)
|
||||||
|
|
||||||
|
|
||||||
|
|||||||
@@ -13,9 +13,12 @@ importing the other.
|
|||||||
"""
|
"""
|
||||||
|
|
||||||
import re
|
import re
|
||||||
|
import shutil
|
||||||
import uuid
|
import uuid
|
||||||
from pathlib import Path
|
from pathlib import Path
|
||||||
|
|
||||||
|
import anyio
|
||||||
|
|
||||||
from app.domain.errors import NotFoundError
|
from app.domain.errors import NotFoundError
|
||||||
from app.domain.ports import TrackRepository
|
from app.domain.ports import TrackRepository
|
||||||
|
|
||||||
@@ -60,6 +63,12 @@ def hls_segment_path(root: Path, track_id: uuid.UUID, name: str) -> Path | None:
|
|||||||
return hls_dir(root, track_id) / name
|
return hls_dir(root, track_id) / name
|
||||||
|
|
||||||
|
|
||||||
|
def remove_track_cache(root: Path, track_id: uuid.UUID) -> None:
|
||||||
|
"""Delete every cached rendition for a track (Opus + HLS). Best-effort — used
|
||||||
|
when a track is deleted so its transcode cache doesn't dangle forever."""
|
||||||
|
shutil.rmtree(track_cache_dir(root, track_id), ignore_errors=True)
|
||||||
|
|
||||||
|
|
||||||
class TranscodeService:
|
class TranscodeService:
|
||||||
"""Request-side cache lookups for transcoded renditions."""
|
"""Request-side cache lookups for transcoded renditions."""
|
||||||
|
|
||||||
@@ -83,14 +92,16 @@ class TranscodeService:
|
|||||||
if bitrate is None:
|
if bitrate is None:
|
||||||
return None
|
return None
|
||||||
path = opus_path(self._root, track_id, bitrate)
|
path = opus_path(self._root, track_id, bitrate)
|
||||||
return path if path.exists() else None
|
exists = await anyio.to_thread.run_sync(path.exists)
|
||||||
|
return path if exists else None
|
||||||
|
|
||||||
async def hls_playlist(self, track_id: uuid.UUID) -> Path | None:
|
async def hls_playlist(self, track_id: uuid.UUID) -> Path | None:
|
||||||
"""Cached HLS playlist if generated, else ``None`` (validates the track
|
"""Cached HLS playlist if generated, else ``None`` (validates the track
|
||||||
exists so an unknown id 404s rather than silently missing)."""
|
exists so an unknown id 404s rather than silently missing)."""
|
||||||
await self._require_streamable(track_id)
|
await self._require_streamable(track_id)
|
||||||
path = hls_playlist_path(self._root, track_id)
|
path = hls_playlist_path(self._root, track_id)
|
||||||
return path if path.exists() else None
|
exists = await anyio.to_thread.run_sync(path.exists)
|
||||||
|
return path if exists else None
|
||||||
|
|
||||||
def hls_segment(self, track_id: uuid.UUID, name: str) -> Path | None:
|
def hls_segment(self, track_id: uuid.UUID, name: str) -> Path | None:
|
||||||
path = hls_segment_path(self._root, track_id, name)
|
path = hls_segment_path(self._root, track_id, name)
|
||||||
|
|||||||
@@ -7,6 +7,7 @@ quality renditions; HLS segments are AAC-in-MPEG-TS (broad player support).
|
|||||||
"""
|
"""
|
||||||
|
|
||||||
import asyncio
|
import asyncio
|
||||||
|
import uuid
|
||||||
from pathlib import Path
|
from pathlib import Path
|
||||||
|
|
||||||
import anyio
|
import anyio
|
||||||
@@ -30,7 +31,10 @@ class FfmpegTranscoder:
|
|||||||
|
|
||||||
async def to_opus(self, src: Path, dest: Path, *, bitrate_kbps: int) -> None:
|
async def to_opus(self, src: Path, dest: Path, *, bitrate_kbps: int) -> None:
|
||||||
await anyio.to_thread.run_sync(_mkdir, dest.parent)
|
await anyio.to_thread.run_sync(_mkdir, dest.parent)
|
||||||
tmp = dest.with_suffix(dest.suffix + ".part")
|
# Per-writer temp name (not a shared ``.part``) so two concurrent jobs for
|
||||||
|
# the same rendition can't interleave into one file — each writes its own
|
||||||
|
# temp and the last atomic replace wins, both leaving a valid output.
|
||||||
|
tmp = dest.with_name(f"{dest.name}.{uuid.uuid4().hex}.part")
|
||||||
await self._run(
|
await self._run(
|
||||||
"-i", str(src),
|
"-i", str(src),
|
||||||
"-vn", "-c:a", "libopus", "-b:a", f"{bitrate_kbps}k",
|
"-vn", "-c:a", "libopus", "-b:a", f"{bitrate_kbps}k",
|
||||||
|
|||||||
@@ -7,9 +7,13 @@ enqueues (e.g. several plays before the first finishes) are cheap. The DB sessio
|
|||||||
is released before ffmpeg runs; the track entity is a detached value object.
|
is released before ffmpeg runs; the track entity is a detached value object.
|
||||||
"""
|
"""
|
||||||
|
|
||||||
|
import shutil
|
||||||
import uuid
|
import uuid
|
||||||
|
from pathlib import Path
|
||||||
from typing import Any
|
from typing import Any
|
||||||
|
|
||||||
|
import anyio
|
||||||
|
|
||||||
from app.application.transcode_service import (
|
from app.application.transcode_service import (
|
||||||
HLS_BITRATE,
|
HLS_BITRATE,
|
||||||
bitrate_for_quality,
|
bitrate_for_quality,
|
||||||
@@ -27,6 +31,15 @@ from app.infrastructure.transcode.ffmpeg import FfmpegTranscoder
|
|||||||
log = get_logger("worker.transcode")
|
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(
|
async def transcode_track(
|
||||||
_ctx: dict[str, Any],
|
_ctx: dict[str, Any],
|
||||||
*,
|
*,
|
||||||
@@ -57,7 +70,12 @@ async def transcode_track(
|
|||||||
await transcoder.to_opus(src, dest, bitrate_kbps=bitrate)
|
await transcoder.to_opus(src, dest, bitrate_kbps=bitrate)
|
||||||
did_opus = True
|
did_opus = True
|
||||||
if hls and not hls_playlist_path(root, tid).exists():
|
if hls and not hls_playlist_path(root, tid).exists():
|
||||||
await transcoder.to_hls(src, hls_dir(root, tid), bitrate_kbps=HLS_BITRATE)
|
# 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
|
did_hls = True
|
||||||
|
|
||||||
log.info("transcode_done", track_id=track_id, opus=did_opus, hls=did_hls)
|
log.info("transcode_done", track_id=track_id, opus=did_opus, hls=did_hls)
|
||||||
|
|||||||
Reference in New Issue
Block a user