mirror of
https://github.com/navidrome/navidrome.git
synced 2026-10-11 20:07:11 +02:00
* refactor(persistence): adopt generic deluan/rest repository API Pin deluan/rest to the refactor branch. REST-facing repository methods take a context and return typed values. Drop DataStore.Resource and ResourceRepository; the native API names typed repositories directly through a per-request adapter that later commits remove. * refactor(persistence): base repository helpers take a context * refactor(persistence): LibraryRepository takes a context per call * refactor(persistence): PropertyRepository takes a context per call * refactor(persistence): UserPropsRepository takes a context per call * refactor(persistence): TranscodingRepository takes a context per call * refactor(persistence): ShareRepository takes a context per call * refactor(persistence): PlayerRepository takes a context per call * refactor(persistence): RadioRepository takes a context per call * refactor(persistence): PlayQueueRepository takes a context per call * refactor(persistence): Tag and Genre repositories take a context per call * refactor(persistence): PluginRepository takes a context per call * refactor(persistence): Scrobble repositories take a context per call * refactor(persistence): FolderRepository takes a context per call * refactor(persistence): Artwork repositories take a context per call * refactor(persistence): UserRepository takes a context per call * refactor(persistence): ArtistRepository takes a context per call ReadAll no longer rewrites the shared sort mappings for the role filter; it works on a per-call copy. * test(persistence): assert artist role sort sanitization in ReadAll * refactor(persistence): AlbumRepository takes a context per call * test(persistence): pass the test context to album repository helpers * refactor(persistence): MediaFileRepository takes a context per call * refactor(persistence): Playlist repositories take a context per call * refactor(persistence): build all repositories once per store * refactor(core): REST repository wrappers are built once * refactor(persistence): repositories are stateless Remove the context field from the base repository and the per-request REST adapter. Enable the containedctx linter so no repository can hold a request context again. * chore(lint): skip containedctx in test files * refactor: share simplifications from the stateless repositories sweep Add deleteOwnedAll on sqlRepository and use it in player/share Delete to remove the duplicated bulk-delete loop; have Share.Repository() return model.ShareRepository so subsonic sharing.go drops its repeated type assertions. * chore(core): assert REST wrappers implement Persistable * chore: reformat imports * perf(persistence): build repositories on first use Each transaction store used to construct all 21 repositories up front, paying for filter and sort mapping setup the block never touched. Fields are now sync.OnceValue thunks, so a store only builds what it uses. * fix(persistence): clean plugin references per deleted user A bulk user delete that fails on a later id had already removed the earlier rows but skipped their plugin cleanup. Cleanup now runs right after each successful delete. * fix(core): unload disabled plugins even when a user delete fails A bulk delete can fail on a later id after earlier users were removed and their plugins auto-disabled. The wrapper returned before unloading, leaving those plugins running until the next successful delete or a restart. * chore(deps): pin deluan/rest to v1.0.1 Replaces the pseudo-version of the refactor branch with the tagged release. REST error messages now name the bare type (Artist, not model.Artist). * test: use the spec context instead of context.Background() Replace the context.Background()/context.TODO() calls this branch added to tests with the spec's ctx, GinkgoT().Context(), or t/b.Context(), so repository calls are bound to the running spec's lifetime. * test: declare the spec context once per Describe Set ctx from GinkgoT().Context() first in each top-level BeforeEach and reuse it, building user contexts on top of it instead of repeating inline calls.
305 lines
10 KiB
Go
305 lines
10 KiB
Go
package stream
|
|
|
|
import (
|
|
"context"
|
|
"errors"
|
|
"fmt"
|
|
"io"
|
|
"mime"
|
|
"net/http"
|
|
"os"
|
|
"strconv"
|
|
"sync"
|
|
"time"
|
|
|
|
"github.com/navidrome/navidrome/conf"
|
|
"github.com/navidrome/navidrome/consts"
|
|
"github.com/navidrome/navidrome/core/ffmpeg"
|
|
"github.com/navidrome/navidrome/log"
|
|
"github.com/navidrome/navidrome/model"
|
|
"github.com/navidrome/navidrome/model/request"
|
|
"github.com/navidrome/navidrome/utils/cache"
|
|
"github.com/navidrome/navidrome/utils/req"
|
|
)
|
|
|
|
type MediaStreamer interface {
|
|
NewStream(ctx context.Context, mf *model.MediaFile, req Request) (*Stream, error)
|
|
}
|
|
|
|
type TranscodingCache cache.FileCache
|
|
|
|
func NewMediaStreamer(ds model.DataStore, t ffmpeg.FFmpeg, cache TranscodingCache) MediaStreamer {
|
|
return &mediaStreamer{
|
|
ds: ds,
|
|
transcoder: t,
|
|
cache: cache,
|
|
limiter: NewTranscodeLimiter(conf.Server.Transcoding.MaxConcurrent, conf.Server.Transcoding.MaxConcurrentPerUser),
|
|
}
|
|
}
|
|
|
|
type mediaStreamer struct {
|
|
ds model.DataStore
|
|
transcoder ffmpeg.FFmpeg
|
|
cache cache.FileCache
|
|
limiter TranscodeLimiter
|
|
}
|
|
|
|
type streamJob struct {
|
|
ms *mediaStreamer
|
|
mf *model.MediaFile
|
|
filePath string
|
|
format string
|
|
bitRate int
|
|
sampleRate int
|
|
bitDepth int
|
|
channels int
|
|
offset int
|
|
}
|
|
|
|
func (j *streamJob) Key() string {
|
|
return fmt.Sprintf("%s.%s.%d.%d.%d.%d.%s.%d", j.mf.ID, j.mf.UpdatedAt.Format(time.RFC3339Nano), j.bitRate, j.sampleRate, j.bitDepth, j.channels, j.format, j.offset)
|
|
}
|
|
|
|
// NewStream creates a Stream for the given MediaFile and Request. It handles both raw streaming (no transcoding)
|
|
// and transcoded streaming based on the requested format and bitrate. It also logs detailed information about
|
|
// the streaming request and whether the transcoding result was served from cache or not.
|
|
func (ms *mediaStreamer) NewStream(ctx context.Context, mf *model.MediaFile, req Request) (*Stream, error) {
|
|
var format string
|
|
var bitRate int
|
|
var cached bool
|
|
defer func() {
|
|
log.Info(ctx, "Streaming file", "title", mf.Title, "artist", mf.Artist, "format", format, "cached", cached,
|
|
"bitRate", bitRate, "sampleRate", req.SampleRate, "bitDepth", req.BitDepth, "channels", req.Channels,
|
|
"user", userName(ctx), "transcoding", format != "raw",
|
|
"originalFormat", mf.Suffix, "originalBitRate", mf.BitRate)
|
|
}()
|
|
|
|
format = req.Format
|
|
bitRate = req.BitRate
|
|
if format == "" || format == "raw" {
|
|
format = "raw"
|
|
bitRate = 0
|
|
}
|
|
s := &Stream{ctx: ctx, mf: mf, format: format, bitRate: bitRate}
|
|
filePath := mf.AbsolutePath()
|
|
|
|
if format == "raw" {
|
|
log.Debug(ctx, "Streaming RAW file", "id", mf.ID, "path", filePath,
|
|
"requestBitrate", req.BitRate, "requestFormat", req.Format, "requestOffset", req.Offset,
|
|
"originalBitrate", mf.BitRate, "originalFormat", mf.Suffix,
|
|
"selectedBitrate", bitRate, "selectedFormat", format)
|
|
f, err := os.Open(filePath)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
s.ReadCloser = f
|
|
s.Seeker = f
|
|
s.format = mf.Suffix
|
|
return s, nil
|
|
}
|
|
|
|
job := &streamJob{
|
|
ms: ms,
|
|
mf: mf,
|
|
filePath: filePath,
|
|
format: format,
|
|
bitRate: bitRate,
|
|
sampleRate: req.SampleRate,
|
|
bitDepth: req.BitDepth,
|
|
channels: req.Channels,
|
|
offset: req.Offset,
|
|
}
|
|
r, err := ms.cache.Get(ctx, job)
|
|
if err != nil {
|
|
// Rate-limit rejections are already logged at warn level by the
|
|
// producer; treating them as cache failures here would both
|
|
// double-log and mask actual cache problems.
|
|
if !errors.Is(err, ErrTooManyTranscodes) {
|
|
log.Error(ctx, "Error accessing transcoding cache", "id", mf.ID, err)
|
|
}
|
|
return nil, err
|
|
}
|
|
cached = r.Cached
|
|
|
|
s.ReadCloser = r
|
|
s.Seeker = r.Seeker
|
|
|
|
log.Debug(ctx, "Streaming TRANSCODED file", "id", mf.ID, "path", filePath,
|
|
"requestBitrate", req.BitRate, "requestFormat", req.Format, "requestOffset", req.Offset,
|
|
"originalBitrate", mf.BitRate, "originalFormat", mf.Suffix,
|
|
"selectedBitrate", bitRate, "selectedFormat", format, "cached", cached, "seekable", s.Seekable())
|
|
|
|
return s, nil
|
|
}
|
|
|
|
type Stream struct {
|
|
ctx context.Context //nolint:containedctx // stream outlives the call that built it; Read has no ctx
|
|
mf *model.MediaFile
|
|
bitRate int
|
|
format string
|
|
io.ReadCloser
|
|
io.Seeker
|
|
}
|
|
|
|
func (s *Stream) Seekable() bool { return s.Seeker != nil }
|
|
func (s *Stream) Duration() float32 { return s.mf.Duration }
|
|
func (s *Stream) ContentType() string { return mime.TypeByExtension("." + s.format) }
|
|
func (s *Stream) Name() string { return s.mf.Title + "." + s.format }
|
|
func (s *Stream) ModTime() time.Time { return s.mf.UpdatedAt }
|
|
func (s *Stream) EstimatedContentLength() int {
|
|
return int(s.mf.Duration * float32(s.bitRate) / 8 * 1024)
|
|
}
|
|
|
|
// Serve writes the stream to the HTTP response. For seekable streams it uses http.ServeContent
|
|
// (supporting range requests). For non-seekable streams it writes directly and logs any errors.
|
|
// Returns the number of bytes written and an error only when it fails with 0 bytes written
|
|
// (meaning the HTTP 200 status has not been flushed yet and the caller can still send an error response).
|
|
// Once bytes are on the wire it panics with http.ErrAbortHandler instead, aborting the response.
|
|
// Empty output (0 bytes, no error) is logged but not treated as an error.
|
|
func (s *Stream) Serve(ctx context.Context, w http.ResponseWriter, r *http.Request) (int64, error) {
|
|
if s.Seekable() {
|
|
http.ServeContent(w, r, s.Name(), s.ModTime(), s)
|
|
return -1, nil
|
|
}
|
|
|
|
w.Header().Set("Accept-Ranges", "none")
|
|
w.Header().Set("Content-Type", s.ContentType())
|
|
|
|
if req.Params(r).BoolOr("estimateContentLength", false) {
|
|
length := strconv.Itoa(s.EstimatedContentLength())
|
|
log.Trace(ctx, "Estimated content-length", "contentLength", length)
|
|
w.Header().Set("Content-Length", length)
|
|
}
|
|
|
|
if r.Method == http.MethodHead {
|
|
go func() { _, _ = io.Copy(io.Discard, s) }()
|
|
return 0, nil
|
|
}
|
|
|
|
id := s.mf.ID
|
|
c, err := io.Copy(w, s)
|
|
if err != nil {
|
|
log.Error(ctx, "Error sending transcoded file", "id", id, err)
|
|
if c == 0 {
|
|
w.Header().Del("Content-Length")
|
|
return 0, fmt.Errorf("sending transcoded file: %w", err)
|
|
}
|
|
// The 200 is already sent, so dropping the connection is the only way to say "truncated".
|
|
panic(http.ErrAbortHandler)
|
|
}
|
|
if c == 0 {
|
|
log.Error(ctx, "Transcoding returned empty output, ffmpeg may have failed. "+
|
|
"Check that ffmpeg supports the requested codec. Enable Trace logging for ffmpeg stderr details",
|
|
"id", id, "format", s.ContentType())
|
|
} else {
|
|
log.Trace(ctx, "Success sending transcoded file", "id", id, "size", c)
|
|
}
|
|
return c, nil
|
|
}
|
|
|
|
// NewStream creates a non-seekable Stream from the given components.
|
|
func NewStream(mf *model.MediaFile, format string, bitRate int, r io.ReadCloser) *Stream {
|
|
return &Stream{
|
|
ctx: context.Background(),
|
|
mf: mf,
|
|
format: format,
|
|
bitRate: bitRate,
|
|
ReadCloser: r,
|
|
}
|
|
}
|
|
|
|
var (
|
|
onceTranscodingCache sync.Once
|
|
instanceTranscodingCache TranscodingCache
|
|
)
|
|
|
|
func GetTranscodingCache() TranscodingCache {
|
|
onceTranscodingCache.Do(func() {
|
|
instanceTranscodingCache = NewTranscodingCache()
|
|
})
|
|
return instanceTranscodingCache
|
|
}
|
|
|
|
func NewTranscodingCache() TranscodingCache {
|
|
return cache.NewFileCache("Transcoding", conf.Server.TranscodingCacheSize,
|
|
consts.TranscodingCacheDir, consts.DefaultTranscodingCacheMaxItems,
|
|
func(ctx context.Context, arg cache.Item) (io.Reader, error) {
|
|
job := arg.(*streamJob)
|
|
command := LookupTranscodeCommand(ctx, job.ms.ds, job.format)
|
|
if command == "" {
|
|
log.Error(ctx, "No transcoding command available", "format", job.format)
|
|
return nil, os.ErrInvalid
|
|
}
|
|
|
|
release, err := job.ms.limiter.Acquire(ctx, limiterKey(ctx))
|
|
if err != nil {
|
|
log.Warn(ctx, "Refusing transcode: concurrent transcode limit reached",
|
|
"id", job.mf.ID, "user", userName(ctx),
|
|
"maxConcurrent", conf.Server.Transcoding.MaxConcurrent,
|
|
"maxPerUser", conf.Server.Transcoding.MaxConcurrentPerUser)
|
|
return nil, err
|
|
}
|
|
|
|
// Choose the context that drives the ffmpeg process.
|
|
//
|
|
// When the limiter is enabled, force the request context so a
|
|
// client disconnect cancels ffmpeg and frees the slot promptly.
|
|
// Otherwise a client could open many transcodes, disconnect
|
|
// immediately, and still leave the configured cap's worth of
|
|
// ffmpeg processes draining in the background — which is exactly
|
|
// the DoS the limiter is meant to prevent.
|
|
//
|
|
// When the limiter is disabled, preserve the legacy behavior
|
|
// governed by Transcoding.EnableCancellation so unchanged configs
|
|
// keep their previous observable behavior.
|
|
var transcodingCtx context.Context
|
|
if job.ms.limiter.Enabled() || conf.Server.Transcoding.EnableCancellation {
|
|
transcodingCtx = ctx
|
|
} else {
|
|
transcodingCtx = request.AddValues(context.Background(), ctx)
|
|
}
|
|
|
|
out, err := job.ms.transcoder.Transcode(transcodingCtx, ffmpeg.TranscodeOptions{
|
|
Command: command,
|
|
Format: job.format,
|
|
FilePath: job.filePath,
|
|
BitRate: job.bitRate,
|
|
SampleRate: job.sampleRate,
|
|
BitDepth: job.bitDepth,
|
|
Channels: job.channels,
|
|
Offset: job.offset,
|
|
Duration: job.mf.Duration,
|
|
})
|
|
if err != nil {
|
|
release()
|
|
log.Error(ctx, "Error starting transcoder", "id", job.mf.ID, err)
|
|
return nil, os.ErrInvalid
|
|
}
|
|
// Tie the slot to the ffmpeg process: copyAndClose calls Close
|
|
// on this reader after io.Copy returns, which is exactly when
|
|
// ffmpeg has exited (either EOF or context cancellation).
|
|
return &releasingReadCloser{ReadCloser: out, release: release}, nil
|
|
})
|
|
}
|
|
|
|
// userName extracts the username from the context for logging purposes.
|
|
func userName(ctx context.Context) string {
|
|
if user, ok := request.UserFrom(ctx); !ok {
|
|
return "UNKNOWN"
|
|
} else {
|
|
return user.UserName
|
|
}
|
|
}
|
|
|
|
// limiterKey returns the per-user bucket key used by the transcode limiter.
|
|
// For anonymous requests (e.g. public shares) it returns the empty string,
|
|
// which signals the limiter to skip the per-user cap entirely — otherwise
|
|
// every anonymous viewer of a public share would collide on the same key
|
|
// and starve each other within MaxConcurrentPerUser slots. The global cap
|
|
// still applies and remains the protection against runaway anonymous load.
|
|
func limiterKey(ctx context.Context) string {
|
|
if user, ok := request.UserFrom(ctx); ok {
|
|
return user.UserName
|
|
}
|
|
return ""
|
|
}
|