Compare commits
2 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 313af3a070 | |||
| 8271de34eb |
@@ -18,11 +18,13 @@ from sqlalchemy.ext.asyncio import AsyncSession
|
|||||||
|
|
||||||
from app.application.auth_service import AuthService
|
from app.application.auth_service import AuthService
|
||||||
from app.application.download_service import DownloadService
|
from app.application.download_service import DownloadService
|
||||||
|
from app.application.lyrics_service import LyricsService
|
||||||
from app.application.metadata_service import MetadataEnrichmentService
|
from app.application.metadata_service import MetadataEnrichmentService
|
||||||
from app.application.remote_library_service import RemoteLibraryService
|
from app.application.remote_library_service import RemoteLibraryService
|
||||||
from app.application.streaming_service import StreamingService
|
from app.application.streaming_service import StreamingService
|
||||||
from app.application.subsonic_auth_service import SubsonicAuthService
|
from app.application.subsonic_auth_service import SubsonicAuthService
|
||||||
from app.application.sync_service import SyncService
|
from app.application.sync_service import SyncService
|
||||||
|
from app.application.transcode_service import TranscodeService
|
||||||
from app.application.upload_service import UploadService
|
from app.application.upload_service import UploadService
|
||||||
from app.application.user_service import UserService
|
from app.application.user_service import UserService
|
||||||
from app.application.user_settings_service import UserSettingsService
|
from app.application.user_settings_service import UserSettingsService
|
||||||
@@ -38,6 +40,7 @@ from app.infrastructure.db.repositories import (
|
|||||||
SqlAlchemyDownloadJobRepository,
|
SqlAlchemyDownloadJobRepository,
|
||||||
SqlAlchemyHistoryRepository,
|
SqlAlchemyHistoryRepository,
|
||||||
SqlAlchemyLikeRepository,
|
SqlAlchemyLikeRepository,
|
||||||
|
SqlAlchemyLyricsRepository,
|
||||||
SqlAlchemyPlaylistRepository,
|
SqlAlchemyPlaylistRepository,
|
||||||
SqlAlchemyRefreshTokenRepository,
|
SqlAlchemyRefreshTokenRepository,
|
||||||
SqlAlchemyTrackRepository,
|
SqlAlchemyTrackRepository,
|
||||||
@@ -46,6 +49,7 @@ from app.infrastructure.db.repositories import (
|
|||||||
)
|
)
|
||||||
from app.infrastructure.metadata.acoustid import AcoustIdHttpClient
|
from app.infrastructure.metadata.acoustid import AcoustIdHttpClient
|
||||||
from app.infrastructure.metadata.fingerprint import FpcalcFingerprinter
|
from app.infrastructure.metadata.fingerprint import FpcalcFingerprinter
|
||||||
|
from app.infrastructure.metadata.lrclib import LrclibHttpClient
|
||||||
from app.infrastructure.metadata.tags import MutagenTagReader
|
from app.infrastructure.metadata.tags import MutagenTagReader
|
||||||
from app.infrastructure.sources.registry import SourceRegistry, build_source_registry
|
from app.infrastructure.sources.registry import SourceRegistry, build_source_registry
|
||||||
from app.infrastructure.storage.provider import get_file_storage
|
from app.infrastructure.storage.provider import get_file_storage
|
||||||
@@ -175,6 +179,28 @@ def get_metadata_service(session: SessionDep, storage: FileStorageDep) -> Metada
|
|||||||
)
|
)
|
||||||
|
|
||||||
|
|
||||||
|
def get_transcode_service(session: SessionDep) -> TranscodeService:
|
||||||
|
"""Request-side cache lookups for transcoded renditions (§6.6). Generation
|
||||||
|
itself runs in the ``transcode_track`` worker, never here."""
|
||||||
|
return TranscodeService(
|
||||||
|
tracks=SqlAlchemyTrackRepository(session),
|
||||||
|
cache_root=get_settings().transcode_cache_path,
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
def get_lyrics_service(session: SessionDep) -> LyricsService:
|
||||||
|
"""Wires the LRCLIB lyrics provider + cache repo (plan §6.7). LRCLIB is
|
||||||
|
keyless, so this is always available; failures degrade to ``not_found``."""
|
||||||
|
settings = get_settings()
|
||||||
|
return LyricsService(
|
||||||
|
lyrics=SqlAlchemyLyricsRepository(session),
|
||||||
|
tracks=SqlAlchemyTrackRepository(session),
|
||||||
|
artists=SqlAlchemyArtistRepository(session),
|
||||||
|
albums=SqlAlchemyAlbumRepository(session),
|
||||||
|
provider=LrclibHttpClient(user_agent=settings.musicbrainz_user_agent),
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
def get_download_service(session: SessionDep, storage: FileStorageDep) -> DownloadService:
|
def get_download_service(session: SessionDep, storage: FileStorageDep) -> DownloadService:
|
||||||
return DownloadService(
|
return DownloadService(
|
||||||
jobs=SqlAlchemyDownloadJobRepository(session),
|
jobs=SqlAlchemyDownloadJobRepository(session),
|
||||||
@@ -198,6 +224,8 @@ def get_remote_library_service(session: SessionDep) -> RemoteLibraryService:
|
|||||||
UploadServiceDep = Annotated[UploadService, Depends(get_upload_service)]
|
UploadServiceDep = Annotated[UploadService, Depends(get_upload_service)]
|
||||||
StreamingServiceDep = Annotated[StreamingService, Depends(get_streaming_service)]
|
StreamingServiceDep = Annotated[StreamingService, Depends(get_streaming_service)]
|
||||||
MetadataServiceDep = Annotated[MetadataEnrichmentService, Depends(get_metadata_service)]
|
MetadataServiceDep = Annotated[MetadataEnrichmentService, Depends(get_metadata_service)]
|
||||||
|
LyricsServiceDep = Annotated[LyricsService, Depends(get_lyrics_service)]
|
||||||
|
TranscodeServiceDep = Annotated[TranscodeService, Depends(get_transcode_service)]
|
||||||
DownloadServiceDep = Annotated[DownloadService, Depends(get_download_service)]
|
DownloadServiceDep = Annotated[DownloadService, Depends(get_download_service)]
|
||||||
RemoteLibraryServiceDep = Annotated[RemoteLibraryService, Depends(get_remote_library_service)]
|
RemoteLibraryServiceDep = Annotated[RemoteLibraryService, Depends(get_remote_library_service)]
|
||||||
|
|
||||||
|
|||||||
@@ -0,0 +1,33 @@
|
|||||||
|
"""Lyrics response schema (§6.7 / Now Playing lyrics panel).
|
||||||
|
|
||||||
|
Returns the raw LRC (``synced``) and/or ``plain`` text; the client parses LRC
|
||||||
|
timestamps for synced highlighting. A miss is a normal 200 with
|
||||||
|
``status="not_found"`` and null text — not an error — so the panel can render a
|
||||||
|
"no lyrics" state.
|
||||||
|
"""
|
||||||
|
|
||||||
|
import uuid
|
||||||
|
|
||||||
|
from pydantic import BaseModel
|
||||||
|
|
||||||
|
from app.domain.entities.lyrics import Lyrics
|
||||||
|
|
||||||
|
|
||||||
|
class LyricsOut(BaseModel):
|
||||||
|
track_id: uuid.UUID
|
||||||
|
status: str
|
||||||
|
source: str | None
|
||||||
|
synced: str | None
|
||||||
|
plain: str | None
|
||||||
|
synced_available: bool
|
||||||
|
|
||||||
|
@classmethod
|
||||||
|
def from_entity(cls, lyrics: Lyrics) -> LyricsOut:
|
||||||
|
return cls(
|
||||||
|
track_id=lyrics.track_id,
|
||||||
|
status=lyrics.status,
|
||||||
|
source=lyrics.source,
|
||||||
|
synced=lyrics.synced,
|
||||||
|
plain=lyrics.plain,
|
||||||
|
synced_available=lyrics.synced is not None,
|
||||||
|
)
|
||||||
@@ -0,0 +1,11 @@
|
|||||||
|
"""Transcode/optimize response schemas (§6.6)."""
|
||||||
|
|
||||||
|
from pydantic import BaseModel
|
||||||
|
|
||||||
|
|
||||||
|
class OptimizeEnqueuedOut(BaseModel):
|
||||||
|
"""Acknowledgement that a transcode job was queued (it runs in the worker)."""
|
||||||
|
|
||||||
|
status: str
|
||||||
|
job_id: str
|
||||||
|
quality: str
|
||||||
+64
-6
@@ -1,30 +1,53 @@
|
|||||||
"""Audio streaming endpoint — direct stream with Range support."""
|
"""Audio streaming — direct byte-range stream, transcoded quality, and HLS.
|
||||||
|
|
||||||
|
``GET /stream/{id}`` streams the master with Range support, or a cached Opus
|
||||||
|
rendition when ``?quality=`` is set (a cache miss falls back to the master and
|
||||||
|
warms the cache in the background — playback never waits on ffmpeg). ``/hls/*``
|
||||||
|
serves the cached HLS rendition (generated by the ``transcode_track`` worker).
|
||||||
|
"""
|
||||||
|
|
||||||
|
import re
|
||||||
import uuid
|
import uuid
|
||||||
from typing import Annotated
|
from typing import Annotated
|
||||||
|
|
||||||
from fastapi import APIRouter, Header
|
from fastapi import APIRouter, Header, Query, Response
|
||||||
from fastapi.responses import StreamingResponse
|
from fastapi.responses import FileResponse, StreamingResponse
|
||||||
|
|
||||||
from app.api.deps import StreamingServiceDep, StreamUser
|
from app.api.deps import StreamingServiceDep, StreamUser, TranscodeServiceDep
|
||||||
|
from app.domain.errors import NotFoundError
|
||||||
|
from app.workers.queue import enqueue_transcode_quiet
|
||||||
|
|
||||||
router = APIRouter(prefix="/stream", tags=["streaming"])
|
router = APIRouter(prefix="/stream", tags=["streaming"])
|
||||||
|
|
||||||
|
_HLS_PLAYLIST_TYPE = "application/vnd.apple.mpegurl"
|
||||||
|
_HLS_SEGMENT_TYPE = "video/mp2t"
|
||||||
|
_OPUS_TYPE = "audio/ogg"
|
||||||
|
_SEGMENT_LINE_RE = re.compile(r"^(seg_\d+\.ts)$", re.MULTILINE)
|
||||||
|
|
||||||
|
|
||||||
@router.get("/{track_id}")
|
@router.get("/{track_id}")
|
||||||
async def stream_track(
|
async def stream_track(
|
||||||
track_id: uuid.UUID,
|
track_id: uuid.UUID,
|
||||||
service: StreamingServiceDep,
|
service: StreamingServiceDep,
|
||||||
|
transcode: TranscodeServiceDep,
|
||||||
_user: StreamUser,
|
_user: StreamUser,
|
||||||
range_header: Annotated[str | None, Header(alias="Range")] = None,
|
range_header: Annotated[str | None, Header(alias="Range")] = None,
|
||||||
) -> StreamingResponse:
|
quality: Annotated[str | None, Query()] = None,
|
||||||
|
) -> Response:
|
||||||
|
# A quality rendition, if one is cached; otherwise fall back to the master
|
||||||
|
# and enqueue generation so the next play gets it (graceful degradation).
|
||||||
|
if quality and quality != "original":
|
||||||
|
cached = await transcode.resolve_quality_file(track_id, quality)
|
||||||
|
if cached is not None:
|
||||||
|
return FileResponse(cached, media_type=_OPUS_TYPE)
|
||||||
|
await enqueue_transcode_quiet(track_id, quality=quality, hls=False)
|
||||||
|
|
||||||
result = await service.open_stream(track_id, range_header)
|
result = await service.open_stream(track_id, range_header)
|
||||||
|
|
||||||
headers = {
|
headers = {
|
||||||
"Accept-Ranges": "bytes",
|
"Accept-Ranges": "bytes",
|
||||||
"Content-Length": str(result.content_length),
|
"Content-Length": str(result.content_length),
|
||||||
}
|
}
|
||||||
|
|
||||||
if result.is_partial:
|
if result.is_partial:
|
||||||
headers["Content-Range"] = f"bytes {result.start}-{result.end}/{result.total_size}"
|
headers["Content-Range"] = f"bytes {result.start}-{result.end}/{result.total_size}"
|
||||||
status_code = 206
|
status_code = 206
|
||||||
@@ -37,3 +60,38 @@ async def stream_track(
|
|||||||
headers=headers,
|
headers=headers,
|
||||||
media_type=result.content_type,
|
media_type=result.content_type,
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
||||||
|
@router.get("/{track_id}/hls/playlist.m3u8")
|
||||||
|
async def stream_hls_playlist(
|
||||||
|
track_id: uuid.UUID,
|
||||||
|
transcode: TranscodeServiceDep,
|
||||||
|
_user: StreamUser,
|
||||||
|
token: Annotated[str | None, Query()] = None,
|
||||||
|
) -> Response:
|
||||||
|
"""Serve the cached HLS playlist. On a miss, kick off generation and 404 so
|
||||||
|
the client retries. Segment URLs are relative; when the request carried a
|
||||||
|
``?token=`` (players can't set an Authorization header), it's appended to
|
||||||
|
each segment line so the segment requests authenticate the same way."""
|
||||||
|
path = await transcode.hls_playlist(track_id)
|
||||||
|
if path is None:
|
||||||
|
await enqueue_transcode_quiet(track_id, hls=True)
|
||||||
|
raise NotFoundError("HLS rendition is being prepared; retry shortly.")
|
||||||
|
|
||||||
|
body = path.read_text()
|
||||||
|
if token:
|
||||||
|
body = _SEGMENT_LINE_RE.sub(rf"\1?token={token}", body)
|
||||||
|
return Response(body, media_type=_HLS_PLAYLIST_TYPE)
|
||||||
|
|
||||||
|
|
||||||
|
@router.get("/{track_id}/hls/{segment}")
|
||||||
|
async def stream_hls_segment(
|
||||||
|
track_id: uuid.UUID,
|
||||||
|
segment: str,
|
||||||
|
transcode: TranscodeServiceDep,
|
||||||
|
_user: StreamUser,
|
||||||
|
) -> FileResponse:
|
||||||
|
path = transcode.hls_segment(track_id, segment)
|
||||||
|
if path is None:
|
||||||
|
raise NotFoundError("Segment not found.")
|
||||||
|
return FileResponse(path, media_type=_HLS_SEGMENT_TYPE)
|
||||||
|
|||||||
+43
-5
@@ -1,7 +1,7 @@
|
|||||||
"""Track endpoints."""
|
"""Track endpoints."""
|
||||||
|
|
||||||
import uuid
|
import uuid
|
||||||
from typing import Any
|
from typing import Annotated, Any
|
||||||
|
|
||||||
from fastapi import APIRouter, Query, Response
|
from fastapi import APIRouter, Query, Response
|
||||||
from fastapi.responses import StreamingResponse
|
from fastapi.responses import StreamingResponse
|
||||||
@@ -12,12 +12,14 @@ from app.api.deps import (
|
|||||||
ArtistRepoDep,
|
ArtistRepoDep,
|
||||||
CurrentUser,
|
CurrentUser,
|
||||||
FileStorageDep,
|
FileStorageDep,
|
||||||
|
LyricsServiceDep,
|
||||||
MetadataServiceDep,
|
MetadataServiceDep,
|
||||||
RemoteLibraryServiceDep,
|
RemoteLibraryServiceDep,
|
||||||
StreamUser,
|
StreamUser,
|
||||||
TrackRepoDep,
|
TrackRepoDep,
|
||||||
)
|
)
|
||||||
from app.api.schemas.download import DownloadJobOut
|
from app.api.schemas.download import DownloadJobOut
|
||||||
|
from app.api.schemas.lyrics import LyricsOut
|
||||||
from app.api.schemas.pagination import PagedResponse
|
from app.api.schemas.pagination import PagedResponse
|
||||||
from app.api.schemas.track import (
|
from app.api.schemas.track import (
|
||||||
MaterializeResponse,
|
MaterializeResponse,
|
||||||
@@ -28,10 +30,12 @@ from app.api.schemas.track import (
|
|||||||
TrackOut,
|
TrackOut,
|
||||||
TrackUpdate,
|
TrackUpdate,
|
||||||
)
|
)
|
||||||
|
from app.api.schemas.transcode import OptimizeEnqueuedOut
|
||||||
|
from app.application.transcode_service import bitrate_for_quality
|
||||||
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
|
from app.domain.errors import NotFoundError, ValidationError
|
||||||
from app.workers.queue import enqueue
|
from app.workers.queue import enqueue, enqueue_transcode
|
||||||
|
|
||||||
router = APIRouter(prefix="/tracks", tags=["tracks"])
|
router = APIRouter(prefix="/tracks", tags=["tracks"])
|
||||||
|
|
||||||
@@ -220,8 +224,42 @@ async def delete_track(
|
|||||||
async def get_similar_tracks(track_id: uuid.UUID, _: CurrentUser) -> Any: ...
|
async def get_similar_tracks(track_id: uuid.UUID, _: CurrentUser) -> Any: ...
|
||||||
|
|
||||||
|
|
||||||
@router.post("/{track_id}/optimize")
|
@router.post("/{track_id}/optimize", status_code=202)
|
||||||
async def optimize_track(track_id: uuid.UUID, _: CurrentUser) -> Any: ...
|
async def optimize_track(
|
||||||
|
track_id: uuid.UUID,
|
||||||
|
track_repo: TrackRepoDep,
|
||||||
|
_: CurrentUser,
|
||||||
|
quality: Annotated[str, Query()] = "high",
|
||||||
|
) -> OptimizeEnqueuedOut:
|
||||||
|
"""Enqueue transcoding of a track into a cached Opus rendition + HLS (§6.6).
|
||||||
|
Heavy ffmpeg work runs in the worker; this only queues it."""
|
||||||
|
if bitrate_for_quality(quality) is None:
|
||||||
|
raise ValidationError(
|
||||||
|
f"Unknown quality '{quality}'; expected one of high, medium, low."
|
||||||
|
)
|
||||||
|
track = await track_repo.get_by_id(track_id)
|
||||||
|
if track is None:
|
||||||
|
raise NotFoundError(f"Track {track_id} not found.")
|
||||||
|
job_id = await enqueue_transcode(track_id, quality=quality, hls=True)
|
||||||
|
return OptimizeEnqueuedOut(status="enqueued", job_id=job_id, quality=quality)
|
||||||
|
|
||||||
|
|
||||||
|
@router.get("/{track_id}/lyrics")
|
||||||
|
async def get_track_lyrics(
|
||||||
|
track_id: uuid.UUID, lyrics: LyricsServiceDep, _: CurrentUser
|
||||||
|
) -> LyricsOut:
|
||||||
|
"""Cached lyrics for the Now Playing panel (§6.7). A miss is a normal 200
|
||||||
|
with ``status="not_found"`` — the provider (LRCLIB) is queried at most once,
|
||||||
|
then the outcome is cached."""
|
||||||
|
return LyricsOut.from_entity(await lyrics.get_lyrics(track_id))
|
||||||
|
|
||||||
|
|
||||||
|
@router.post("/{track_id}/lyrics/refetch")
|
||||||
|
async def refetch_track_lyrics(
|
||||||
|
track_id: uuid.UUID, lyrics: LyricsServiceDep, _: CurrentUser
|
||||||
|
) -> LyricsOut:
|
||||||
|
"""Force a fresh provider lookup, bypassing the cache (user-triggered)."""
|
||||||
|
return LyricsOut.from_entity(await lyrics.get_lyrics(track_id, force=True))
|
||||||
|
|
||||||
|
|
||||||
@router.get("/{track_id}/cover")
|
@router.get("/{track_id}/cover")
|
||||||
|
|||||||
@@ -0,0 +1,97 @@
|
|||||||
|
"""Lyrics service (plan §6.7).
|
||||||
|
|
||||||
|
Get-or-fetch with caching: a track's lyrics are served from the DB when present;
|
||||||
|
on a miss (or an expired ``not_found``) we ask the provider (LRCLIB) once, then
|
||||||
|
cache the outcome. ``not_found`` is cached with a TTL so tracks that genuinely
|
||||||
|
have no lyrics aren't looked up on every play, but can eventually be retried.
|
||||||
|
|
||||||
|
Degrades gracefully: if the provider is unreachable the lookup just yields a
|
||||||
|
``not_found`` — the endpoint still returns 200 with empty lyrics, never an error.
|
||||||
|
"""
|
||||||
|
|
||||||
|
import datetime as dt
|
||||||
|
import uuid
|
||||||
|
|
||||||
|
from app.domain.entities.lyrics import Lyrics
|
||||||
|
from app.domain.errors import NotFoundError
|
||||||
|
from app.domain.ports import (
|
||||||
|
AlbumRepository,
|
||||||
|
ArtistRepository,
|
||||||
|
LyricsProvider,
|
||||||
|
LyricsRepository,
|
||||||
|
TrackRepository,
|
||||||
|
)
|
||||||
|
|
||||||
|
# Re-lookup a cached "not_found" only after this long — long enough not to spam
|
||||||
|
# the provider, short enough that lyrics added upstream eventually surface.
|
||||||
|
_NOT_FOUND_TTL = dt.timedelta(days=7)
|
||||||
|
|
||||||
|
_STATUS_FOUND = "found"
|
||||||
|
_STATUS_NOT_FOUND = "not_found"
|
||||||
|
|
||||||
|
|
||||||
|
class LyricsService:
|
||||||
|
def __init__(
|
||||||
|
self,
|
||||||
|
*,
|
||||||
|
lyrics: LyricsRepository,
|
||||||
|
tracks: TrackRepository,
|
||||||
|
artists: ArtistRepository,
|
||||||
|
albums: AlbumRepository,
|
||||||
|
provider: LyricsProvider,
|
||||||
|
) -> None:
|
||||||
|
self._lyrics = lyrics
|
||||||
|
self._tracks = tracks
|
||||||
|
self._artists = artists
|
||||||
|
self._albums = albums
|
||||||
|
self._provider = provider
|
||||||
|
|
||||||
|
async def get_lyrics(self, track_id: uuid.UUID, *, force: bool = False) -> Lyrics:
|
||||||
|
"""Return cached lyrics, fetching from the provider on a miss/expiry.
|
||||||
|
``force`` (the refetch endpoint) bypasses the cache entirely."""
|
||||||
|
cached = await self._lyrics.get(track_id)
|
||||||
|
if not force and cached is not None and self._is_fresh(cached):
|
||||||
|
return cached
|
||||||
|
return await self._fetch_and_cache(track_id)
|
||||||
|
|
||||||
|
def _is_fresh(self, cached: Lyrics) -> bool:
|
||||||
|
if cached.status == _STATUS_FOUND:
|
||||||
|
return True
|
||||||
|
if cached.status == _STATUS_NOT_FOUND:
|
||||||
|
return dt.datetime.now(dt.UTC) - cached.fetched_at < _NOT_FOUND_TTL
|
||||||
|
# "pending" (never fetched) → not fresh, go fetch.
|
||||||
|
return False
|
||||||
|
|
||||||
|
async def _fetch_and_cache(self, track_id: uuid.UUID) -> Lyrics:
|
||||||
|
track = await self._tracks.get_by_id(track_id)
|
||||||
|
if track is None:
|
||||||
|
raise NotFoundError(f"Track {track_id} not found.")
|
||||||
|
|
||||||
|
artist = await self._artists.get_by_id(track.artist_id)
|
||||||
|
album = (
|
||||||
|
await self._albums.get_by_id(track.album_id)
|
||||||
|
if track.album_id is not None
|
||||||
|
else None
|
||||||
|
)
|
||||||
|
result = await self._provider.fetch(
|
||||||
|
artist=artist.name if artist else "",
|
||||||
|
title=track.title,
|
||||||
|
album=album.title if album else None,
|
||||||
|
duration_seconds=track.duration_seconds,
|
||||||
|
)
|
||||||
|
|
||||||
|
if result is None:
|
||||||
|
return await self._lyrics.upsert(
|
||||||
|
track_id=track_id,
|
||||||
|
synced=None,
|
||||||
|
plain=None,
|
||||||
|
source=None,
|
||||||
|
status=_STATUS_NOT_FOUND,
|
||||||
|
)
|
||||||
|
return await self._lyrics.upsert(
|
||||||
|
track_id=track_id,
|
||||||
|
synced=result.synced,
|
||||||
|
plain=result.plain,
|
||||||
|
source=result.source,
|
||||||
|
status=_STATUS_FOUND,
|
||||||
|
)
|
||||||
@@ -0,0 +1,99 @@
|
|||||||
|
"""Transcode service + cache-path helpers (Group B / plan §6.6).
|
||||||
|
|
||||||
|
Cache layout under ``transcode_cache_path``::
|
||||||
|
|
||||||
|
{track_id}/opus_{kbps}.opus # direct quality renditions
|
||||||
|
{track_id}/hls/playlist.m3u8 # HLS rendition (AAC-in-TS)
|
||||||
|
{track_id}/hls/seg_000.ts …
|
||||||
|
|
||||||
|
The request side (streaming router) only *reads* the cache — misses fall back to
|
||||||
|
the original file and enqueue generation. The worker (``transcode_task``) writes
|
||||||
|
it. Path helpers are module-level so both sides agree on locations without one
|
||||||
|
importing the other.
|
||||||
|
"""
|
||||||
|
|
||||||
|
import re
|
||||||
|
import uuid
|
||||||
|
from pathlib import Path
|
||||||
|
|
||||||
|
from app.domain.errors import NotFoundError
|
||||||
|
from app.domain.ports import TrackRepository
|
||||||
|
|
||||||
|
# Stream-quality name (matches the user-settings ``StreamQuality``) → Opus
|
||||||
|
# bitrate. ``original`` is absent: it means "serve the master, no transcode".
|
||||||
|
QUALITY_BITRATE: dict[str, int] = {"high": 128, "medium": 96, "low": 64}
|
||||||
|
|
||||||
|
# Single HLS rendition bitrate (AAC). One rendition keeps the MVP simple; a
|
||||||
|
# multi-bitrate ladder can come later.
|
||||||
|
HLS_BITRATE = 128
|
||||||
|
|
||||||
|
# Only these segment names may be served, guarding the segment route against
|
||||||
|
# path traversal.
|
||||||
|
_SEGMENT_RE = re.compile(r"^seg_\d{3,}\.ts$")
|
||||||
|
|
||||||
|
|
||||||
|
def bitrate_for_quality(quality: str) -> int | None:
|
||||||
|
"""Opus bitrate for a quality name, or ``None`` for ``original``/unknown."""
|
||||||
|
return QUALITY_BITRATE.get(quality)
|
||||||
|
|
||||||
|
|
||||||
|
def track_cache_dir(root: Path, track_id: uuid.UUID) -> Path:
|
||||||
|
return root / str(track_id)
|
||||||
|
|
||||||
|
|
||||||
|
def opus_path(root: Path, track_id: uuid.UUID, bitrate_kbps: int) -> Path:
|
||||||
|
return track_cache_dir(root, track_id) / f"opus_{bitrate_kbps}.opus"
|
||||||
|
|
||||||
|
|
||||||
|
def hls_dir(root: Path, track_id: uuid.UUID) -> Path:
|
||||||
|
return track_cache_dir(root, track_id) / "hls"
|
||||||
|
|
||||||
|
|
||||||
|
def hls_playlist_path(root: Path, track_id: uuid.UUID) -> Path:
|
||||||
|
return hls_dir(root, track_id) / "playlist.m3u8"
|
||||||
|
|
||||||
|
|
||||||
|
def hls_segment_path(root: Path, track_id: uuid.UUID, name: str) -> Path | None:
|
||||||
|
"""Resolve a segment file, or ``None`` if the name is not a valid segment."""
|
||||||
|
if not _SEGMENT_RE.fullmatch(name):
|
||||||
|
return None
|
||||||
|
return hls_dir(root, track_id) / name
|
||||||
|
|
||||||
|
|
||||||
|
class TranscodeService:
|
||||||
|
"""Request-side cache lookups for transcoded renditions."""
|
||||||
|
|
||||||
|
def __init__(self, *, tracks: TrackRepository, cache_root: Path) -> None:
|
||||||
|
self._tracks = tracks
|
||||||
|
self._root = cache_root
|
||||||
|
|
||||||
|
async def _require_streamable(self, track_id: uuid.UUID) -> None:
|
||||||
|
track = await self._tracks.get_by_id(track_id)
|
||||||
|
if track is None:
|
||||||
|
raise NotFoundError("Track not found.")
|
||||||
|
if track.storage_uri is None:
|
||||||
|
raise NotFoundError("Track is not yet downloaded.")
|
||||||
|
|
||||||
|
async def resolve_quality_file(
|
||||||
|
self, track_id: uuid.UUID, quality: str
|
||||||
|
) -> Path | None:
|
||||||
|
"""Cached Opus file for ``quality`` if present, else ``None`` (caller
|
||||||
|
falls back to the master and enqueues generation). ``original`` → None."""
|
||||||
|
bitrate = bitrate_for_quality(quality)
|
||||||
|
if bitrate is None:
|
||||||
|
return None
|
||||||
|
path = opus_path(self._root, track_id, bitrate)
|
||||||
|
return path if path.exists() else None
|
||||||
|
|
||||||
|
async def hls_playlist(self, track_id: uuid.UUID) -> Path | None:
|
||||||
|
"""Cached HLS playlist if generated, else ``None`` (validates the track
|
||||||
|
exists so an unknown id 404s rather than silently missing)."""
|
||||||
|
await self._require_streamable(track_id)
|
||||||
|
path = hls_playlist_path(self._root, track_id)
|
||||||
|
return path if path.exists() else None
|
||||||
|
|
||||||
|
def hls_segment(self, track_id: uuid.UUID, name: str) -> Path | None:
|
||||||
|
path = hls_segment_path(self._root, track_id, name)
|
||||||
|
if path is None or not path.exists():
|
||||||
|
return None
|
||||||
|
return path
|
||||||
@@ -81,6 +81,9 @@ class Settings(BaseSettings):
|
|||||||
# -- media / storage --------------------------------------------------
|
# -- media / storage --------------------------------------------------
|
||||||
media_path: Path = Path("/data/media")
|
media_path: Path = Path("/data/media")
|
||||||
transcode_cache_path: Path = Path("/data/transcode-cache")
|
transcode_cache_path: Path = Path("/data/transcode-cache")
|
||||||
|
# ffmpeg binary for transcoding/HLS (on PATH in the image); override for a
|
||||||
|
# non-standard location.
|
||||||
|
ffmpeg_path: str = "ffmpeg"
|
||||||
max_parallel_downloads: int = 2
|
max_parallel_downloads: int = 2
|
||||||
# How many times the download worker retries a failed fetch (yt-dlp fails
|
# How many times the download worker retries a failed fetch (yt-dlp fails
|
||||||
# often) before marking the job ``failed`` — exponential backoff between tries.
|
# often) before marking the job ``failed`` — exponential backoff between tries.
|
||||||
|
|||||||
@@ -0,0 +1,41 @@
|
|||||||
|
"""Lyrics value objects (plan §6.7).
|
||||||
|
|
||||||
|
``Lyrics`` is the cached row for a track; ``LyricsResult`` is what a provider
|
||||||
|
(LRCLIB) returns for a lookup. Both cross the domain boundary — no framework
|
||||||
|
imports. Status values mirror ``LyricsStatus`` in the ORM enum ("found" /
|
||||||
|
"not_found" / "pending") but are kept as plain strings here so the domain stays
|
||||||
|
independent of the persistence layer.
|
||||||
|
"""
|
||||||
|
|
||||||
|
import datetime as dt
|
||||||
|
import uuid
|
||||||
|
from dataclasses import dataclass
|
||||||
|
|
||||||
|
|
||||||
|
@dataclass(frozen=True, slots=True)
|
||||||
|
class LyricsResult:
|
||||||
|
"""A provider hit: synced (timestamped LRC) and/or plain text."""
|
||||||
|
|
||||||
|
synced: str | None
|
||||||
|
plain: str | None
|
||||||
|
source: str
|
||||||
|
|
||||||
|
|
||||||
|
@dataclass(frozen=True, slots=True)
|
||||||
|
class Lyrics:
|
||||||
|
"""Cached lyrics for one track. ``status`` is ``found`` / ``not_found`` /
|
||||||
|
``pending``; ``not_found`` is cached too (with a TTL in the service) so a
|
||||||
|
track with no lyrics doesn't hammer the provider on every play."""
|
||||||
|
|
||||||
|
track_id: uuid.UUID
|
||||||
|
synced: str | None
|
||||||
|
plain: str | None
|
||||||
|
source: str | None
|
||||||
|
status: str
|
||||||
|
fetched_at: dt.datetime
|
||||||
|
|
||||||
|
@property
|
||||||
|
def has_lyrics(self) -> bool:
|
||||||
|
return self.status == "found" and (
|
||||||
|
self.synced is not None or self.plain is not None
|
||||||
|
)
|
||||||
@@ -76,6 +76,13 @@ class StorageError(DomainError):
|
|||||||
code = "storage_error"
|
code = "storage_error"
|
||||||
|
|
||||||
|
|
||||||
|
class TranscodeError(DomainError):
|
||||||
|
"""Transcoding (ffmpeg) failed. Raised in the worker; a play falls back to
|
||||||
|
the original file rather than surfacing this."""
|
||||||
|
|
||||||
|
code = "transcode_error"
|
||||||
|
|
||||||
|
|
||||||
class RangeNotSatisfiableError(DomainError):
|
class RangeNotSatisfiableError(DomainError):
|
||||||
"""Requested byte range cannot be satisfied."""
|
"""Requested byte range cannot be satisfied."""
|
||||||
|
|
||||||
|
|||||||
@@ -29,6 +29,7 @@ from app.domain.entities import (
|
|||||||
SubsonicCredentials,
|
SubsonicCredentials,
|
||||||
User,
|
User,
|
||||||
)
|
)
|
||||||
|
from app.domain.entities.lyrics import Lyrics, LyricsResult
|
||||||
from app.domain.entities.settings import UserSettings
|
from app.domain.entities.settings import UserSettings
|
||||||
from app.domain.entities.track import Artist, Track
|
from app.domain.entities.track import Artist, Track
|
||||||
from app.domain.sources import DownloadResult, RawMetadata, SearchResult, SourceFile, SourceInfo
|
from app.domain.sources import DownloadResult, RawMetadata, SearchResult, SourceFile, SourceInfo
|
||||||
@@ -494,3 +495,43 @@ class CoverArtProvider(Protocol):
|
|||||||
|
|
||||||
def is_available(self) -> bool: ...
|
def is_available(self) -> bool: ...
|
||||||
async def fetch_release_group(self, release_group_mbid: str) -> CoverArt | None: ...
|
async def fetch_release_group(self, release_group_mbid: str) -> CoverArt | None: ...
|
||||||
|
|
||||||
|
|
||||||
|
class Transcoder(Protocol):
|
||||||
|
"""Transcodes an audio file with ffmpeg (plan §6.6 / Group B). ``to_opus``
|
||||||
|
writes a single Opus rendition; ``to_hls`` writes an HLS playlist + segments
|
||||||
|
(AAC-in-TS) into ``out_dir``. Both raise ``TranscodeError`` on failure; heavy
|
||||||
|
work always runs in a worker, never the request cycle."""
|
||||||
|
|
||||||
|
async def to_opus(self, src: Path, dest: Path, *, bitrate_kbps: int) -> None: ...
|
||||||
|
async def to_hls(self, src: Path, out_dir: Path, *, bitrate_kbps: int) -> None: ...
|
||||||
|
|
||||||
|
|
||||||
|
class LyricsProvider(Protocol):
|
||||||
|
"""Fetches lyrics from an external database (LRCLIB) by artist/title/album/
|
||||||
|
duration. Returns a hit or ``None`` (no match / service down), never raising."""
|
||||||
|
|
||||||
|
async def fetch(
|
||||||
|
self,
|
||||||
|
*,
|
||||||
|
artist: str,
|
||||||
|
title: str,
|
||||||
|
album: str | None,
|
||||||
|
duration_seconds: int | None,
|
||||||
|
) -> LyricsResult | None: ...
|
||||||
|
|
||||||
|
|
||||||
|
class LyricsRepository(Protocol):
|
||||||
|
"""Cached lyrics, one row per track. ``upsert`` also caches a ``not_found``
|
||||||
|
(empty text) so misses aren't re-fetched until the service's TTL lapses."""
|
||||||
|
|
||||||
|
async def get(self, track_id: uuid.UUID) -> Lyrics | None: ...
|
||||||
|
async def upsert(
|
||||||
|
self,
|
||||||
|
*,
|
||||||
|
track_id: uuid.UUID,
|
||||||
|
synced: str | None,
|
||||||
|
plain: str | None,
|
||||||
|
source: str | None,
|
||||||
|
status: str,
|
||||||
|
) -> Lyrics: ...
|
||||||
|
|||||||
@@ -7,6 +7,7 @@ from app.infrastructure.db.repositories.download_job_repository import (
|
|||||||
)
|
)
|
||||||
from app.infrastructure.db.repositories.history_repository import SqlAlchemyHistoryRepository
|
from app.infrastructure.db.repositories.history_repository import SqlAlchemyHistoryRepository
|
||||||
from app.infrastructure.db.repositories.like_repository import SqlAlchemyLikeRepository
|
from app.infrastructure.db.repositories.like_repository import SqlAlchemyLikeRepository
|
||||||
|
from app.infrastructure.db.repositories.lyrics_repository import SqlAlchemyLyricsRepository
|
||||||
from app.infrastructure.db.repositories.playlist_repository import SqlAlchemyPlaylistRepository
|
from app.infrastructure.db.repositories.playlist_repository import SqlAlchemyPlaylistRepository
|
||||||
from app.infrastructure.db.repositories.refresh_token_repository import (
|
from app.infrastructure.db.repositories.refresh_token_repository import (
|
||||||
SqlAlchemyRefreshTokenRepository,
|
SqlAlchemyRefreshTokenRepository,
|
||||||
@@ -23,6 +24,7 @@ __all__ = [
|
|||||||
"SqlAlchemyDownloadJobRepository",
|
"SqlAlchemyDownloadJobRepository",
|
||||||
"SqlAlchemyHistoryRepository",
|
"SqlAlchemyHistoryRepository",
|
||||||
"SqlAlchemyLikeRepository",
|
"SqlAlchemyLikeRepository",
|
||||||
|
"SqlAlchemyLyricsRepository",
|
||||||
"SqlAlchemyPlaylistRepository",
|
"SqlAlchemyPlaylistRepository",
|
||||||
"SqlAlchemyRefreshTokenRepository",
|
"SqlAlchemyRefreshTokenRepository",
|
||||||
"SqlAlchemyTrackRepository",
|
"SqlAlchemyTrackRepository",
|
||||||
|
|||||||
@@ -0,0 +1,72 @@
|
|||||||
|
"""Lyrics repository — adapter over ``AsyncSession``.
|
||||||
|
|
||||||
|
One cached row per track (``track_id`` unique). ``upsert`` refreshes the row and
|
||||||
|
bumps ``fetched_at`` so the service's TTL is measured from the last fetch.
|
||||||
|
"""
|
||||||
|
|
||||||
|
import uuid
|
||||||
|
|
||||||
|
from sqlalchemy import func, select
|
||||||
|
from sqlalchemy.dialects.postgresql import insert as pg_insert
|
||||||
|
from sqlalchemy.ext.asyncio import AsyncSession
|
||||||
|
|
||||||
|
from app.domain.entities.lyrics import Lyrics
|
||||||
|
from app.infrastructure.db.models.lyrics import LyricsModel
|
||||||
|
|
||||||
|
|
||||||
|
def _to_entity(row: LyricsModel) -> Lyrics:
|
||||||
|
return Lyrics(
|
||||||
|
track_id=row.track_id,
|
||||||
|
synced=row.synced,
|
||||||
|
plain=row.plain,
|
||||||
|
source=row.source,
|
||||||
|
status=row.status,
|
||||||
|
fetched_at=row.fetched_at,
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
class SqlAlchemyLyricsRepository:
|
||||||
|
def __init__(self, session: AsyncSession) -> None:
|
||||||
|
self._session = session
|
||||||
|
|
||||||
|
async def get(self, track_id: uuid.UUID) -> Lyrics | None:
|
||||||
|
row = await self._session.scalar(
|
||||||
|
select(LyricsModel).where(LyricsModel.track_id == track_id)
|
||||||
|
)
|
||||||
|
return _to_entity(row) if row is not None else None
|
||||||
|
|
||||||
|
async def upsert(
|
||||||
|
self,
|
||||||
|
*,
|
||||||
|
track_id: uuid.UUID,
|
||||||
|
synced: str | None,
|
||||||
|
plain: str | None,
|
||||||
|
source: str | None,
|
||||||
|
status: str,
|
||||||
|
) -> Lyrics:
|
||||||
|
values = {
|
||||||
|
"track_id": track_id,
|
||||||
|
"synced": synced,
|
||||||
|
"plain": plain,
|
||||||
|
"source": source,
|
||||||
|
"status": status,
|
||||||
|
"fetched_at": func.now(),
|
||||||
|
}
|
||||||
|
stmt = (
|
||||||
|
pg_insert(LyricsModel)
|
||||||
|
.values(**values)
|
||||||
|
.on_conflict_do_update(
|
||||||
|
index_elements=[LyricsModel.track_id],
|
||||||
|
set_={
|
||||||
|
"synced": synced,
|
||||||
|
"plain": plain,
|
||||||
|
"source": source,
|
||||||
|
"status": status,
|
||||||
|
"fetched_at": func.now(),
|
||||||
|
},
|
||||||
|
)
|
||||||
|
.returning(LyricsModel)
|
||||||
|
)
|
||||||
|
row = (await self._session.scalars(stmt)).one()
|
||||||
|
await self._session.flush()
|
||||||
|
return _to_entity(row)
|
||||||
@@ -0,0 +1,100 @@
|
|||||||
|
"""LrclibHttpClient — fetches lyrics from LRCLIB (plan §6.7).
|
||||||
|
|
||||||
|
LRCLIB is a free, keyless lyrics database. ``/api/get`` does an exact match on
|
||||||
|
artist+track+album+duration; if that misses we fall back to ``/api/search`` and
|
||||||
|
take the best-scoring hit. Graceful degradation: any network/parse error →
|
||||||
|
``fetch`` returns ``None`` (the service then caches a ``not_found``), never
|
||||||
|
raising. No API key is needed, so this provider is always "available".
|
||||||
|
"""
|
||||||
|
|
||||||
|
import httpx
|
||||||
|
|
||||||
|
from app.core.logging import get_logger
|
||||||
|
from app.domain.entities.lyrics import LyricsResult
|
||||||
|
|
||||||
|
log = get_logger(__name__)
|
||||||
|
|
||||||
|
_BASE_URL = "https://lrclib.net"
|
||||||
|
_TIMEOUT_SECONDS = 10.0
|
||||||
|
_SOURCE = "lrclib"
|
||||||
|
|
||||||
|
|
||||||
|
class LrclibHttpClient:
|
||||||
|
"""Implements :class:`app.domain.ports.LyricsProvider`."""
|
||||||
|
|
||||||
|
def __init__(self, *, user_agent: str, base_url: str = _BASE_URL) -> None:
|
||||||
|
self._user_agent = user_agent
|
||||||
|
self._base_url = base_url.rstrip("/")
|
||||||
|
|
||||||
|
async def fetch(
|
||||||
|
self,
|
||||||
|
*,
|
||||||
|
artist: str,
|
||||||
|
title: str,
|
||||||
|
album: str | None,
|
||||||
|
duration_seconds: int | None,
|
||||||
|
) -> LyricsResult | None:
|
||||||
|
try:
|
||||||
|
async with httpx.AsyncClient(
|
||||||
|
timeout=_TIMEOUT_SECONDS,
|
||||||
|
headers={"User-Agent": self._user_agent},
|
||||||
|
base_url=self._base_url,
|
||||||
|
) as client:
|
||||||
|
hit = await self._get(client, artist, title, album, duration_seconds)
|
||||||
|
if hit is None:
|
||||||
|
hit = await self._search(client, artist, title)
|
||||||
|
except (httpx.HTTPError, ValueError) as exc:
|
||||||
|
log.warning("lrclib.fetch_failed", error=str(exc))
|
||||||
|
return None
|
||||||
|
return hit
|
||||||
|
|
||||||
|
async def _get(
|
||||||
|
self,
|
||||||
|
client: httpx.AsyncClient,
|
||||||
|
artist: str,
|
||||||
|
title: str,
|
||||||
|
album: str | None,
|
||||||
|
duration_seconds: int | None,
|
||||||
|
) -> LyricsResult | None:
|
||||||
|
"""Exact match via ``/api/get`` (404 when nothing matches exactly)."""
|
||||||
|
params = {"artist_name": artist, "track_name": title}
|
||||||
|
if album:
|
||||||
|
params["album_name"] = album
|
||||||
|
if duration_seconds is not None:
|
||||||
|
params["duration"] = str(duration_seconds)
|
||||||
|
resp = await client.get("/api/get", params=params)
|
||||||
|
if resp.status_code == httpx.codes.NOT_FOUND:
|
||||||
|
return None
|
||||||
|
resp.raise_for_status()
|
||||||
|
return _to_result(resp.json())
|
||||||
|
|
||||||
|
async def _search(
|
||||||
|
self, client: httpx.AsyncClient, artist: str, title: str
|
||||||
|
) -> LyricsResult | None:
|
||||||
|
"""Fuzzy fallback via ``/api/search`` — take the first usable hit."""
|
||||||
|
resp = await client.get(
|
||||||
|
"/api/search", params={"artist_name": artist, "track_name": title}
|
||||||
|
)
|
||||||
|
resp.raise_for_status()
|
||||||
|
results = resp.json()
|
||||||
|
if not isinstance(results, list):
|
||||||
|
return None
|
||||||
|
for item in results:
|
||||||
|
result = _to_result(item)
|
||||||
|
if result is not None:
|
||||||
|
return result
|
||||||
|
return None
|
||||||
|
|
||||||
|
|
||||||
|
def _to_result(payload: object) -> LyricsResult | None:
|
||||||
|
"""Map an LRCLIB record to a ``LyricsResult``. Instrumental tracks and empty
|
||||||
|
records yield ``None`` (nothing worth caching as "found")."""
|
||||||
|
if not isinstance(payload, dict):
|
||||||
|
return None
|
||||||
|
if payload.get("instrumental"):
|
||||||
|
return None
|
||||||
|
synced = payload.get("syncedLyrics") or None
|
||||||
|
plain = payload.get("plainLyrics") or None
|
||||||
|
if synced is None and plain is None:
|
||||||
|
return None
|
||||||
|
return LyricsResult(synced=synced, plain=plain, source=_SOURCE)
|
||||||
@@ -0,0 +1 @@
|
|||||||
|
"""ffmpeg-based transcoding adapters (Group B / plan §6.6)."""
|
||||||
@@ -0,0 +1,67 @@
|
|||||||
|
"""FfmpegTranscoder — Opus + HLS renditions via the ffmpeg CLI.
|
||||||
|
|
||||||
|
Implements :class:`app.domain.ports.Transcoder`. Runs ffmpeg as a subprocess
|
||||||
|
(only ever from a worker — CLAUDE.md: no heavy work in the request cycle) and
|
||||||
|
raises :class:`TranscodeError` on a non-zero exit. Opus is used for the direct
|
||||||
|
quality renditions; HLS segments are AAC-in-MPEG-TS (broad player support).
|
||||||
|
"""
|
||||||
|
|
||||||
|
import asyncio
|
||||||
|
from pathlib import Path
|
||||||
|
|
||||||
|
import anyio
|
||||||
|
|
||||||
|
from app.core.logging import get_logger
|
||||||
|
from app.domain.errors import TranscodeError
|
||||||
|
|
||||||
|
log = get_logger(__name__)
|
||||||
|
|
||||||
|
# HLS segment length. 10s is a common VOD default — few requests, quick seeks.
|
||||||
|
_HLS_SEGMENT_SECONDS = "10"
|
||||||
|
|
||||||
|
|
||||||
|
def _mkdir(path: Path) -> None:
|
||||||
|
path.mkdir(parents=True, exist_ok=True)
|
||||||
|
|
||||||
|
|
||||||
|
class FfmpegTranscoder:
|
||||||
|
def __init__(self, ffmpeg_path: str = "ffmpeg") -> None:
|
||||||
|
self._ffmpeg = ffmpeg_path
|
||||||
|
|
||||||
|
async def to_opus(self, src: Path, dest: Path, *, bitrate_kbps: int) -> None:
|
||||||
|
await anyio.to_thread.run_sync(_mkdir, dest.parent)
|
||||||
|
tmp = dest.with_suffix(dest.suffix + ".part")
|
||||||
|
await self._run(
|
||||||
|
"-i", str(src),
|
||||||
|
"-vn", "-c:a", "libopus", "-b:a", f"{bitrate_kbps}k",
|
||||||
|
"-f", "opus", str(tmp),
|
||||||
|
)
|
||||||
|
# Publish atomically so a concurrent reader never sees a half-written file.
|
||||||
|
await anyio.to_thread.run_sync(tmp.replace, dest)
|
||||||
|
|
||||||
|
async def to_hls(self, src: Path, out_dir: Path, *, bitrate_kbps: int) -> None:
|
||||||
|
await anyio.to_thread.run_sync(_mkdir, out_dir)
|
||||||
|
playlist = out_dir / "playlist.m3u8"
|
||||||
|
segments = out_dir / "seg_%03d.ts"
|
||||||
|
await self._run(
|
||||||
|
"-i", str(src),
|
||||||
|
"-vn", "-c:a", "aac", "-b:a", f"{bitrate_kbps}k",
|
||||||
|
"-f", "hls",
|
||||||
|
"-hls_time", _HLS_SEGMENT_SECONDS,
|
||||||
|
"-hls_playlist_type", "vod",
|
||||||
|
"-hls_segment_filename", str(segments),
|
||||||
|
str(playlist),
|
||||||
|
)
|
||||||
|
|
||||||
|
async def _run(self, *args: str) -> None:
|
||||||
|
cmd = [self._ffmpeg, "-y", "-nostdin", "-loglevel", "error", *args]
|
||||||
|
proc = await asyncio.create_subprocess_exec(
|
||||||
|
*cmd,
|
||||||
|
stdout=asyncio.subprocess.DEVNULL,
|
||||||
|
stderr=asyncio.subprocess.PIPE,
|
||||||
|
)
|
||||||
|
_, stderr = await proc.communicate()
|
||||||
|
if proc.returncode != 0:
|
||||||
|
detail = stderr.decode(errors="replace").strip()[-500:]
|
||||||
|
log.error("ffmpeg_failed", returncode=proc.returncode, error=detail)
|
||||||
|
raise TranscodeError(f"ffmpeg exited {proc.returncode}: {detail}")
|
||||||
@@ -14,6 +14,7 @@ from app.workers.tasks.download_task import download_track
|
|||||||
from app.workers.tasks.enrich_task import enrich_track
|
from app.workers.tasks.enrich_task import enrich_track
|
||||||
from app.workers.tasks.import_task import scan_local_folder
|
from app.workers.tasks.import_task import scan_local_folder
|
||||||
from app.workers.tasks.materialize_task import materialize_track
|
from app.workers.tasks.materialize_task import materialize_track
|
||||||
|
from app.workers.tasks.transcode_task import transcode_track
|
||||||
|
|
||||||
log = get_logger("worker")
|
log = get_logger("worker")
|
||||||
|
|
||||||
@@ -35,6 +36,7 @@ class WorkerSettings:
|
|||||||
download_track,
|
download_track,
|
||||||
materialize_track,
|
materialize_track,
|
||||||
cleanup_storage,
|
cleanup_storage,
|
||||||
|
transcode_track,
|
||||||
]
|
]
|
||||||
on_startup = startup
|
on_startup = startup
|
||||||
on_shutdown = shutdown
|
on_shutdown = shutdown
|
||||||
|
|||||||
@@ -59,6 +59,28 @@ async def enqueue_materialize(job_id: uuid.UUID) -> None:
|
|||||||
log.warning("materialize_enqueue_failed", job_id=str(job_id))
|
log.warning("materialize_enqueue_failed", job_id=str(job_id))
|
||||||
|
|
||||||
|
|
||||||
|
async def enqueue_transcode(
|
||||||
|
track_id: uuid.UUID, *, quality: str = "high", hls: bool = True
|
||||||
|
) -> str:
|
||||||
|
"""Enqueue a transcode job (§6.6). Unlike the best-effort follow-ups below,
|
||||||
|
this is user/stream-driven, so it surfaces a job id (and a 503 via the caller
|
||||||
|
if the queue is down) rather than swallowing failures."""
|
||||||
|
return await enqueue(
|
||||||
|
"transcode_track", track_id=str(track_id), quality=quality, hls=hls
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
async def enqueue_transcode_quiet(
|
||||||
|
track_id: uuid.UUID, *, quality: str = "high", hls: bool = False
|
||||||
|
) -> None:
|
||||||
|
"""Best-effort transcode enqueue for a streaming cache miss — never blocks or
|
||||||
|
fails the play (the master is served meanwhile)."""
|
||||||
|
try:
|
||||||
|
await enqueue_transcode(track_id, quality=quality, hls=hls)
|
||||||
|
except DependencyUnavailableError:
|
||||||
|
log.warning("transcode_enqueue_failed", track_id=str(track_id))
|
||||||
|
|
||||||
|
|
||||||
async def enqueue_enrich(track_id: uuid.UUID) -> None:
|
async def enqueue_enrich(track_id: uuid.UUID) -> None:
|
||||||
"""Best-effort enqueue of metadata enrichment for a freshly stored track.
|
"""Best-effort enqueue of metadata enrichment for a freshly stored track.
|
||||||
|
|
||||||
|
|||||||
@@ -0,0 +1,64 @@
|
|||||||
|
"""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 uuid
|
||||||
|
from typing import Any
|
||||||
|
|
||||||
|
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")
|
||||||
|
|
||||||
|
|
||||||
|
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():
|
||||||
|
await transcoder.to_hls(src, hls_dir(root, tid), bitrate_kbps=HLS_BITRATE)
|
||||||
|
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}
|
||||||
Reference in New Issue
Block a user